diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs index 00090c983..6af8c9c50 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs @@ -2,6 +2,9 @@ using CodeSpace.Core.Persistence.Entities; using CodeSpace.Messages.Enums; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Npgsql; using System.Text.RegularExpressions; namespace CodeSpace.Core.Services.Agents.AgentRunLogging; @@ -18,10 +21,11 @@ public sealed partial class AgentRunLogCaptureRecoveryService : IAgentRunLogCapt private readonly DbContextOptions _dbOptions; private readonly IAgentRunLogService _logs; private readonly AgentRunLogCaptureRecoveryOptions _options; + private readonly ILogger _logger; - public AgentRunLogCaptureRecoveryService(DbContextOptions dbOptions, IAgentRunLogService logs) : this(dbOptions, logs, Defaults) { } + public AgentRunLogCaptureRecoveryService(DbContextOptions dbOptions, IAgentRunLogService logs, ILogger logger) : this(dbOptions, logs, Defaults, logger) { } - internal AgentRunLogCaptureRecoveryService(DbContextOptions dbOptions, IAgentRunLogService logs, AgentRunLogCaptureRecoveryOptions options) + internal AgentRunLogCaptureRecoveryService(DbContextOptions dbOptions, IAgentRunLogService logs, AgentRunLogCaptureRecoveryOptions options, ILogger? logger = null) { if (options.BatchSize is <= 0 or > 500 || options.MaxConcurrency is <= 0 or > 32 || options.MaxConcurrency > options.BatchSize || options.LeaseDuration <= options.OperationTimeout + options.OperationTimeout + MinimumLeaseMargin || options.OperationTimeout <= TimeSpan.Zero @@ -31,6 +35,7 @@ internal AgentRunLogCaptureRecoveryService(DbContextOptions _dbOptions = dbOptions; _logs = logs; _options = options; + _logger = logger ?? NullLogger.Instance; } public async Task DeclareAsync(AgentRunLogCaptureDeclarationRequest request, CancellationToken cancellationToken) @@ -110,11 +115,26 @@ private async Task RecoverClaimAsync(RecoveryClaim claim, Ca } using var settlement = new CancellationTokenSource(_options.OperationTimeout); + var attempt = new SettlementAttempt(outcome); try { - return await SettleAsync(claim, outcome, settlement.Token).ConfigureAwait(false); + return await SettleAsync(claim, attempt, settlement.Token).ConfigureAwait(false); } - catch (Exception) when (!cancellationToken.IsCancellationRequested) { return new RecoverySettlement(outcome.State, true); } + catch (Exception exception) when (!cancellationToken.IsCancellationRequested) + { + LogUnsettledClaim(claim, attempt.Outcome, exception); + return new RecoverySettlement(outcome.State, true); + } + } + + // A settlement that raises is treated as unwritten: its claim keeps this wave's owner and lease, so nothing re-claims + // it until the lease expires. It is counted as a lost lease, and its cause is visible only here. The outcome named is + // the one the settlement was writing, which may have replaced the observed one; the guard judges that write. + private void LogUnsettledClaim(RecoveryClaim claim, RecoveryOutcome outcome, Exception exception) + { + var refusal = DatabaseRefusal(exception); + + _logger.LogWarning(exception, "Agent run {RunId} log capture intent {IntentId} could not settle as {Outcome} with last_error_code {OutcomeCode} (SQLSTATE {SqlState}: {MessageText}); its claim stays leased until its recovery lease expires, then a later wave re-claims it", claim.AgentRunId, claim.Id, outcome.State, outcome.ErrorCode, refusal?.SqlState, refusal?.MessageText); } private async Task> ClaimBatchAsync(Guid ownerId, DateTimeOffset cutoff, int limit, CancellationToken cancellationToken) @@ -242,8 +262,9 @@ private async Task FailStreamAsync(RecoveryClaim claim, AgentRu : RecoveryOutcome.Retry(stream.Id, claim.State, $"fail-{Code(problem.Code)}", "The stream health transition could not yet be persisted."); } - private async Task SettleAsync(RecoveryClaim claim, RecoveryOutcome outcome, CancellationToken cancellationToken) + private async Task SettleAsync(RecoveryClaim claim, SettlementAttempt attempt, CancellationToken cancellationToken) { + var outcome = attempt.Outcome; await using var db = CreateDb(); await using var transaction = await db.Database.BeginTransactionAsync(cancellationToken).ConfigureAwait(false); var run = await db.AgentRun.FromSqlInterpolated($"SELECT agent_run.*, xmin FROM agent_run WHERE team_id = {claim.TeamId} AND id = {claim.AgentRunId} FOR UPDATE") @@ -278,6 +299,7 @@ private async Task SettleAsync(RecoveryClaim claim, Recovery && ((isManifest ? row.VerificationStalledAttempts : row.RecoveryAttemptCount) >= _options.RetryPolicy.MaxAttempts || now - (isManifest ? row.LastVerificationProgressAt ?? row.RecoveryStartedAt!.Value : row.RecoveryStartedAt!.Value) >= _options.RetryPolicy.MaxAge)) settled = RecoveryOutcome.Indeterminate(outcome.StreamId, "recovery-exhausted", $"Recovery exhausted its bounded attempts or age after '{outcome.ErrorCode}'."); + attempt.Outcome = settled; row.StreamId = settled.StreamId ?? row.StreamId; row.State = settled.State; @@ -338,6 +360,14 @@ private static bool ExactClaim(AgentRunLogStream value, RecoveryClaim expected) private static string Code(T value) where T : struct, Enum => string.Concat(value.ToString().Select((character, index) => char.IsUpper(character) && index > 0 ? $"-{char.ToLowerInvariant(character)}" : char.ToLowerInvariant(character).ToString())); private static AgentRunLogRecoveryClaimRef Fence(RecoveryClaim claim) => new(claim.Id, claim.RecoveryOwnerId, claim.RecoveryFenceEpoch); + private static PostgresException? DatabaseRefusal(Exception exception) + { + for (Exception? current = exception; current != null; current = current.InnerException) + if (current is PostgresException postgres) return postgres; + + return null; + } + private sealed record RecoveryClaim(Guid Id, Guid TeamId, Guid AgentRunId, long WorkerFenceEpoch, Guid CaptureSessionId, string StreamKind, string ContentType, string? ContentEncoding, string CaptureSource, Guid? StreamId, AgentRunLogCaptureIntentState State, DateTimeOffset NextRecoveryAt, DateTimeOffset? TerminalObservedAt, Guid RecoveryOwnerId, long RecoveryFenceEpoch) @@ -348,6 +378,12 @@ private sealed record RecoveryClaim(Guid Id, Guid TeamId, Guid AgentRunId, long private sealed record RecoverySettlement(AgentRunLogCaptureIntentState State, bool LostLease); + /// The outcome a settlement is writing: the observed one until the settlement supersedes or exhausts it. + private sealed class SettlementAttempt(RecoveryOutcome outcome) + { + public RecoveryOutcome Outcome { get; set; } = outcome; + } + private sealed record RecoveryOutcome(Guid? StreamId, AgentRunLogCaptureIntentState State, string? ErrorCode, string? ErrorMessage, RecoveryRetryDirective? RetryDirective) { public bool Terminal => State is AgentRunLogCaptureIntentState.Completed or AgentRunLogCaptureIntentState.CaptureFailed or AgentRunLogCaptureIntentState.Superseded or AgentRunLogCaptureIntentState.ExternalStateIndeterminate; diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs index a956d10ad..405fc58cc 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs @@ -1,3 +1,4 @@ +using System.Collections.Concurrent; using Autofac; using CodeSpace.Core.Persistence.Db; using CodeSpace.Core.Persistence.Entities; @@ -9,7 +10,9 @@ using CodeSpace.Messages.Enums; using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.Diagnostics; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; +using Npgsql; using Shouldly; namespace CodeSpace.IntegrationTests.Workflows; @@ -624,6 +627,96 @@ private static async Task SettleRetryAfterPauseAsync(CodeSpaceDbContext db, Guid await transaction.CommitAsync(); } + [Theory] + [InlineData(8, AgentRunLogCaptureIntentState.SourceFinalized, "complete-backend-unavailable")] // the observed retry is the write refused + [InlineData(1, AgentRunLogCaptureIntentState.ExternalStateIndeterminate, "recovery-exhausted")] // the settlement exhausts the observed retry, and that write is refused + public async Task A_settlement_the_database_refuses_is_logged_with_the_write_it_refused_and_still_waits_out_its_lease(int maxAttempts, AgentRunLogCaptureIntentState refusedState, string refusedCode) + { + // A settlement that raises leaves its claim leased and idle until the lease expires, and is counted as a lost + // lease. When the cause is a guard refusal, the code and the database contract disagree, and without the cause + // in the log that reads as an unreproducible flake. The guard's clauses branch on the state and code the write + // carries, and a settlement can replace the outcome it observed before it writes, so the log must name the write. + var scene = await SeedOwnedIntentBehindDueNeighbourAsync(); + var refusal = new RefusedCaptureSettlementInterceptor(scene.IntentId); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(new AlwaysRetryableCompleteLogService(scene.Logs), new RecoveryTestOptions { MaxAttempts = maxAttempts }, refusal, log); + + var summary = await ReconcileUntilFaultedAsync(recovery, () => refusal.Refused, scene.Owned); + + var entry = await ShouldWaitOutItsLeaseAloneAsync(scene, summary, log); + entry.Properties["Outcome"].ShouldBe(refusedState); + entry.Properties["OutcomeCode"].ShouldBe(refusedCode); + entry.Properties["SqlState"].ShouldBe(PostgresErrorCodes.RaiseException); + entry.Properties["MessageText"].ShouldBeOfType().ShouldContain($"(id={scene.IntentId})"); + entry.Exception.ShouldBeOfType().InnerException.ShouldBeOfType(); + } + + [Fact] + public async Task A_settlement_that_outlives_its_budget_is_logged_without_a_database_cause_and_still_waits_out_its_lease() + { + // The likeliest cause in production is not a Postgres error at all: the settlement outlives its budget waiting on + // a row lock or a slow read, and is cancelled. The exception is then the whole cause, SqlState and MessageText are + // null, and naming them must not throw from inside the catch, where it would fault the wave and stop its claims. + var scene = await SeedOwnedIntentBehindDueNeighbourAsync(); + var held = new SlowCaptureSettlementInterceptor(scene.Owned.AgentRunId, Timeout.InfiniteTimeSpan); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(new AlwaysRetryableCompleteLogService(scene.Logs), interceptor: held, logger: log); + + var summary = await ReconcileUntilFaultedAsync(recovery, () => held.Held, scene.Owned); + + var entry = await ShouldWaitOutItsLeaseAloneAsync(scene, summary, log); + entry.Properties["Outcome"].ShouldBe(AgentRunLogCaptureIntentState.SourceFinalized); + entry.Properties["OutcomeCode"].ShouldBe("complete-backend-unavailable"); + entry.Properties["SqlState"].ShouldBeNull(); + entry.Properties["MessageText"].ShouldBeNull(); + entry.Exception.ShouldBeAssignableTo(); + } + + /// + /// Seeds this test's terminal run with a finalized stream, so the next wave settles its intent, behind another + /// tenant's due neighbour that the same waves settle first through the same seam. + /// + private async Task SeedOwnedIntentBehindDueNeighbourAsync() + { + var neighbour = await SeedDueFinalizedStreamNeighbourAsync(); + var owned = await SeedWorldAsync(); + var logs = LogService(); + var sessionId = Guid.NewGuid(); + await Recovery(logs).DeclareAsync(Declaration(owned, sessionId, 7, AgentRunLogKinds.StandardOutput), CancellationToken.None); + await SeedFinalizedTerminalStreamAsync(owned, logs, sessionId); + + return new SettlementScene(neighbour, owned, (await IntentAsync(owned)).Id, logs); + } + + /// + /// Asserts that the faulted settlement kept today's outcome — counted as a lost lease, its claim left leased and + /// unwritten — and that the seam touched only this test's intent, then returns the one entry that names it. + /// + private async Task ShouldWaitOutItsLeaseAloneAsync(SettlementScene scene, AgentRunLogCaptureRecoverySummary summary, RecordedRecoveryLog log) + { + // This is the wave that reached the faulted intent, so it cannot have been crowded out; a stranger's lost lease + // could only satisfy the tally falsely, never fail it. + summary.LostLease.ShouldBeGreaterThanOrEqualTo(1, "a settlement that raises is still counted as a lost lease"); + + var intent = await IntentAsync(scene.Owned); + intent.RecoveryOwnerId.ShouldNotBeNull("the faulted write did not land, so the claim stays leased until its lease expires"); + intent.RecoveryAttemptCount.ShouldBe(1); + intent.State.ShouldBe(AgentRunLogCaptureIntentState.Expected); + intent.LastErrorCode.ShouldBeNull(); + + var stranger = await IntentAsync(scene.Neighbour); + stranger.RecoveryAttemptCount.ShouldBeGreaterThan(0, "the neighbour ahead of this intent was never claimed, so nothing showed the seam leaves it alone"); + stranger.RecoveryOwnerId.ShouldBeNull("the seam faulted a neighbour's settlement; it must fault only this test's intent"); + log.About(stranger.Id).ShouldBeEmpty(); + + var entry = log.About(scene.IntentId).ShouldHaveSingleItem("the faulted settlement must be logged once, naming its intent"); + entry.Level.ShouldBe(LogLevel.Warning); + entry.Properties["RunId"].ShouldBe(scene.Owned.AgentRunId); + entry.Properties["{OriginalFormat}"].ShouldBeOfType().ShouldContain("until its recovery lease expires"); + + return entry; + } + [Fact] public async Task Terminal_grace_uses_database_observation_time_not_positive_or_negative_application_clock_skew() { @@ -811,6 +904,30 @@ private async Task ReleaseAndReconcileUntilCompleted return seen; } + /// + /// Reconciles until a wave reaches THIS test's faulted settlement — is the seam's signal — + /// and returns that wave's summary. A wave is deployment-wide and bounded, so one start is not guaranteed to reach it. + /// + private static async Task ReconcileUntilFaultedAsync(AgentRunLogCaptureRecoveryService recovery, Func faulted, World world) + { + var deadline = DateTimeOffset.UtcNow + TimeSpan.FromSeconds(10); + faulted().ShouldBeFalse("the seam fired before any wave was started to reach this test's own settlement"); + + while (DateTimeOffset.UtcNow < deadline) + { + var summary = await recovery.ReconcileAsync(CancellationToken.None); + + if (faulted()) return summary; + + await Task.Delay(25); + } + + throw new Xunit.Sdk.XunitException( + $"No reconcile wave reached the settlement of agent run {world.AgentRunId}'s capture intent, so its seam never faulted it. " + + "Reconcile waves are deployment-wide and bounded, so check whether earlier tests left enough due intents to crowd this one out, " + + "then whether the seam still recognises the settlement's SQL."); + } + private async Task IntentAsync(World world) { using var scope = _fixture.BeginScope(); @@ -819,7 +936,7 @@ private async Task IntentAsync(World world) .SingleAsync(value => value.AgentRunId == world.AgentRunId); } - private AgentRunLogCaptureRecoveryService Recovery(IAgentRunLogService logs, RecoveryTestOptions? options = null, IInterceptor? interceptor = null) + private AgentRunLogCaptureRecoveryService Recovery(IAgentRunLogService logs, RecoveryTestOptions? options = null, IInterceptor? interceptor = null, ILogger? logger = null) { options ??= new RecoveryTestOptions(); using var scope = _fixture.BeginScope(); @@ -828,7 +945,7 @@ private AgentRunLogCaptureRecoveryService Recovery(IAgentRunLogService logs, Rec return new AgentRunLogCaptureRecoveryService(dbOptions, logs, new AgentRunLogCaptureRecoveryOptions(20, options.MaxConcurrency, TimeSpan.FromSeconds(6), options.OperationTimeout, - new AgentRunLogCaptureRetryPolicy(options.BaseDelay, options.MaxDelay, options.MaxAttempts, options.MaxAge, options.TerminalGrace))); + new AgentRunLogCaptureRetryPolicy(options.BaseDelay, options.MaxDelay, options.MaxAttempts, options.MaxAge, options.TerminalGrace)), logger); } private sealed record RecoveryTestOptions @@ -948,6 +1065,32 @@ await logs.FinalizeSourceAsync(new AgentRunLogFinalizeSourceRequest private sealed record World(Guid TeamId, Guid ActorId, Guid AgentRunId); + private sealed record SettlementScene(World Neighbour, World Owned, Guid IntentId, IAgentRunLogService Logs); + + /// + /// The recovery's log entries, kept as their structured properties rather than a rendered string, so what an operator + /// is handed — which intent, which run, which database refusal — is the thing under test. + /// + private sealed class RecordedRecoveryLog : ILogger + { + private readonly ConcurrentBag _entries = []; + + /// The entries naming one intent. The sweep is deployment-wide, so whatever else it met in this shared database is not this test's business. + public IReadOnlyList About(Guid intentId) => _entries.Where(entry => Equals(entry.Properties.GetValueOrDefault("IntentId"), intentId)).ToList(); + + public IDisposable? BeginScope(TState state) where TState : notnull => null; + public bool IsEnabled(LogLevel logLevel) => true; + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) + { + if (state is not IReadOnlyList> properties) return; + + _entries.Add(new RecordedEntry(logLevel, properties.ToDictionary(property => property.Key, property => property.Value, StringComparer.Ordinal), exception)); + } + } + + private sealed record RecordedEntry(LogLevel Level, IReadOnlyDictionary Properties, Exception? Exception); + private sealed class EmptyCas : IArtifactCasRuntimeCoordinator { public Task PutAsync(ArtifactCasTransferRequest request, CancellationToken cancellationToken) => Task.FromResult(new ArtifactCasTransferResult.Rejected(null, new ArtifactCasProblem(ArtifactCasProblemCode.Unsupported, false))); diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs new file mode 100644 index 000000000..a549f458d --- /dev/null +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs @@ -0,0 +1,52 @@ +using System.Data.Common; +using System.Text.RegularExpressions; +using Microsoft.EntityFrameworkCore.Diagnostics; + +namespace CodeSpace.IntegrationTests.Workflows.Infrastructure; + +/// +/// Makes agent_run_log_capture_intent_guard() refuse every capture-recovery settlement of +/// with its own P0001, the refusal the service meets whenever its code and the database contract disagree. The write +/// keeps every column the settlement computed except its revision, which skips one, so the guard refuses it whatever +/// outcome it carries, a retry or a terminal one. A settlement releases its claim and a claim takes one, so only a write +/// that assigns recovery_owner_id NULL is touched. Only the recovery's own options carry it, and only that intent's write +/// is touched, so a neighbour's intent in the same wave settles. +/// +internal sealed partial class RefusedCaptureSettlementInterceptor(Guid intentId) : DbCommandInterceptor +{ + private int _refused; + + public bool Refused => Volatile.Read(ref _refused) == 1; + + public override ValueTask> ReaderExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) + { + if (OwnSettlementRevision(command) is { } revision) + { + revision.Value = (long)revision.Value! + 1; + Volatile.Write(ref _refused, 1); + } + + return ValueTask.FromResult(result); + } + + /// The revision parameter of this intent's settlement write, or null for any other command. + private DbParameter? OwnSettlementRevision(DbCommand command) + { + if (!IntentUpdate().IsMatch(command.CommandText) || !command.Parameters.Cast().Any(parameter => parameter.Value is Guid id && id == intentId)) return null; + + return Assigned(command, "recovery_owner_id") is { Value: DBNull } ? Assigned(command, "revision") : null; + } + + private static DbParameter? Assigned(DbCommand command, string column) + { + var assignment = Assignment().Matches(command.CommandText).FirstOrDefault(match => match.Groups["column"].Value == column); + + return assignment == null ? null : command.Parameters.Cast().Single(parameter => parameter.ParameterName.TrimStart('@') == assignment.Groups["name"].Value); + } + + [GeneratedRegex(@"UPDATE ""?agent_run_log_capture_intent""? SET ", RegexOptions.CultureInvariant)] + private static partial Regex IntentUpdate(); + + [GeneratedRegex(@"""?(?\w+)""? = @(?\w+)", RegexOptions.CultureInvariant)] + private static partial Regex Assignment(); +}