Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -18,10 +21,11 @@ public sealed partial class AgentRunLogCaptureRecoveryService : IAgentRunLogCapt
private readonly DbContextOptions<CodeSpaceDbContext> _dbOptions;
private readonly IAgentRunLogService _logs;
private readonly AgentRunLogCaptureRecoveryOptions _options;
private readonly ILogger<AgentRunLogCaptureRecoveryService> _logger;

public AgentRunLogCaptureRecoveryService(DbContextOptions<CodeSpaceDbContext> dbOptions, IAgentRunLogService logs) : this(dbOptions, logs, Defaults) { }
public AgentRunLogCaptureRecoveryService(DbContextOptions<CodeSpaceDbContext> dbOptions, IAgentRunLogService logs, ILogger<AgentRunLogCaptureRecoveryService> logger) : this(dbOptions, logs, Defaults, logger) { }

internal AgentRunLogCaptureRecoveryService(DbContextOptions<CodeSpaceDbContext> dbOptions, IAgentRunLogService logs, AgentRunLogCaptureRecoveryOptions options)
internal AgentRunLogCaptureRecoveryService(DbContextOptions<CodeSpaceDbContext> dbOptions, IAgentRunLogService logs, AgentRunLogCaptureRecoveryOptions options, ILogger<AgentRunLogCaptureRecoveryService>? 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
Expand All @@ -31,6 +35,7 @@ internal AgentRunLogCaptureRecoveryService(DbContextOptions<CodeSpaceDbContext>
_dbOptions = dbOptions;
_logs = logs;
_options = options;
_logger = logger ?? NullLogger<AgentRunLogCaptureRecoveryService>.Instance;
}

public async Task<AgentRunLogCaptureDeclarationResult> DeclareAsync(AgentRunLogCaptureDeclarationRequest request, CancellationToken cancellationToken)
Expand Down Expand Up @@ -110,11 +115,26 @@ private async Task<RecoverySettlement> 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<IReadOnlyList<RecoveryClaim>> ClaimBatchAsync(Guid ownerId, DateTimeOffset cutoff, int limit, CancellationToken cancellationToken)
Expand Down Expand Up @@ -242,8 +262,9 @@ private async Task<RecoveryOutcome> 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<RecoverySettlement> SettleAsync(RecoveryClaim claim, RecoveryOutcome outcome, CancellationToken cancellationToken)
private async Task<RecoverySettlement> 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")
Expand Down Expand Up @@ -278,6 +299,7 @@ private async Task<RecoverySettlement> 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;
Expand Down Expand Up @@ -338,6 +360,14 @@ private static bool ExactClaim(AgentRunLogStream value, RecoveryClaim expected)
private static string Code<T>(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)
Expand All @@ -348,6 +378,12 @@ private sealed record RecoveryClaim(Guid Id, Guid TeamId, Guid AgentRunId, long

private sealed record RecoverySettlement(AgentRunLogCaptureIntentState State, bool LostLease);

/// <summary>The outcome a settlement is writing: the observed one until the settlement supersedes or exhausts it.</summary>
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;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
using System.Collections.Concurrent;
using Autofac;
using CodeSpace.Core.Persistence.Db;
using CodeSpace.Core.Persistence.Entities;
Expand All @@ -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;
Expand Down Expand Up @@ -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<string>().ShouldContain($"(id={scene.IntentId})");
entry.Exception.ShouldBeOfType<DbUpdateException>().InnerException.ShouldBeOfType<PostgresException>();
}

[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<OperationCanceledException>();
}

/// <summary>
/// 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.
/// </summary>
private async Task<SettlementScene> 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);
}

/// <summary>
/// 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.
/// </summary>
private async Task<RecordedEntry> 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<string>().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()
{
Expand Down Expand Up @@ -811,6 +904,30 @@ private async Task<AgentRunLogCaptureIntent[]> ReleaseAndReconcileUntilCompleted
return seen;
}

/// <summary>
/// Reconciles until a wave reaches THIS test's faulted settlement — <paramref name="faulted"/> 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.
/// </summary>
private static async Task<AgentRunLogCaptureRecoverySummary> ReconcileUntilFaultedAsync(AgentRunLogCaptureRecoveryService recovery, Func<bool> 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<AgentRunLogCaptureIntent> IntentAsync(World world)
{
using var scope = _fixture.BeginScope();
Expand All @@ -819,7 +936,7 @@ private async Task<AgentRunLogCaptureIntent> 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<AgentRunLogCaptureRecoveryService>? logger = null)
{
options ??= new RecoveryTestOptions();
using var scope = _fixture.BeginScope();
Expand All @@ -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
Expand Down Expand Up @@ -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);

/// <summary>
/// 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.
/// </summary>
private sealed class RecordedRecoveryLog : ILogger<AgentRunLogCaptureRecoveryService>
{
private readonly ConcurrentBag<RecordedEntry> _entries = [];

/// <summary>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.</summary>
public IReadOnlyList<RecordedEntry> About(Guid intentId) => _entries.Where(entry => Equals(entry.Properties.GetValueOrDefault("IntentId"), intentId)).ToList();

public IDisposable? BeginScope<TState>(TState state) where TState : notnull => null;
public bool IsEnabled(LogLevel logLevel) => true;

public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func<TState, Exception?, string> formatter)
{
if (state is not IReadOnlyList<KeyValuePair<string, object?>> 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<string, object?> Properties, Exception? Exception);

private sealed class EmptyCas : IArtifactCasRuntimeCoordinator
{
public Task<ArtifactCasTransferResult> PutAsync(ArtifactCasTransferRequest request, CancellationToken cancellationToken) => Task.FromResult<ArtifactCasTransferResult>(new ArtifactCasTransferResult.Rejected(null, new ArtifactCasProblem(ArtifactCasProblemCode.Unsupported, false)));
Expand Down
Loading
Loading