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> Action, CancellationTokenSource Cancel); private readonly Channel queue = Channel.CreateBounded(new BoundedChannelOptions(QueueCapacity) { SingleReader = true, FullMode = BoundedChannelFullMode.Wait }); private readonly ConcurrentDictionary cancellations = new(); private readonly object admissionGate = new(); private bool stopping; public OperationRecord Submit(ServiceIdentity identity, string accountId, AgentCapability capability, string? idempotencyKey, string canonicalParameters, Func 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> 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 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(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); } }