159 lines
7.8 KiB
C#
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);
|
|
}
|
|
}
|