Files
rogee c7c0ab273f
Build web service image / build (push) Successful in 1m9s
feat: validate single-client broadcast operations
2026-09-19 14:33:26 +08:00

159 lines
7.8 KiB
C#

using System.Collections.Concurrent;
using System.Security.Cryptography;
using System.Text;
using System.Threading.Channels;
using Microsoft.Extensions.Hosting;
namespace WxAgent.Service;
public sealed record OperationExecutionResult(string? Details = null, string? ErrorCode = null, bool Unconfirmed = false);
public sealed class OperationQueue(OperationStore store, ServiceSecurity security) : BackgroundService
{
private const int QueueCapacity = 100;
private sealed record Work(OperationRecord Record, ServiceIdentity Identity, string Permission,
Func<CancellationToken, Task<OperationExecutionResult>> Action, CancellationTokenSource Cancel);
private readonly Channel<Work> queue = Channel.CreateBounded<Work>(new BoundedChannelOptions(QueueCapacity)
{ SingleReader = true, FullMode = BoundedChannelFullMode.Wait });
private readonly ConcurrentDictionary<string, CancellationTokenSource> cancellations = new();
private readonly object admissionGate = new();
private bool stopping;
public OperationRecord Submit(ServiceIdentity identity, string accountId, AgentCapability capability,
string? idempotencyKey, string canonicalParameters, Func<CancellationToken, Task> action) =>
SubmitResult(identity, accountId, capability, idempotencyKey, canonicalParameters, null,
async cancellationToken =>
{
await action(cancellationToken).ConfigureAwait(false);
return new OperationExecutionResult();
});
public OperationRecord SubmitResult(ServiceIdentity identity, string accountId, AgentCapability capability,
string? idempotencyKey, string canonicalParameters, string? initialDetails,
Func<CancellationToken, Task<OperationExecutionResult>> action)
{
security.RequireCurrent(identity, capability.Permission, accountId);
if (!capability.Enabled) throw new ServiceException("CapabilityDisabled", 409, capability.DisabledReason ?? "Capability unavailable.");
if (capability.TimeoutSeconds is < 1 or > 300) throw new ServiceException("InvalidRequest", 400, "Invalid execution budget.");
if (idempotencyKey is { Length: < 1 or > 128 } || (capability.HasSideEffects && string.IsNullOrWhiteSpace(idempotencyKey)))
throw new ServiceException("InvalidRequest", 400, "Side effects require a nonempty idempotency key of at most 128 characters.");
var digest = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(canonicalParameters)));
lock (admissionGate)
{
if (stopping) throw new ServiceException("Unavailable", 503, "Agent is stopping.");
var (record, created) = store.Enqueue(identity.PrincipalId, accountId, capability.Operation,
idempotencyKey, digest, capability.HasSideEffects, TimeSpan.FromSeconds(capability.TimeoutSeconds), QueueCapacity, initialDetails);
if (!created) return record;
var cancel = new CancellationTokenSource();
cancellations[record.Id] = cancel;
if (!queue.Writer.TryWrite(new Work(record, identity, capability.Permission, action, cancel)))
{
cancellations.TryRemove(record.Id, out _);
cancel.Dispose();
store.Transition(record.Id, "Failed", "admission", "QueueFull");
throw new ServiceException("QueueFull", 429, "Agent queue is full.");
}
return record;
}
}
public Page<OperationSummary> List(ServiceIdentity identity, string? accountId, int limit, int offset)
{
security.RequireCurrent(identity, "read");
var page = store.List(identity.PrincipalId, accountId, limit, offset);
return new Page<OperationSummary>(page.Items.Select(OperationSummary.From).ToArray(), page.Limit,
page.Offset, page.HasMore, page.NextOffset);
}
public OperationRecord Get(ServiceIdentity identity, string id)
{
security.RequireCurrent(identity, "read");
return store.Get(id, identity.PrincipalId);
}
public OperationRecord Cancel(ServiceIdentity identity, string id)
{
lock (admissionGate)
{
var record = Get(identity, id);
if (cancellations.TryGetValue(id, out var cancel)) cancel.Cancel();
if (record.State == "Queued") store.Transition(id, "Cancelled", "queued", "Cancelled");
}
return Get(identity, id);
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
try
{
await foreach (var work in queue.Reader.ReadAllAsync(stoppingToken))
{
var started = false;
try
{
if (store.Get(work.Record.Id, work.Identity.PrincipalId).State != "Queued") continue;
var remaining = work.Record.ExpiresAt - DateTimeOffset.UtcNow;
if (remaining <= TimeSpan.Zero) throw new ServiceException("Timeout", 408, "Queue budget expired.");
using var budget = CancellationTokenSource.CreateLinkedTokenSource(work.Cancel.Token, stoppingToken);
budget.CancelAfter(remaining);
lock (admissionGate)
{
budget.Token.ThrowIfCancellationRequested();
security.RequireCurrent(work.Identity, work.Permission, work.Record.AccountId);
store.Transition(work.Record.Id, "Running", "execution");
started = true;
}
// Never release this slot with WaitAsync: the actual action must have stopped first.
var result = await work.Action(budget.Token);
if (result.ErrorCode is not null)
{
store.Transition(work.Record.Id, result.Unconfirmed ? "Unconfirmed" : "Failed",
"execution", result.ErrorCode, result.Details);
continue;
}
budget.Token.ThrowIfCancellationRequested();
security.RequireCurrent(work.Identity, work.Permission, work.Record.AccountId);
store.Transition(work.Record.Id, "Succeeded", "complete", details: result.Details);
}
catch (Exception e)
{
var unknown = started && work.Record.HasSideEffects;
var cancelled = e is OperationCanceledException;
var error = e is ServiceException se ? se.Code : cancelled ? "Cancelled" : "ExecutionFailed";
store.Transition(work.Record.Id, unknown ? "Unconfirmed" : cancelled ? "Cancelled" : "Failed",
started ? "execution" : "admission", error);
}
finally
{
lock (admissionGate)
{
cancellations.TryRemove(work.Record.Id, out _);
work.Cancel.Dispose();
}
}
}
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { }
finally
{
lock (admissionGate)
{
stopping = true;
queue.Writer.TryComplete();
while (queue.Reader.TryRead(out var queued))
{
store.Transition(queued.Record.Id, "Cancelled", "shutdown", "AgentStopping");
cancellations.TryRemove(queued.Record.Id, out _);
queued.Cancel.Dispose();
}
}
}
}
public override Task StopAsync(CancellationToken cancellationToken)
{
lock (admissionGate) { stopping = true; queue.Writer.TryComplete(); }
return base.StopAsync(cancellationToken);
}
}