diff --git a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs index 6af8c9c50..b75959538 100644 --- a/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs +++ b/backend/src/CodeSpace.Core/Services/Agents/AgentRunLogging/AgentRunLogCaptureRecoveryService.cs @@ -95,7 +95,20 @@ public async Task ReconcileAsync(Cancellation else counts.Record(settlement.State); } } - return counts.Summary(); + var summary = counts.Summary(); + + LogReconciledWave(summary); + + return summary; + } + + // The recurring job discards the summary, so this line is the only reader of the wave's tally, lost leases included. + // A wave that claimed nothing has nothing to report, and every idle worker would otherwise log it each minute. + private void LogReconciledWave(AgentRunLogCaptureRecoverySummary summary) + { + if (summary.Claimed == 0) return; + + _logger.LogInformation("Agent run log capture recovery claimed {Claimed} intent(s): {Completed} completed, {CaptureFailed} capture-failed, {Superseded} superseded, {ExternalStateIndeterminate} indeterminate, {Retried} retried, {LostLease} lost their lease or could not settle", summary.Claimed, summary.Completed, summary.CaptureFailed, summary.Superseded, summary.ExternalStateIndeterminate, summary.Retried, summary.LostLease); } private async Task RecoverClaimAsync(RecoveryClaim claim, CancellationToken cancellationToken) @@ -109,9 +122,10 @@ private async Task RecoverClaimAsync(RecoveryClaim claim, Ca { outcome = RecoveryOutcome.Retry(claim.StreamId, claim.State, "recovery-operation-timeout", "The bounded log recovery operation timed out."); } - catch (Exception) + catch (Exception exception) { outcome = RecoveryOutcome.Retry(claim.StreamId, claim.State, "recovery-operation-exception", "The log recovery operation raised an unexpected error."); + LogFailedRecovery(claim, outcome, exception); } using var settlement = new CancellationTokenSource(_options.OperationTimeout); @@ -127,6 +141,17 @@ private async Task RecoverClaimAsync(RecoveryClaim claim, Ca } } + // An unexpected error from the recovery step becomes a typed retry whose last_error_code names no cause, so a + // deterministic one repeats on every retry until the intent is exhausted, and its cause is visible only here. The + // line is written before the settlement runs, so it names the step's outcome and promises nothing the settlement + // decides: the settlement may supersede or exhaust the intent instead, or write nothing at all. + private void LogFailedRecovery(RecoveryClaim claim, RecoveryOutcome outcome, Exception exception) + { + var refusal = DatabaseRefusal(exception); + + _logger.LogWarning(exception, "Agent run {RunId} log capture intent {IntentId} could not be recovered: the recovery step raised an unexpected error (SQLSTATE {SqlState}: {MessageText}); the step's outcome is a typed retry as {Outcome} with last_error_code {OutcomeCode}, which the settlement writes unless the settlement supersedes or exhausts the intent, or cannot settle it", claim.AgentRunId, claim.Id, refusal?.SqlState, refusal?.MessageText, outcome.State, outcome.ErrorCode); + } + // 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. @@ -273,7 +298,10 @@ private async Task SettleAsync(RecoveryClaim claim, Settleme .SingleOrDefaultAsync(cancellationToken).ConfigureAwait(false); var now = await DatabaseClockAsync(db, cancellationToken).ConfigureAwait(false); if (run == null || row == null || row.RecoveryOwnerId != claim.RecoveryOwnerId || row.RecoveryFenceEpoch != claim.RecoveryFenceEpoch || row.RecoveryLeaseExpiresAt <= now) + { + LogLostClaim(claim, outcome, LostClaimCause(claim, run, row), null); return new RecoverySettlement(outcome.State, true); + } var manifestStreamId = await db.AgentRunLogStream.Where(value => value.TeamId == row.TeamId && value.AgentRunId == row.AgentRunId && value.StreamKind == row.StreamKind && value.SchemaVersion == 3) .Select(value => (Guid?)value.Id).SingleOrDefaultAsync(cancellationToken).ConfigureAwait(false); @@ -318,7 +346,28 @@ private async Task SettleAsync(RecoveryClaim claim, Settleme await transaction.CommitAsync(cancellationToken).ConfigureAwait(false); return new RecoverySettlement(settled.State, false); } - catch (DbUpdateConcurrencyException) { return new RecoverySettlement(settled.State, true); } + catch (DbUpdateConcurrencyException exception) + { + LogLostClaim(claim, settled, "row-version-changed", exception); + return new RecoverySettlement(settled.State, true); + } + } + + /// Why a settlement whose claim no longer holds lost it; call it only once the lease re-check has failed. + private static string LostClaimCause(RecoveryClaim claim, AgentRun? run, AgentRunLogCaptureIntent? row) + { + if (run == null || row == null) return "claim-row-missing"; + + return row.RecoveryOwnerId != claim.RecoveryOwnerId || row.RecoveryFenceEpoch != claim.RecoveryFenceEpoch ? "reclaimed" : "lease-expired"; + } + + // A lease outlives both bounded steps of a claim, so a live worker loses one only when a step overran the bound its + // cancellation set, the process stalled, or the row changed under the settlement's own lock. The fence keeps that + // safe, but the discarded settlement is redone by a later claim, which spends another recovery attempt, and only + // this line says which claim it was and why. + private void LogLostClaim(RecoveryClaim claim, RecoveryOutcome outcome, string cause, Exception? exception) + { + _logger.LogWarning(exception, "Agent run {RunId} log capture intent {IntentId} lost its recovery claim ({Cause}) before it could settle as {Outcome} with last_error_code {OutcomeCode}; the settlement is discarded, and the claim's current holder, or the wave that re-claims it once its lease expires, settles the intent instead", claim.AgentRunId, claim.Id, cause, outcome.State, outcome.ErrorCode); } private CodeSpaceDbContext CreateDb() => new(_dbOptions); diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs index 405fc58cc..389ebf3f8 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/AgentRunLogCaptureRecoveryFlowTests.cs @@ -698,25 +698,237 @@ private async Task ShouldWaitOutItsLeaseAloneAsync(SettlementScen // 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); + await ShouldBeLeftLeasedAndUnwrittenAsync(scene.Owned); + await ShouldHaveLeftTheNeighbourAloneAsync(scene, log); + + 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; + } + + private async Task ShouldBeLeftLeasedAndUnwrittenAsync(World owned) + { + var intent = await IntentAsync(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(); + } + /// The neighbour ahead of this test's intent went through the same seam: it was claimed, settled, and never logged. + private async Task ShouldHaveLeftTheNeighbourAloneAsync(SettlementScene scene, RecordedRecoveryLog log) + { 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"); + [Fact] + public async Task A_recovery_step_the_database_refuses_is_logged_with_its_cause_and_still_settles_as_a_typed_retry() + { + // An error the recovery step raises becomes a retry whose last_error_code, recovery-operation-exception, names no + // cause. A deterministic refusal then repeats on every retry until the intent is exhausted into + // ExternalStateIndeterminate, and without the cause in the log nothing says why. + var scene = await SeedOwnedIntentBehindDueNeighbourAsync(); + var refusal = new RefusedCaptureRecoveryReadInterceptor(scene.Owned.TeamId); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(scene.Logs, interceptor: refusal, logger: log); + + var entry = await ShouldSettleTheStepsRetryAloneAsync(scene, recovery, log, AgentRunLogCaptureIntentState.Expected, "recovery-operation-exception"); + + refusal.Refused.ShouldBeTrue("the seam never refused this intent's recovery read, so nothing was raised to log"); + entry.Properties["SqlState"].ShouldBe(PostgresErrorCodes.RaiseException); + entry.Properties["MessageText"].ShouldBeOfType().ShouldContain(scene.Owned.TeamId.ToString()); + entry.Exception.ShouldBeOfType(); + } + + [Theory] + [InlineData(8, AgentRunLogCaptureIntentState.Expected, "recovery-operation-exception")] // the settlement writes the step's retry + [InlineData(1, AgentRunLogCaptureIntentState.ExternalStateIndeterminate, "recovery-exhausted")] // the settlement exhausts the step's retry on its last attempt + public async Task A_recovery_step_that_fails_outside_the_database_is_logged_without_a_database_cause_and_its_settlement_decides_the_retry(int maxAttempts, AgentRunLogCaptureIntentState settledState, string settledCode) + { + // A provider fault inside CompleteAsync carries no SQLSTATE. The exception is then the whole cause, and naming the + // absent database fields must not throw from inside the catch, where it would fault the wave and stop its claims. + // A deterministic fault repeats until the settlement exhausts the intent on its last attempt, and the line is + // written before that settlement runs, so it must not promise that attempt another retry. + var scene = await SeedOwnedIntentBehindDueNeighbourAsync(); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(new ThrowingCompleteLogService(scene.Logs, scene.Owned.AgentRunId), new RecoveryTestOptions { MaxAttempts = maxAttempts }, logger: log); + + var entry = await ShouldSettleTheStepsRetryAloneAsync(scene, recovery, log, settledState, settledCode); + + entry.Properties["SqlState"].ShouldBeNull(); + entry.Properties["MessageText"].ShouldBeNull(); + entry.Exception.ShouldBeOfType().Message.ShouldContain(scene.Owned.AgentRunId.ToString()); + } + + /// + /// Reconciles until THIS test's intent settles the typed retry an unexpected recovery error becomes, as that retry or as + /// what its settlement replaced it with, asserts that the settlement released its claim as it did before the error was + /// logged and that the seam left the neighbour alone, then returns the one entry that names the intent. + /// + private async Task ShouldSettleTheStepsRetryAloneAsync(SettlementScene scene, AgentRunLogCaptureRecoveryService recovery, RecordedRecoveryLog log, AgentRunLogCaptureIntentState settledState, string settledCode) + { + var intent = await ReconcileUntilAsync(recovery, scene.Owned, value => value.LastErrorCode == settledCode, $"settled the typed retry an unexpected recovery error becomes as {settledState}/{settledCode}"); + + intent.State.ShouldBe(settledState); + intent.RecoveryOwnerId.ShouldBeNull("a settlement releases its claim"); + intent.RecoveryAttemptCount.ShouldBe(1); + await ShouldHaveLeftTheNeighbourAloneAsync(scene, log); + + var entry = log.About(scene.IntentId).ShouldHaveSingleItem("the unexpected recovery error 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"); + + // The line is written before the settlement runs, so it names the step's outcome and leaves the rest to the settlement. + entry.Properties["Outcome"].ShouldBe(AgentRunLogCaptureIntentState.Expected); + entry.Properties["OutcomeCode"].ShouldBe("recovery-operation-exception"); + entry.Properties["{OriginalFormat}"].ShouldBeOfType().ShouldContain("unless the settlement supersedes or exhausts the intent"); return entry; } + [Theory] + [InlineData(false, "lease-expired")] // nothing re-claimed the intent: the claim is still this wave's, but its lease is gone + [InlineData(true, "reclaimed")] // another worker re-claimed the intent once the lease expired, and settled it first + public async Task A_settlement_that_outlived_its_lease_is_logged_with_why_and_counted_in_the_summary_its_wave_logs(bool reclaimedByAnotherWorker, string cause) + { + // A lease outlives both bounded steps of a claim, so a live worker loses one only when a step overruns the bound + // its cancellation sets — here a provider call that ignores cancellation. The fence discards the late settlement, + // which is right, but it was counted as a lost lease with no log line, and nothing read that tally either. + var scene = await SeedOwnedIntentBehindDueNeighbourAsync(); + var stall = new StallFirstOwnCompleteLogService(scene.Logs, scene.Owned.AgentRunId); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(stall, logger: log); + var wave = await ReconcileUntilPausedAsync(recovery, stall.Entered, scene.Owned); + + if (reclaimedByAnotherWorker) + await ReconcileUntilAsync(Recovery(scene.Logs), scene.Owned, value => value.State == AgentRunLogCaptureIntentState.Completed, "re-claimed and completed by another worker once the stalled claim's lease expired", LeaseWait); + else + await WaitForLeaseToExpireAsync(scene.Owned); + + stall.Release(); + var summary = await wave; + + summary.LostLease.ShouldBeGreaterThanOrEqualTo(1, "the late settlement is discarded and counted as a lost lease"); + log.Summaries.Last().ShouldBe(summary, "the wave must log the tally it returns; nothing else reads its lost leases"); + await ShouldHaveLeftTheNeighbourAloneAsync(scene, log); + + var entry = log.About(scene.IntentId).ShouldHaveSingleItem("the lost claim must be logged once, naming its intent"); + entry.Level.ShouldBe(LogLevel.Warning); + entry.Properties["RunId"].ShouldBe(scene.Owned.AgentRunId); + entry.Properties["Cause"].ShouldBe(cause); + entry.Properties["Outcome"].ShouldBe(AgentRunLogCaptureIntentState.SourceFinalized); + entry.Properties["OutcomeCode"].ShouldBe("complete-backend-unavailable"); + entry.Exception.ShouldBeNull(); + } + + [Theory] + [InlineData(8, AgentRunLogCaptureIntentState.SourceFinalized, "complete-backend-unavailable")] // the observed retry is the write discarded + [InlineData(1, AgentRunLogCaptureIntentState.ExternalStateIndeterminate, "recovery-exhausted")] // the settlement exhausts the observed retry, and that write is discarded + public async Task A_settlement_whose_row_changed_under_its_lock_is_logged_with_the_write_it_discarded_and_still_waits_out_its_lease(int maxAttempts, AgentRunLogCaptureIntentState discardedState, string discardedCode) + { + // The settlement locks the intent before it reads it, so its write matches no row only if the row changed under + // that lock — a contract the service relies on and nothing else checks. The write is counted as a lost lease, and + // its claim stays leased until the lease expires, with no log line. A settlement can replace the outcome it + // observed before it writes, so the line must name the write it discarded, not the outcome it observed. + var scene = await SeedOwnedIntentBehindDueNeighbourAsync(); + var stale = new RefusedCaptureSettlementInterceptor(scene.IntentId, "xmin"); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(new AlwaysRetryableCompleteLogService(scene.Logs), new RecoveryTestOptions { MaxAttempts = maxAttempts }, stale, log); + + var summary = await ReconcileUntilFaultedAsync(recovery, () => stale.Refused, scene.Owned); + + summary.LostLease.ShouldBeGreaterThanOrEqualTo(1, "a write that matches no row is counted as a lost lease"); + log.Summaries.Last().ShouldBe(summary, "the wave must log the tally it returns; nothing else reads its lost leases"); + await ShouldBeLeftLeasedAndUnwrittenAsync(scene.Owned); + await ShouldHaveLeftTheNeighbourAloneAsync(scene, log); + + var entry = log.About(scene.IntentId).ShouldHaveSingleItem("the lost claim must be logged once, naming its intent"); + entry.Level.ShouldBe(LogLevel.Warning); + entry.Properties["RunId"].ShouldBe(scene.Owned.AgentRunId); + entry.Properties["Cause"].ShouldBe("row-version-changed"); + entry.Properties["Outcome"].ShouldBe(discardedState); + entry.Properties["OutcomeCode"].ShouldBe(discardedCode); + entry.Exception.ShouldBeOfType(); + } + + [Fact] + public async Task A_wave_logs_one_summary_when_it_claimed_work_and_none_when_it_claimed_nothing() + { + // The recurring job discards the summary, so the wave's own log line is the only place its tally is read. A wave + // that claimed nothing has nothing to report, and one line a minute from every idle worker would bury the rest. + 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); + var log = new RecordedRecoveryLog(); + var recovery = Recovery(logs, logger: log); + + var claimingWaves = await ReconcileUntilIdleAsync(recovery, owned); + + claimingWaves.ShouldBeGreaterThanOrEqualTo(1, "this test's own due intent was never claimed, so no wave had work to report"); + log.Summaries.Count.ShouldBe(claimingWaves, "every wave that claimed work logs its summary once, and the idle wave logs none"); + } + + /// + /// Reconciles until a wave claims nothing, after THIS test's intent has been claimed, and returns how many waves claimed + /// work first. Waves are deployment-wide, so earlier tests' due intents are claimed too, and each settles terminal or + /// schedules its retry past the next wave's cutoff, so a wave with nothing due follows. + /// + private async Task ReconcileUntilIdleAsync(AgentRunLogCaptureRecoveryService recovery, World world) + { + var deadline = DateTimeOffset.UtcNow + TimeSpan.FromSeconds(10); + var claimingWaves = 0; + (await IntentAsync(world)).RecoveryAttemptCount.ShouldBe(0, "the intent was claimed before any wave was started to claim it"); + + while (DateTimeOffset.UtcNow < deadline) + { + var summary = await recovery.ReconcileAsync(CancellationToken.None); + + if (summary.Claimed == 0 && (await IntentAsync(world)).RecoveryAttemptCount > 0) return claimingWaves; + + if (summary.Claimed > 0) claimingWaves++; + } + + throw new Xunit.Sdk.XunitException( + $"No reconcile wave claimed nothing after agent run {world.AgentRunId}'s capture intent was claimed, across {claimingWaves} claiming wave(s). " + + "Reconcile waves are deployment-wide, so check whether an earlier test left intents that fall due again before every next wave's cutoff."); + } + + /// How long a test waits on a 6 s recovery lease to expire, with room for a loaded database. + private static readonly TimeSpan LeaseWait = TimeSpan.FromSeconds(20); + + /// Polls the database clock until THIS test's claimed intent's recovery lease has expired. + private async Task WaitForLeaseToExpireAsync(World world) + { + var deadline = DateTimeOffset.UtcNow + LeaseWait; + (await LeaseExpiredAsync(world)).ShouldBeFalse("waiting for the lease to expire is meaningless when it had expired before the wait began"); + + while (!await LeaseExpiredAsync(world)) + { + if (DateTimeOffset.UtcNow >= deadline) + throw new Xunit.Sdk.XunitException($"The recovery lease on agent run {world.AgentRunId}'s capture intent never expired by the database clock; check its recovery_lease_expires_at against clock_timestamp()."); + + await Task.Delay(50); + } + } + + private async Task LeaseExpiredAsync(World world) + { + using var scope = _fixture.BeginScope(); + + return await scope.Resolve().Database + .SqlQuery($"SELECT COALESCE(recovery_lease_expires_at <= clock_timestamp(), FALSE) AS \"Value\" FROM agent_run_log_capture_intent WHERE agent_run_id = {world.AgentRunId}").SingleAsync(); + } + [Fact] public async Task Terminal_grace_uses_database_observation_time_not_positive_or_negative_application_clock_skew() { @@ -1073,11 +1285,21 @@ private sealed record SettlementScene(World Neighbour, World Owned, Guid IntentI /// private sealed class RecordedRecoveryLog : ILogger { - private readonly ConcurrentBag _entries = []; + private readonly ConcurrentQueue _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(); + /// The wave summaries, in the order the waves logged them, rebuilt from their structured properties. + public IReadOnlyList Summaries => _entries.Where(entry => entry.Properties.ContainsKey("Claimed")).Select(Summary).ToList(); + + private static AgentRunLogCaptureRecoverySummary Summary(RecordedEntry entry) => new(Count(entry, "Claimed"), Count(entry, "Completed"), Count(entry, "CaptureFailed"), Count(entry, "Superseded"), Count(entry, "Retried"), Count(entry, "LostLease")) + { + ExternalStateIndeterminate = Count(entry, "ExternalStateIndeterminate"), + }; + + private static int Count(RecordedEntry entry, string name) => entry.Properties[name].ShouldBeOfType(); + public IDisposable? BeginScope(TState state) where TState : notnull => null; public bool IsEnabled(LogLevel logLevel) => true; @@ -1085,7 +1307,7 @@ public void Log(LogLevel logLevel, EventId eventId, TState state, Except { if (state is not IReadOnlyList> properties) return; - _entries.Add(new RecordedEntry(logLevel, properties.ToDictionary(property => property.Key, property => property.Value, StringComparer.Ordinal), exception)); + _entries.Enqueue(new RecordedEntry(logLevel, properties.ToDictionary(property => property.Key, property => property.Value, StringComparer.Ordinal), exception)); } } @@ -1135,6 +1357,59 @@ public Task CompleteAsync(AgentRunLogCompleteRequest public Task ReadRangeAsync(AgentRunLogRangeRequest request, CancellationToken cancellationToken) => inner.ReadRangeAsync(request, cancellationToken); } + /// + /// Fails the owning run's every CompleteAsync outside the database, the way a provider or CAS fault would. Only that + /// run's call fails: a wave is deployment-wide, so a neighbour's call runs through this same instance and completes. + /// + private sealed class ThrowingCompleteLogService(IAgentRunLogService inner, Guid agentRunId) : IAgentRunLogService + { + public Task OpenAsync(AgentRunLogOpenRequest request, CancellationToken cancellationToken) => inner.OpenAsync(request, cancellationToken); + public Task AppendAsync(AgentRunLogAppendRequest request, CancellationToken cancellationToken) => inner.AppendAsync(request, cancellationToken); + public Task FinalizeSourceAsync(AgentRunLogFinalizeSourceRequest request, CancellationToken cancellationToken) => inner.FinalizeSourceAsync(request, cancellationToken); + public Task CompleteAsync(AgentRunLogCompleteRequest request, CancellationToken cancellationToken) => request.AgentRunId == agentRunId + ? Task.FromException(new InvalidOperationException($"The completion provider failed for agent run {agentRunId}.")) + : inner.CompleteAsync(request, cancellationToken); + public Task FailCaptureAsync(AgentRunLogFailCaptureRequest request, CancellationToken cancellationToken) => inner.FailCaptureAsync(request, cancellationToken); + public Task RecordOwnerLossAsync(AgentRunLogOwnerLossRequest request, CancellationToken cancellationToken) => inner.RecordOwnerLossAsync(request, cancellationToken); + public Task GetMetadataAsync(Guid teamId, Guid streamId, CancellationToken cancellationToken) => inner.GetMetadataAsync(teamId, streamId, cancellationToken); + public Task> ListMetadataAsync(Guid teamId, Guid agentRunId, CancellationToken cancellationToken) => inner.ListMetadataAsync(teamId, agentRunId, cancellationToken); + public Task> ListCaptureHeadsAsync(Guid teamId, Guid agentRunId, CancellationToken cancellationToken) => inner.ListCaptureHeadsAsync(teamId, agentRunId, cancellationToken); + public Task ReadRangeAsync(AgentRunLogRangeRequest request, CancellationToken cancellationToken) => inner.ReadRangeAsync(request, cancellationToken); + } + + /// + /// Holds the owning run's first CompleteAsync until released and IGNORES its cancellation, as a provider call that does + /// not honour its bound would, then answers it as retryable. That is how a live worker overruns its lease. Only the + /// owning run's first call is held: a wave is deployment-wide, so a neighbour's call runs through this same instance, + /// and holding a stranger's would stall the wave before this test's intent is claimed. Every later call passes through. + /// + private sealed class StallFirstOwnCompleteLogService(IAgentRunLogService inner, Guid agentRunId) : IAgentRunLogService + { + private readonly TaskCompletionSource _entered = new(TaskCreationOptions.RunContinuationsAsynchronously); + private readonly TaskCompletionSource _release = new(TaskCreationOptions.RunContinuationsAsynchronously); + private int _stalled; + + public Task Entered => _entered.Task; + public void Release() => _release.TrySetResult(); + public Task OpenAsync(AgentRunLogOpenRequest request, CancellationToken cancellationToken) => inner.OpenAsync(request, cancellationToken); + public Task AppendAsync(AgentRunLogAppendRequest request, CancellationToken cancellationToken) => inner.AppendAsync(request, cancellationToken); + public Task FinalizeSourceAsync(AgentRunLogFinalizeSourceRequest request, CancellationToken cancellationToken) => inner.FinalizeSourceAsync(request, cancellationToken); + public async Task CompleteAsync(AgentRunLogCompleteRequest request, CancellationToken cancellationToken) + { + if (request.AgentRunId != agentRunId || Interlocked.Exchange(ref _stalled, 1) == 1) return await inner.CompleteAsync(request, cancellationToken); + + _entered.TrySetResult(); + await _release.Task; + return new AgentRunLogCompleteResult.Rejected(new AgentRunLogProblem(AgentRunLogProblemCode.BackendUnavailable, true)); + } + public Task FailCaptureAsync(AgentRunLogFailCaptureRequest request, CancellationToken cancellationToken) => inner.FailCaptureAsync(request, cancellationToken); + public Task RecordOwnerLossAsync(AgentRunLogOwnerLossRequest request, CancellationToken cancellationToken) => inner.RecordOwnerLossAsync(request, cancellationToken); + public Task GetMetadataAsync(Guid teamId, Guid streamId, CancellationToken cancellationToken) => inner.GetMetadataAsync(teamId, streamId, cancellationToken); + public Task> ListMetadataAsync(Guid teamId, Guid agentRunId, CancellationToken cancellationToken) => inner.ListMetadataAsync(teamId, agentRunId, cancellationToken); + public Task> ListCaptureHeadsAsync(Guid teamId, Guid agentRunId, CancellationToken cancellationToken) => inner.ListCaptureHeadsAsync(teamId, agentRunId, cancellationToken); + public Task ReadRangeAsync(AgentRunLogRangeRequest request, CancellationToken cancellationToken) => inner.ReadRangeAsync(request, cancellationToken); + } + private sealed class BlockingCompleteLogService(IAgentRunLogService inner) : IAgentRunLogService { public Task OpenAsync(AgentRunLogOpenRequest request, CancellationToken cancellationToken) => inner.OpenAsync(request, cancellationToken); diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureRecoveryReadInterceptor.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureRecoveryReadInterceptor.cs new file mode 100644 index 000000000..96c7de3aa --- /dev/null +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureRecoveryReadInterceptor.cs @@ -0,0 +1,32 @@ +using System.Data.Common; +using Microsoft.EntityFrameworkCore.Diagnostics; + +namespace CodeSpace.IntegrationTests.Workflows.Infrastructure; + +/// +/// Makes the database refuse every capture-session read the capture recovery step makes for . The +/// server raises the P0001, so the service meets the PostgresException a failing read would hand it, not one built here. +/// Only the recovery step reads agent_run_log_capture_session through the recovery's own options, so the claim +/// and the settlement around it are untouched. Every test run has its own team and only that team's read is refused, so a +/// neighbour's intent in the same wave recovers. +/// +internal sealed class RefusedCaptureRecoveryReadInterceptor(Guid teamId) : 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 (!ReadsOwnCaptureSession(command)) return ValueTask.FromResult(result); + + command.Parameters.Clear(); + command.CommandText = $"DO $$ BEGIN RAISE EXCEPTION 'capture-session read refused for team {teamId}'; END $$"; + Volatile.Write(ref _refused, 1); + + return ValueTask.FromResult(result); + } + + private bool ReadsOwnCaptureSession(DbCommand command) => command.CommandText.Contains("FROM agent_run_log_capture_session", StringComparison.Ordinal) + && command.Parameters.Cast().Any(parameter => parameter.Value is Guid id && id == teamId); +} diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs index a549f458d..1e47e5d1d 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/Infrastructure/RefusedCaptureSettlementInterceptor.cs @@ -5,14 +5,16 @@ 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. +/// Makes every capture-recovery settlement of fail at the database by skipping one value of the +/// its write binds. The write keeps every other column the settlement computed, so it fails +/// whatever outcome it carries, a retry or a terminal one. With revision, the value it assigns, +/// agent_run_log_capture_intent_guard() refuses it with its own P0001, the refusal the service meets whenever its +/// code and the database contract disagree. With xmin, the row version its WHERE clause expects, the write matches +/// no row, which is how a row that changed under the settlement's lock would look. 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 +internal sealed partial class RefusedCaptureSettlementInterceptor(Guid intentId, string column = "revision") : DbCommandInterceptor { private int _refused; @@ -20,23 +22,31 @@ internal sealed partial class RefusedCaptureSettlementInterceptor(Guid intentId) public override ValueTask> ReaderExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) { - if (OwnSettlementRevision(command) is { } revision) + if (OwnSettlementValue(command) is { } forged) { - revision.Value = (long)revision.Value! + 1; + forged.Value = Skipped(forged.Value); 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) + /// The parameter this intent's settlement write binds to the forged column, or null for any other command. + private DbParameter? OwnSettlementValue(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; + return Assigned(command, "recovery_owner_id") is { Value: DBNull } ? Assigned(command, column) : null; } + // Each arm keeps its own type: a switch whose arms widen to one numeric type would bind the row version as a long. + private static object Skipped(object? value) => value switch + { + long revision => (object)(revision + 1), + uint xmin => (object)(xmin + 1u), + _ => throw new InvalidOperationException($"The settlement bound {value?.GetType().Name ?? "null"} where a revision or row version was expected."), + }; + private static DbParameter? Assigned(DbCommand command, string column) { var assignment = Assignment().Matches(command.CommandText).FirstOrDefault(match => match.Groups["column"].Value == column);