diff --git a/backend/src/CodeSpace.Core/Persistence/DbUpFiles/0237_durable_retention_record_deadlines.sql b/backend/src/CodeSpace.Core/Persistence/DbUpFiles/0237_durable_retention_record_deadlines.sql new file mode 100644 index 000000000..1ef743a5b --- /dev/null +++ b/backend/src/CodeSpace.Core/Persistence/DbUpFiles/0237_durable_retention_record_deadlines.sql @@ -0,0 +1,34 @@ +-- 0237_durable_retention_record_deadlines.sql +-- +-- The quarantine marker for the record plane that gets a reaper in this change: settled cleanup receipts. +-- +-- WHY EACH PLANE NEEDS ITS OWN. The retention protocol has two independent waits — an age floor measured from the +-- record's own terminal instant, and a quarantine measured from the FIRST observation that nothing cites it. The +-- second wait only exists if it is written down; a sweep that decided "uncited" in memory and deleted on the same +-- tick has one wait, not two, whatever its code says. `agent_run_log_stream` got this in 0235; this table has nowhere +-- to record it at all. +-- +-- WHY ONE COLUMN IS ENOUGH. The cursor excludes a permanently-uncollectable receipt — a pinned one, an orphan — in the +-- CLAIM query rather than settling it, so the only rows claimed and not collected are the ones waiting out a deadline +-- this column holds. That is what makes one column enough here where the log stream also needed its modification time: +-- nothing else can occupy a batch slot for ever. +-- +-- It is a nullable ADD COLUMN with no default: metadata-only, no table rewrite, and NULL is exactly what an older +-- binary writes and reads. `agent_run_cleanup_receipt` carries no trigger, so no guard has to learn about it. +-- +-- THE INDEX takes a SHARE lock for the length of its build — concurrent reads continue, concurrent INSERTs wait. It is +-- built in-script rather than CONCURRENTLY because DbUp runs each script in one transaction and CONCURRENTLY cannot +-- appear inside one; the table holds a handful of rows per abandoned run, so the build is short. Its shape is the +-- claim query's exactly: the per-team head of the settled queue, oldest first. +-- +-- Rollback: DROP INDEX ix_agent_run_cleanup_receipt_retention; ALTER TABLE agent_run_cleanup_receipt DROP COLUMN +-- retain_until. Nothing reads either from an older binary. + +ALTER TABLE agent_run_cleanup_receipt ADD COLUMN retain_until timestamptz NULL; + +CREATE INDEX ix_agent_run_cleanup_receipt_retention + ON agent_run_cleanup_receipt (team_id, recorded_at, id) + WHERE outcome IN ('Completed', 'Compensated'); + +COMMENT ON COLUMN agent_run_cleanup_receipt.retain_until IS + 'The earliest instant this receipt may be reclaimed, written by the first sweep that found it settled, past its rule and cited by nobody. NULL means no sweep has proposed it.'; diff --git a/backend/src/CodeSpace.Core/Persistence/Entities/AgentRunCleanupReceiptRecord.cs b/backend/src/CodeSpace.Core/Persistence/Entities/AgentRunCleanupReceiptRecord.cs index 575beca83..7c2b7303d 100644 --- a/backend/src/CodeSpace.Core/Persistence/Entities/AgentRunCleanupReceiptRecord.cs +++ b/backend/src/CodeSpace.Core/Persistence/Entities/AgentRunCleanupReceiptRecord.cs @@ -21,5 +21,8 @@ public sealed class AgentRunCleanupReceiptRecord : IEntity public DateTimeOffset RecordedAt { get; set; } public string? ErrorCode { get; set; } + /// The earliest instant the retention plane may reclaim this receipt, written by the first sweep that found it settled and cited by nobody. Null until then, and on every row an older binary wrote. + public DateTimeOffset? RetainUntil { get; set; } + public AgentRun AgentRun { get; set; } = default!; } diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/CleanupReceiptRetentionCursor.cs b/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/CleanupReceiptRetentionCursor.cs new file mode 100644 index 000000000..50ea508c9 --- /dev/null +++ b/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/CleanupReceiptRetentionCursor.cs @@ -0,0 +1,201 @@ +using CodeSpace.Core.DependencyInjection; +using CodeSpace.Core.Persistence.Db; +using CodeSpace.Core.Persistence.Entities; +using CodeSpace.Messages.Agents.Recovery; +using CodeSpace.Messages.Retention; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; + +namespace CodeSpace.Core.Services.Workflows.Retention.Cursors; + +/// +/// Reclaims a settled agent_run_cleanup_receipt row that nothing cites any more. +/// +/// A receipt is a ledger row, not bytes, so reclaiming it means deleting it. There is nothing to purge +/// first and no tombstone to leave: the receipt's whole content is "what became of this resource", and a row that is +/// gone reads as a run that left nothing behind — which is exactly what a settled receipt already says. Deletion is +/// admitted: agent_run_cleanup_receipt carries no trigger, and its only foreign key points AT the run, so a +/// receipt can go while its run stays. +/// +/// Which outcomes are settled. Exactly the two RunCleanupReceipt.IsSettled already names: +/// and . Neither of the other +/// two is, and each for its own reason. is addressed to a sweep on the host +/// that owes the work and can still become Compensated; an orphan a retention pass removed is a resource nobody will +/// ever be told about again. is REWRITABLE — the ledger's upsert fences only +/// Completed and Compensated, and the orphan sweep answers Unknown when the host cannot even attempt a teardown, so a +/// live leak can move from Orphaned to Unknown and leave the orphan queue for good. It is also READ: the Room counts +/// Unknown receipts of the sweepable kinds and renders "N resources with an unknown cleanup state" with no age bound +/// at all, so draining them per-team-head would walk that count down to a card that quietly disappears — a wrong +/// number where the truth used to be. +/// +/// Starvation. A pinned receipt and an unsettled one are both excluded by the CLAIM query rather than +/// settled, so neither occupies a batch slot — which is why this plane needs one deadline column and not two. A keep +/// the cursor CAN reach — an unanswerable citation question — pushes the same deadline forward instead. Every claim +/// predicate is repeated in the deleting statement for the reason every fenced reaper repeats its guards: the claim +/// and the deletion are different transactions, and a pin, or a re-opened orphan, can land between them. +/// +public sealed class CleanupReceiptRetentionCursor : IDurableRetentionCursor, IScopedDependency +{ + /// + /// Every place a settled cleanup receipt can still be cited. One entry, and the list exists so that a SECOND one + /// is a deliberate edit: a citer missing from it makes the cursor delete a row something still points at, which is + /// the one failure mode of this class. Public and pinned by a test, like the log-stream cursor's. + /// + /// Deliberately NOT citers, each for a stated reason. RoomProjector.AgentRecoveryAsync reads every + /// receipt of a turn's agents, but by RUN and as an aggregate — a row that is gone folds to "nothing unsettled", + /// which is what a settled receipt already meant. AgentRunOrphanReaper reads only Orphaned rows, which this + /// cursor never claims. + /// + public static readonly IReadOnlyList<(string Table, string Column)> CitationSites = + [ + ("paired_qualification_result_pin", "pinned_id"), + ]; + + /// + /// The outcomes a receipt is SETTLED in — the same two RunCleanupReceipt.IsSettled names, and the same two + /// the ledger's upsert refuses to overwrite. Every absence is deliberate and pinned by a test. + /// + /// is the same set as the raw SQL sees it. The claim query cannot + /// parameterise an IN list of enum names, so the two forms exist side by side — and a test compares them, + /// because two copies of one rule drift silently. + /// + internal static readonly RunResourceOutcome[] SettledOutcomes = [RunResourceOutcome.Completed, RunResourceOutcome.Compensated]; + + /// The literal the claim and the delete use. Compared against by a test that reads this file's own SQL. + internal const string SettledOutcomeNames = "'Completed', 'Compensated'"; + + /// The pin kind the claim and the probe ask about, as the raw SQL spells it. Pinned against by the same test. + internal const string PinnedKindName = "CleanupReceipt"; + + private readonly DbContextOptions _dbOptions; + private readonly ILogger _logger; + + public CleanupReceiptRetentionCursor(DbContextOptions dbOptions, ILogger logger) + { + _dbOptions = dbOptions; + _logger = logger; + } + + public DurableRecordClass Class => DurableRecordClass.CleanupReceipt; + + /// + /// The per-team head of the eligible queue, oldest first. DISTINCT ON (team_id) is materialized first so one + /// tenant with a long backlog cannot starve the others, and every guard is repeated on the outer select because + /// quals inside a materialized CTE are not re-evaluated against a row another worker changed in between. + /// + public async Task> ClaimAsync(DurableRetentionSweepWindow window, int limit, CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(window); + await using var db = CreateDb(); + + var rows = await db.AgentRunCleanupReceipt.FromSqlInterpolated($$""" + WITH eligible AS MATERIALIZED ( + SELECT DISTINCT ON (receipt.team_id) receipt.id, receipt.recorded_at + FROM agent_run_cleanup_receipt receipt + WHERE receipt.outcome IN ('Completed', 'Compensated') + AND receipt.recorded_at <= {{window.TerminalBefore}} + AND (receipt.retain_until IS NULL OR receipt.retain_until <= {{window.Now}}) + AND NOT EXISTS (SELECT 1 FROM paired_qualification_result_pin pin WHERE pin.kind = 'CleanupReceipt' AND pin.pinned_id = receipt.id) + ORDER BY receipt.team_id, receipt.recorded_at, receipt.id + ) + SELECT receipt.* FROM agent_run_cleanup_receipt receipt + JOIN eligible ON eligible.id = receipt.id + WHERE receipt.outcome IN ('Completed', 'Compensated') + AND receipt.recorded_at <= {{window.TerminalBefore}} + AND (receipt.retain_until IS NULL OR receipt.retain_until <= {{window.Now}}) + ORDER BY receipt.recorded_at, receipt.id + LIMIT {{limit}} + """).AsNoTracking().ToListAsync(cancellationToken).ConfigureAwait(false); + + return rows.Select(receipt => new DurableRetentionCandidate(receipt.Id, receipt.TeamId, 0, receipt.RecordedAt, receipt.RetainUntil)).ToArray(); + } + + /// Fail-closed: any failure to probe the citation site answers indeterminate, which every consumer reads as keep. + public async Task ClassifyAsync(DurableRetentionCandidate candidate, CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(candidate); + + try + { + await using var db = CreateDb(); + var pinned = await db.PairedQualificationResultPin.AsNoTracking() + .AnyAsync(pin => pin.Kind == DurablePinKind.CleanupReceipt && pin.PinnedId == candidate.Id, cancellationToken).ConfigureAwait(false); + + return pinned ? DurableReferenceVerdict.Referenced : DurableReferenceVerdict.Unreferenced; + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } + catch (Exception ex) + { + _logger.LogWarning(ex, "Cleanup receipt {ReceiptId}: citation sites could not be probed; the receipt is kept", candidate.Id); + + return DurableReferenceVerdict.Indeterminate; + } + } + + /// + /// A keep this cursor cannot record any other way pushes the deadline FORWARD by the class's recheck interval — + /// the row has no modification time of its own, so without that a receipt whose citation question could not be + /// answered would be re-claimed on every tick for ever. + /// + /// A CITATION is the exception, and honouring it is not optional: the decision clears the marker on + /// purpose, because a quarantine that elapsed while something still pointed at the row never waited for anything. + /// Writing a deadline there instead would let the row be collected on the FIRST uncited observation after the + /// citation went away — the second wait skipped, which is the one thing this column exists to prevent. It costs + /// no hot loop: a still-cited receipt is excluded by the claim query, not re-claimed. + /// + public async Task SettleAsync(DurableRetentionSweepWindow window, DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(window); + ArgumentNullException.ThrowIfNull(candidate); + ArgumentNullException.ThrowIfNull(decision); + + if (decision.Action == DurableRetentionAction.Collect) return await CollectAsync(window, candidate, cancellationToken).ConfigureAwait(false); + + var deadline = decision.Action switch + { + DurableRetentionAction.Quarantine => decision.RetainUntil, + DurableRetentionAction.Referenced => null, + _ => window.RecheckAt, + }; + + return await StampAsync(window, candidate, deadline, cancellationToken).ConfigureAwait(false); + } + + /// + /// The delete, with every condition that admitted it repeated in its own statement. The outcome guard is what + /// keeps a receipt an orphan sweep re-opened between the claim and now — Orphaned again, addressed to a host — + /// from being removed by a decision taken while it still read as settled. + /// + private async Task CollectAsync(DurableRetentionSweepWindow window, DurableRetentionCandidate candidate, CancellationToken cancellationToken) + { + await using var db = CreateDb(); + + var deleted = await db.AgentRunCleanupReceipt + .Where(receipt => receipt.Id == candidate.Id && receipt.TeamId == candidate.TeamId + && SettledOutcomes.Contains(receipt.Outcome) + && receipt.RecordedAt <= window.TerminalBefore + && receipt.RetainUntil != null && receipt.RetainUntil <= window.Now + && !db.PairedQualificationResultPin.Any(pin => pin.Kind == DurablePinKind.CleanupReceipt && pin.PinnedId == receipt.Id)) + .ExecuteDeleteAsync(cancellationToken).ConfigureAwait(false); + + if (deleted == 1) + _logger.LogInformation("Cleanup receipt {ReceiptId} reclaimed for team {TeamId}: settled, past its rule, and cited by nothing", candidate.Id, candidate.TeamId); + + return deleted == 1; + } + + /// The deadline, written only onto a row that is still the one the claim saw — the receipt has no revision, so its identity, its outcome and its own age are the fence. + private async Task StampAsync(DurableRetentionSweepWindow window, DurableRetentionCandidate candidate, DateTimeOffset? retainUntil, CancellationToken cancellationToken) + { + await using var db = CreateDb(); + + var updated = await db.AgentRunCleanupReceipt + .Where(receipt => receipt.Id == candidate.Id && receipt.TeamId == candidate.TeamId + && SettledOutcomes.Contains(receipt.Outcome) && receipt.RecordedAt <= window.TerminalBefore) + .ExecuteUpdateAsync(set => set.SetProperty(receipt => receipt.RetainUntil, retainUntil), cancellationToken).ConfigureAwait(false); + + return updated == 1; + } + + private CodeSpaceDbContext CreateDb() => new(_dbOptions); +} diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/LogStreamRetentionCursor.cs b/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/LogStreamRetentionCursor.cs index 9acbd7f88..40789f67c 100644 --- a/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/LogStreamRetentionCursor.cs +++ b/backend/src/CodeSpace.Core/Services/Workflows/Retention/Cursors/LogStreamRetentionCursor.cs @@ -140,8 +140,9 @@ public async Task ClassifyAsync(DurableRetentionCandida } } - public async Task SettleAsync(DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken) + public async Task SettleAsync(DurableRetentionSweepWindow window, DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken) { + ArgumentNullException.ThrowIfNull(window); ArgumentNullException.ThrowIfNull(candidate); ArgumentNullException.ThrowIfNull(decision); diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionPolicy.cs b/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionPolicy.cs index 146c80fd4..70528a1ee 100644 --- a/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionPolicy.cs +++ b/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionPolicy.cs @@ -7,26 +7,50 @@ namespace CodeSpace.Core.Services.Workflows.Retention; /// are kept. Values are committed here and changed by a pull request — there is no environment override, because a /// mistyped retention window is unrecoverable data loss and a code review is the control that belongs in front of it. /// -/// One class, because one cursor exists. Every other plane a reaper will eventually reach — cleanup receipts, -/// capture gaps, qualification evidence, budget reservations, terminal transfer intents — gets its rule in the same -/// commit as the cursor that acts on it, so this table can always be read as "what is actually reclaimed". +/// A class is listed here together with the cursor that sweeps it, never ahead of one, so this table can always +/// be read as "what is actually reclaimed". Three planes that might be expected are deliberately absent, each for a +/// reason that survived being looked for: +/// +/// Migration 0146's guard refuses to delete a run's capture-gap evidence, because a removable gap makes a +/// complete manifest reachable by deleting the evidence for it; a gap can only honestly go with the manifest whose +/// verdict it qualifies. Migration 0219's trigger refuses to delete or update a paired-qualification result, and a result's +/// citers are not enumerable in columns at all — a qualification receipt records its cohort and its metrics as JSON, +/// which is the shape this plane's charter excludes. And a terminal budget reservation is read UNWINDOWED by the +/// Room: RoomProjector.BudgetAsync sums every reservation of a run to state what it committed, so reclaiming +/// the oldest rows of a long run would leave that figure quietly wrong rather than absent. Reclaiming spend records +/// needs a durable per-run summary to survive them first, which is a money-plane change and not a retention one. /// /// An unregistered class is NOT an error a cursor can shrug off: returns null and the decision /// settles Indeterminate, which keeps the record forever. That is what makes removing a class from this table safe. /// public static class DurableRetentionPolicy { + private static readonly TimeSpan Quarantine = TimeSpan.FromHours(24); + private static readonly TimeSpan Recheck = TimeSpan.FromHours(24); + /// /// Thirty days after the stream reached its terminal capture state, then a day of quarantine, and a day between /// looks at a stream that was kept. Deliberately far longer than any run and longer than any Room a human is still /// reading: the bytes this reclaims are an archive nobody opened, and the floor costs storage rather than evidence. /// public static readonly DurableRetentionRule LogStream = - new(DurableRecordClass.LogStream, TimeSpan.FromDays(30), TimeSpan.FromHours(24), TimeSpan.FromHours(24)); + new(DurableRecordClass.LogStream, TimeSpan.FromDays(30), Quarantine, Recheck); + + /// + /// Thirty days after the receipt was recorded, which is at or after the run's own terminal instant — a receipt is + /// only ever written when a run is abandoned, and a compensating sweep writes later still. Measuring from the row + /// keeps the floor on a column the claim query can read, and can never make it earlier than the rule it states. + /// + public static readonly DurableRetentionRule CleanupReceipt = + new(DurableRecordClass.CleanupReceipt, TimeSpan.FromDays(30), Quarantine, Recheck); /// The committed table, public so a test can pin every literal value in it. public static readonly IReadOnlyDictionary Rules = - new Dictionary { [LogStream.Class] = LogStream }; + new Dictionary + { + [LogStream.Class] = LogStream, + [CleanupReceipt.Class] = CleanupReceipt, + }; /// The rule for , or null when this build registers none — which every consumer reads as keep. public static DurableRetentionRule? For(DurableRecordClass value) => Rules.TryGetValue(value, out var rule) ? rule : null; diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionReaper.cs b/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionReaper.cs index 11e4010af..eb9e3b0ac 100644 --- a/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionReaper.cs +++ b/backend/src/CodeSpace.Core/Services/Workflows/Retention/DurableRetentionReaper.cs @@ -31,7 +31,12 @@ namespace CodeSpace.Core.Services.Workflows.Retention; /// public sealed class DurableRetentionReaper : IDurableRetentionReaper { - /// Records claimed per cursor per sweep. The cadence is hourly and every rule is measured in days, so the ceiling changes only how promptly a backlog drains — never what is collected. + /// + /// Records claimed per cursor per sweep, and PER CURSOR rather than shared across them: one budget spent in class + /// order would let a plane with a standing backlog starve every plane after it for ever, while the sweep reported + /// a healthy claim count. The cadence is hourly and every rule is measured in days, so the ceiling changes only + /// how promptly one plane's backlog drains — never what is collected. + /// private const int BatchSize = 200; /// Records asked for per claim. A fair cursor returns the per-tenant head, so the batch above is reached by asking repeatedly rather than in one query, and a tenant with a long backlog never holds the whole batch. @@ -72,12 +77,13 @@ private async Task SweepCursorAsync(IDurableRetentionCursor cursor, DateTimeOffs return; } - var window = new DurableRetentionSweepWindow(now, now.Subtract(rule.MinimumAge), now.Subtract(rule.RecheckInterval)); + var window = new DurableRetentionSweepWindow(now, now.Subtract(rule.MinimumAge), now.Subtract(rule.RecheckInterval), now.Add(rule.RecheckInterval)); var seen = new HashSet(); + var claimedHere = 0; - while (counts.Claimed < BatchSize) + while (claimedHere < BatchSize) { - var limit = Math.Min(ClaimSize, BatchSize - counts.Claimed); + var limit = Math.Min(ClaimSize, BatchSize - claimedHere); var claimed = await cursor.ClaimAsync(window, limit, cancellationToken).ConfigureAwait(false); // A cursor that claims fairly hands back one record per tenant, so the batch is reached by asking again — // and the ids already settled this tick are excluded, because a settlement that writes nothing (a drain @@ -86,16 +92,17 @@ private async Task SweepCursorAsync(IDurableRetentionCursor cursor, DateTimeOffs if (candidates.Count == 0) break; + claimedHere += candidates.Count; counts.Claimed += candidates.Count; foreach (var candidate in candidates) - counts.Record(await SweepCandidateAsync(cursor, rule, candidate, now, cancellationToken).ConfigureAwait(false)); + counts.Record(await SweepCandidateAsync(cursor, window, rule, candidate, cancellationToken).ConfigureAwait(false)); } } /// One candidate, start to finish. Every exit that is not a completed settlement keeps the record. - private async Task SweepCandidateAsync(IDurableRetentionCursor cursor, DurableRetentionRule rule, DurableRetentionCandidate candidate, - DateTimeOffset now, CancellationToken cancellationToken) + private async Task SweepCandidateAsync(IDurableRetentionCursor cursor, DurableRetentionSweepWindow window, DurableRetentionRule rule, + DurableRetentionCandidate candidate, CancellationToken cancellationToken) { using var operation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); operation.CancelAfter(CandidateTimeout); @@ -103,9 +110,9 @@ private async Task SweepCandidateAsync(IDurableRetention try { var verdict = await cursor.ClassifyAsync(candidate, operation.Token).ConfigureAwait(false); - var decision = DurableRetentionDecision.Decide(rule, new DurableRetentionObservation(candidate.TerminalAt, candidate.RetainUntil, verdict, now)); + var decision = DurableRetentionDecision.Decide(rule, new DurableRetentionObservation(candidate.TerminalAt, candidate.RetainUntil, verdict, window.Now)); - return await ApplyAsync(cursor, candidate, decision, operation.Token).ConfigureAwait(false); + return await ApplyAsync(cursor, window, candidate, decision, operation.Token).ConfigureAwait(false); } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } catch (Exception ex) @@ -121,9 +128,9 @@ private async Task SweepCandidateAsync(IDurableRetention /// every tick for ever. A settlement the cursor could not complete — a lost race, a refused removal, a drain that /// needs another sweep — is reported as the keep it is. /// - private static async Task ApplyAsync(IDurableRetentionCursor cursor, DurableRetentionCandidate candidate, - DurableRetentionDecision decision, CancellationToken cancellationToken) => - await cursor.SettleAsync(candidate, decision, cancellationToken).ConfigureAwait(false) + private static async Task ApplyAsync(IDurableRetentionCursor cursor, DurableRetentionSweepWindow window, + DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken) => + await cursor.SettleAsync(window, candidate, decision, cancellationToken).ConfigureAwait(false) ? decision.Action : DurableRetentionAction.Indeterminate; diff --git a/backend/src/CodeSpace.Core/Services/Workflows/Retention/IDurableRetentionCursor.cs b/backend/src/CodeSpace.Core/Services/Workflows/Retention/IDurableRetentionCursor.cs index 847f002b5..92af2cbb7 100644 --- a/backend/src/CodeSpace.Core/Services/Workflows/Retention/IDurableRetentionCursor.cs +++ b/backend/src/CodeSpace.Core/Services/Workflows/Retention/IDurableRetentionCursor.cs @@ -29,8 +29,9 @@ public sealed record DurableRetentionCandidate(Guid Id, Guid TeamId, long Revisi /// /// The database clock at sweep start. Every deadline compared against it was written by a database clock too. /// The age floor: a record that went terminal after this is not a candidate at all. -/// The deferral: a record last settled after this was already looked at recently and is left alone. -public sealed record DurableRetentionSweepWindow(DateTimeOffset Now, DateTimeOffset TerminalBefore, DateTimeOffset RecheckBefore); +/// The deferral, read backwards: a record last settled after this was already looked at recently and is left alone. +/// The same deferral, written forwards: where a cursor that keeps a record for a reason it cannot record puts its next look. +public sealed record DurableRetentionSweepWindow(DateTimeOffset Now, DateTimeOffset TerminalBefore, DateTimeOffset RecheckBefore, DateTimeOffset RecheckAt); /// /// One plane's answers to the three questions the reaper asks: which records are candidates, does anything still cite @@ -58,10 +59,15 @@ public interface IDurableRetentionCursor Task ClassifyAsync(DurableRetentionCandidate candidate, CancellationToken cancellationToken); /// - /// Applies under the candidate's revision fence. False means nothing was settled — - /// the record moved under this sweep, or the work it needed could not be finished — and the loop reports it as a + /// Applies under the candidate's own fence. False means nothing was settled — the + /// record moved under this sweep, or the work it needed could not be finished — and the loop reports it as a /// keep. A cursor decides for itself whether an unfinished settlement also defers the record; work that is making /// progress must NOT, or a drain that needs several passes would take one recheck interval per pass. + /// + /// is the SAME window the claim used, handed back so a deleting statement can + /// repeat the time predicates that admitted the row. The claim and the deletion are different transactions, and a + /// record whose terminal instant or deadline moved in between must match nothing rather than be deleted on the + /// strength of a decision taken about the row it used to be. /// - Task SettleAsync(DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken); + Task SettleAsync(DurableRetentionSweepWindow window, DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken); } diff --git a/backend/src/CodeSpace.Messages/Retention/DurableRetention.cs b/backend/src/CodeSpace.Messages/Retention/DurableRetention.cs index b05b57fec..f5e3c61f4 100644 --- a/backend/src/CodeSpace.Messages/Retention/DurableRetention.cs +++ b/backend/src/CodeSpace.Messages/Retention/DurableRetention.cs @@ -14,6 +14,9 @@ public enum DurableRecordClass { /// An agent_run_log_stream head and the routed CAS bytes its segments name. LogStream = 1, + + /// An agent_run_cleanup_receipt row (migration 0229) that settled — what became of one resource of one abandoned run. + CleanupReceipt = 2, } /// diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRecordRetentionFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRecordRetentionFlowTests.cs new file mode 100644 index 000000000..65072050a --- /dev/null +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRecordRetentionFlowTests.cs @@ -0,0 +1,492 @@ +using Autofac; +using CodeSpace.Core.Persistence.Db; +using CodeSpace.Core.Persistence.Entities; +using CodeSpace.Core.Services.Sessions.Room; +using CodeSpace.Core.Services.Workflows.Retention; +using CodeSpace.Core.Services.Workflows.Retention.Cursors; +using CodeSpace.IntegrationTests.Infrastructure; +using CodeSpace.Messages.Agents.Recovery; +using CodeSpace.Messages.Enums; +using CodeSpace.Messages.Retention; +using Microsoft.EntityFrameworkCore; +using Shouldly; + +namespace CodeSpace.IntegrationTests.Workflows; + +/// +/// The one record plane the retention reaper reclaims by DELETING rows rather than bytes: settled cleanup receipts. A +/// receipt is a ledger row whose whole content is a statement about something that is over, and it is protected by the +/// same two waits as every other class — an age floor from the record's own terminal instant, then a quarantine from +/// the first observation that nothing cites it. +/// +/// Every assertion is about a row this class created. A sweep is deployment-wide over a shared database, so its +/// tallies count other classes' rows as readily as these; the counter-examples assert that THEIR row survived, and the +/// positive controls that THEIRS is gone. +/// +/// The rows staged here carry no bytes and no artifact placements, so this class leaves nothing behind for the +/// bounded sweeps elsewhere in the suite to trip over — what it does not reclaim, it deletes in teardown. +/// +[Collection(PostgresCollection.Name)] +[Trait("Category", "Integration")] +public sealed class DurableRecordRetentionFlowTests : IAsyncLifetime +{ + private readonly PostgresFixture _fixture; + private readonly List _teams = []; + + public DurableRecordRetentionFlowTests(PostgresFixture fixture) => _fixture = fixture; + + /// The positive control. Two sweeps: the first records the quarantine deadline and removes nothing, the second collects. + [Theory] + [InlineData(RunResourceOutcome.Completed)] + [InlineData(RunResourceOutcome.Compensated)] + public async Task A_settled_receipt_past_its_rule_is_reclaimed(RunResourceOutcome outcome) + { + var world = await SeedWorldAsync(); + var receipt = await ReceiptAsync(world, outcome, RunResourceKind.Spool); + await AgeReceiptAsync(receipt, TimeSpan.FromDays(31)); + + await SweepAsync(); + + (await ReceiptExistsAsync(receipt)).ShouldBeTrue("the sweep that first noticed the receipt must not delete it"); + (await RetainUntilAsync(receipt)).ShouldNotBeNull("the quarantine deadline has to be durable — an in-memory one is no wait at all"); + + await ElapseQuarantineAsync(receipt); + await SweepAsync(); + + (await ReceiptExistsAsync(receipt)).ShouldBeFalse("both waits elapsed with nothing citing it"); + } + + /// + /// An orphan is addressed to a sweep on the host that owes the teardown and can still become Compensated. + /// Mutation: admit Orphaned into the claim and the allow-list, and this reds — with a leaked resource nobody will + /// ever be told about. + /// + [Fact] + public async Task An_orphaned_receipt_is_never_claimed_however_old() + { + var world = await SeedWorldAsync(); + var orphan = await ReceiptAsync(world, RunResourceOutcome.Orphaned, RunResourceKind.EgressSubnet, ownerHost: "some-dead-host"); + await AgeReceiptAsync(orphan, TimeSpan.FromDays(400)); + + (await ClaimedAsync()).ShouldNotContain(orphan); + + await SweepAsync(); + await SweepAsync(); + + (await ReceiptExistsAsync(orphan)).ShouldBeTrue("the row naming the host that still owes this teardown is the only record of it"); + (await RetainUntilAsync(orphan)).ShouldBeNull("and it is never even quarantined, because no sweep may propose collecting it"); + } + + /// + /// An Unknown receipt is neither settled nor beyond reach: the ledger's upsert fences only Completed and + /// Compensated, so a live orphan the owning host cannot even attempt to tear down moves Orphaned → Unknown and + /// leaves the orphan queue for good. Mutation: put Unknown back in the reclaimable set and this reds — with a leak + /// that had left the one queue that watches it now deleted by policy as well. + /// + [Fact] + public async Task An_unknown_receipt_is_never_claimed_however_old() + { + var world = await SeedWorldAsync(); + var unknown = await ReceiptAsync(world, RunResourceOutcome.Unknown, RunResourceKind.Cgroup); + await AgeReceiptAsync(unknown, TimeSpan.FromDays(400)); + + (await ClaimedAsync()).ShouldNotContain(unknown); + + await SweepAsync(); + await SweepAsync(); + + (await ReceiptExistsAsync(unknown)).ShouldBeTrue("nothing has settled this resource; the row is what says so"); + (await RetainUntilAsync(unknown)).ShouldBeNull(); + } + + /// + /// The Room surface, which is why the sweepable outcomes are exactly the two the ledger calls settled. The + /// recovery card counts Unknown receipts of the kinds a sweep CAN reach, over every receipt of the run and with no + /// age bound at all — so a plane that drained them would walk that count down over days and then remove the card, + /// putting a wrong number where the truth used to be. This asserts the card reads identically before and after a + /// sweep old enough to have reclaimed everything it may. + /// + [Fact] + public async Task The_rooms_recovery_card_reads_identically_before_and_after_a_sweep() + { + var world = await SeedWorldAsync(); + await ReceiptAsync(world, RunResourceOutcome.Unknown, RunResourceKind.Cgroup); + await ReceiptAsync(world, RunResourceOutcome.Unknown, RunResourceKind.Workspace); + await ReceiptAsync(world, RunResourceOutcome.Orphaned, RunResourceKind.EgressSubnet, ownerHost: "dead-host"); + var settled = await ReceiptAsync(world, RunResourceOutcome.Completed, RunResourceKind.Spool); + foreach (var receipt in await ReceiptsOfAsync(world)) await AgeReceiptAsync(receipt, TimeSpan.FromDays(400)); + + var before = RoomProjector.SummarizeRecovery(await RunCleanupReceiptsAsync(world)).ShouldNotBeNull(); + before.UnknownCount.ShouldBe(2, "the premise: the card is counting rows this sweep would otherwise be free to take"); + + await SweepAsync(); + // Every receipt's deadline, not only the one that may go: a test that elapsed the quarantine of the settled + // row alone would pass with Unknown rows merely waiting rather than protected. + foreach (var receipt in await ReceiptsOfAsync(world)) await ElapseQuarantineAsync(receipt); + await SweepAsync(); + + (await ReceiptExistsAsync(settled)).ShouldBeFalse("the settled receipt IS reclaimed — otherwise this test passes by reclaiming nothing at all"); + var after = RoomProjector.SummarizeRecovery(await RunCleanupReceiptsAsync(world)).ShouldNotBeNull(); + after.UnknownCount.ShouldBe(before.UnknownCount, "an operator's unresolved-cleanup count must not fall because a retention pass ran"); + after.OrphanedCount.ShouldBe(before.OrphanedCount); + after.Detail.ShouldBe(before.Detail, "and the sentence the card renders is the same sentence"); + } + + /// + /// A receipt a sealed qualification result cites is evidence the result was computed from. Mutation: drop the pin + /// predicate from the claim and the probe, and this reds with the evidence gone. + /// + [Fact] + public async Task A_pinned_receipt_past_its_rule_survives() + { + var world = await SeedWorldAsync(); + var receipt = await ReceiptAsync(world, RunResourceOutcome.Completed, RunResourceKind.Workspace); + await PinAsync(world, DurablePinKind.CleanupReceipt, receipt); + await AgeReceiptAsync(receipt, TimeSpan.FromDays(400)); + + (await ClaimedAsync()).ShouldNotContain(receipt, "a cited row must not even occupy a batch slot"); + await SweepAsync(); + + (await ReceiptExistsAsync(receipt)).ShouldBeTrue("a sealed result still cites this receipt; no elapsed window outranks that"); + (await Cursor().ClassifyAsync(CandidateFor(world, receipt), CancellationToken.None)) + .ShouldBe(DurableReferenceVerdict.Referenced, "and the verdict says WHY it was kept, not merely that it was"); + + // A citation lands between a claim and its classification — the one path on which a still-cited row IS + // settled. The marker it leaves behind decides what happens after the citation goes away: a deadline written + // here would let the very next uncited observation collect the row, skipping the second wait entirely. + var quarantined = await ReceiptAsync(world, RunResourceOutcome.Completed, RunResourceKind.Spool); + await AgeReceiptAsync(quarantined, TimeSpan.FromDays(31)); + await SweepAsync(); + (await RetainUntilAsync(quarantined)).ShouldNotBeNull("the premise: this row carries a quarantine a citation must now contradict"); + + // The deadline passes and the row is claimed again — the moment a pin can land between the claim and the + // classification, which is the only path on which a cited receipt is settled at all. + await ElapseQuarantineAsync(quarantined); + var claimed = (await CandidatesAsync()).Single(row => row.Id == quarantined); + var settled = await Cursor().SettleAsync(Window(), claimed, DurableRetentionDecision.Referenced(), CancellationToken.None); + + settled.ShouldBeTrue(); + (await RetainUntilAsync(quarantined)).ShouldBeNull( + "a citation CLEARS the quarantine it contradicts; leaving a deadline behind would collect the row on the first uncited observation after the pin went away"); + } + + /// The age floor at the CALL SITE: the window the loop computes from the class's rule is what keeps a young receipt out of the batch. + [Fact] + public async Task A_receipt_inside_its_age_floor_is_never_claimed() + { + var world = await SeedWorldAsync(); + var young = await ReceiptAsync(world, RunResourceOutcome.Completed, RunResourceKind.Cgroup); + await AgeReceiptAsync(young, DurableRetentionPolicy.CleanupReceipt.MinimumAge - TimeSpan.FromDays(1)); + + (await ClaimedAsync()).ShouldNotContain(young); + + await AgeReceiptAsync(young, TimeSpan.FromDays(2)); + + (await ClaimedAsync()).ShouldContain(young, "the counter-example is worth nothing unless the same row IS claimed once its floor has passed"); + } + + /// + /// The deleting statement repeats every predicate that admitted the row, because the claim and the delete are + /// different transactions. Mutation: drop the time predicates from CollectAsync's Where and a decision + /// taken about the row a receipt USED to be deletes the row it has become. + /// + [Theory] + [InlineData("outcome = 'Orphaned', owner_host = 'a-host-that-came-back'")] + [InlineData("recorded_at = now()")] + [InlineData("retain_until = now() + interval '30 days'")] + public async Task A_receipt_that_changed_after_the_claim_is_not_deleted(string moved) + { + var world = await SeedWorldAsync(); + var receipt = await ReceiptAsync(world, RunResourceOutcome.Completed, RunResourceKind.Spool); + await AgeReceiptAsync(receipt, TimeSpan.FromDays(31)); + await SweepAsync(); + await ElapseQuarantineAsync(receipt); + + // The claim, taken while the row still reads collectable, and then the row moves under it. + var candidate = (await CandidatesAsync()).Single(row => row.Id == receipt); + await MoveAsync(receipt, moved); + + var settled = await Cursor().SettleAsync(Window(), candidate, DurableRetentionDecision.Collect(DateTimeOffset.UtcNow.AddDays(-1)), CancellationToken.None); + + settled.ShouldBeFalse("the decision was taken about a row that no longer exists in that shape"); + (await ReceiptExistsAsync(receipt)).ShouldBeTrue(); + } + + /// + /// A receipt whose citation question could not be answered has nowhere else to record that it was looked at, so + /// the sweep pushes its own deadline forward. Mutation: settle only Quarantine and Collect, and an unanswerable + /// row is re-claimed on every tick for ever while the sweep reports a healthy claim count. + /// + [Fact] + public async Task A_receipt_the_sweep_kept_for_an_unanswerable_reason_is_deferred() + { + var world = await SeedWorldAsync(); + var receipt = await ReceiptAsync(world, RunResourceOutcome.Completed, RunResourceKind.McpSocket); + await AgeReceiptAsync(receipt, TimeSpan.FromDays(31)); + var candidate = (await CandidatesAsync()).Single(row => row.Id == receipt); + + var settled = await Cursor().SettleAsync(Window(), candidate, DurableRetentionDecision.Indeterminate("probe-unreachable", null), CancellationToken.None); + + settled.ShouldBeTrue(); + (await RetainUntilAsync(receipt)).ShouldNotBeNull().ShouldBeGreaterThan(DateTimeOffset.UtcNow.AddHours(23), + "the next look is a day away, not the next tick"); + (await ClaimedAsync()).ShouldNotContain(receipt); + } + + /// + /// The claim is fair per tenant, so one team's backlog can never fill a batch; the loop reaches the rest by asking + /// again. Mutation: claim once instead of looping and only the oldest receipt of the team is ever swept. + /// + [Fact] + public async Task One_sweep_reaches_every_eligible_receipt_of_a_team_not_just_its_oldest() + { + var world = await SeedWorldAsync(); + var receipts = new List(); + + foreach (var kind in new[] { RunResourceKind.Spool, RunResourceKind.Workspace, RunResourceKind.Cgroup }) + { + var receipt = await ReceiptAsync(world, RunResourceOutcome.Completed, kind); + await AgeReceiptAsync(receipt, TimeSpan.FromDays(31)); + receipts.Add(receipt); + } + + await SweepAsync(); + + foreach (var receipt in receipts) + (await RetainUntilAsync(receipt)).ShouldNotBeNull($"receipt {receipt} of the same team was never reached by the sweep"); + } + + /// + /// The batch budget is per cursor, not shared. Mutation: spend one budget across them in class order and the + /// plane behind a saturated one is never swept at all — while the sweep reports a healthy claim count, which is + /// what makes it invisible. + /// + [Fact] + public async Task A_saturated_cursor_does_not_starve_the_one_behind_it() + { + using var scope = _fixture.BeginScope(); + var saturating = new StubCursor(DurableRecordClass.LogStream, candidatesPerClaim: 25, claims: 20); + var behind = new StubCursor(DurableRecordClass.CleanupReceipt, candidatesPerClaim: 1, claims: 1); + var reaper = new DurableRetentionReaper(scope.Resolve>(), [saturating, behind], + Microsoft.Extensions.Logging.Abstractions.NullLogger.Instance); + + var summary = await reaper.SweepAsync(CancellationToken.None); + + saturating.Settled.ShouldBe(200, "the first plane is offered more than one batch and must be held to exactly one"); + behind.Settled.ShouldBe(1, "and the plane behind it still gets its own batch, rather than the remainder of somebody else's"); + summary.Claimed.ShouldBe(201); + } + + /// A cursor with no database behind it: it hands out as many candidates as it was told to and records what the loop did with them. + private sealed class StubCursor(DurableRecordClass recordClass, int candidatesPerClaim, int claims) : IDurableRetentionCursor + { + private int _remainingClaims = claims; + + public DurableRecordClass Class => recordClass; + public int Settled { get; private set; } + + public Task> ClaimAsync(DurableRetentionSweepWindow window, int limit, CancellationToken cancellationToken) + { + if (_remainingClaims-- <= 0) return Task.FromResult>([]); + + var terminal = window.TerminalBefore.AddDays(-1); + + return Task.FromResult>( + Enumerable.Range(0, Math.Min(limit, candidatesPerClaim)).Select(_ => new DurableRetentionCandidate(Guid.NewGuid(), Guid.NewGuid(), 0, terminal, null)).ToList()); + } + + public Task ClassifyAsync(DurableRetentionCandidate candidate, CancellationToken cancellationToken) => + Task.FromResult(DurableReferenceVerdict.Unreferenced); + + public Task SettleAsync(DurableRetentionSweepWindow window, DurableRetentionCandidate candidate, DurableRetentionDecision decision, CancellationToken cancellationToken) + { + Settled++; + + return Task.FromResult(true); + } + } + + // ── Staging ──────────────────────────────────────────────────────────────────────────────────────────────────── + + private async Task ReceiptAsync(World world, RunResourceOutcome outcome, RunResourceKind kind, string? ownerHost = null) + { + using var scope = _fixture.BeginScope(); + var db = scope.Resolve(); + var id = Guid.NewGuid(); + db.AgentRunCleanupReceipt.Add(new AgentRunCleanupReceiptRecord + { + Id = id, TeamId = world.TeamId, AgentRunId = world.AgentRunId, FenceEpoch = 7, Kind = kind, Outcome = outcome, + OwnerHost = ownerHost, ResourceKey = $"{kind}/{id:N}", RecordedByHost = "retention-test", RecordedAt = DateTimeOffset.UtcNow, + ErrorCode = outcome is RunResourceOutcome.Orphaned or RunResourceOutcome.Unknown ? "left-behind" : null, + }); + await db.SaveChangesAsync(); + + return id; + } + + /// A pin of the shape a sealed qualification result writes, without standing up a whole campaign to get one. + private async Task PinAsync(World world, DurablePinKind kind, Guid target) + { + using var scope = _fixture.BeginScope(); + var db = scope.Resolve(); + await db.Database.ExecuteSqlRawAsync( + "INSERT INTO paired_qualification_result_pin (id, result_id, kind, pinned_id, pinned_artifact_id, pinned_at) VALUES ({0}, {1}, {2}, {3}, NULL, now())", + [Guid.NewGuid(), await SealedResultAsync(world, Guid.NewGuid()), kind.ToString(), target]); + } + + private async Task AgeReceiptAsync(Guid receipt, TimeSpan age) => await ExecuteAsync( + "UPDATE agent_run_cleanup_receipt SET recorded_at = recorded_at - {0}::interval WHERE id = {1}", [$"{age.TotalSeconds} seconds", receipt]); + + private async Task ElapseQuarantineAsync(Guid receipt) + { + using var scope = _fixture.BeginScope(); + await scope.Resolve().Database.ExecuteSqlRawAsync( + "UPDATE agent_run_cleanup_receipt SET retain_until = now() - interval '1 second' WHERE id = {0}", [receipt]); + } + + /// Moves the row under a claim already taken — the race the deleting statement's repeated predicates exist for. + private async Task MoveAsync(Guid receipt, string change) => await ExecuteAsync( + $"UPDATE agent_run_cleanup_receipt SET {change} WHERE id = {{0}}", [receipt]); + + private async Task ExecuteAsync(string sql, object[] parameters) + { + using var scope = _fixture.BeginScope(); + (await scope.Resolve().Database.ExecuteSqlRawAsync(sql, parameters)).ShouldBe(1); + } + + // ── Reads ────────────────────────────────────────────────────────────────────────────────────────────────────── + + private async Task SweepAsync() + { + using var scope = _fixture.BeginScope(); + await scope.Resolve().SweepAsync(CancellationToken.None); + } + + private IDurableRetentionCursor Cursor() + { + using var scope = _fixture.BeginScope(); + + return scope.Resolve>().Single(cursor => cursor.Class == DurableRecordClass.CleanupReceipt); + } + + /// The window the loop would build right now, from the class's own rule. + private static DurableRetentionSweepWindow Window() + { + var rule = DurableRetentionPolicy.CleanupReceipt; + var now = DateTimeOffset.UtcNow; + + return new DurableRetentionSweepWindow(now, now - rule.MinimumAge, now - rule.RecheckInterval, now + rule.RecheckInterval); + } + + private async Task> CandidatesAsync() => + await Cursor().ClaimAsync(Window(), 200, CancellationToken.None); + + private async Task> ClaimedAsync() => (await CandidatesAsync()).Select(candidate => candidate.Id).ToList(); + + private static DurableRetentionCandidate CandidateFor(World world, Guid id) => new(id, world.TeamId, 0, DateTimeOffset.UtcNow.AddDays(-400), null); + + private async Task ReceiptExistsAsync(Guid receipt) + { + using var scope = _fixture.BeginScope(); + + return await scope.Resolve().AgentRunCleanupReceipt.AsNoTracking().AnyAsync(row => row.Id == receipt); + } + + private async Task RetainUntilAsync(Guid receipt) + { + using var scope = _fixture.BeginScope(); + + return await scope.Resolve().AgentRunCleanupReceipt.AsNoTracking().Where(row => row.Id == receipt).Select(row => row.RetainUntil).SingleAsync(); + } + + private async Task> ReceiptsOfAsync(World world) + { + using var scope = _fixture.BeginScope(); + + return await scope.Resolve().AgentRunCleanupReceipt.AsNoTracking() + .Where(row => row.AgentRunId == world.AgentRunId).Select(row => row.Id).ToListAsync(); + } + + /// The run's receipts as the Room's own reader hands them to its fold. + private async Task> RunCleanupReceiptsAsync(World world) + { + using var scope = _fixture.BeginScope(); + + return await scope.Resolve() + .ForRunsAsync(world.TeamId, [world.AgentRunId], CancellationToken.None); + } + + // ── The world ────────────────────────────────────────────────────────────────────────────────────────────────── + + private async Task SeedWorldAsync() + { + var actorId = Guid.NewGuid(); + var teamId = Guid.NewGuid(); + var runId = Guid.NewGuid(); + var now = DateTimeOffset.UtcNow; + using var scope = _fixture.BeginScope(); + var db = scope.Resolve(); + db.User.Add(new User { Id = actorId, Email = $"record-retention-{actorId:N}@test.local", Name = "Record Retention" }); + db.Team.Add(new Team { Id = teamId, Slug = $"record-retention-{teamId:N}", Name = "Record Retention", Kind = TeamKind.Workspace }); + db.TeamMembership.Add(new TeamMembership { Id = Guid.NewGuid(), TeamId = teamId, UserId = actorId, Role = TeamRole.Owner }); + await db.SaveChangesAsync(); + db.AgentRun.Add(new AgentRun + { + Id = runId, TeamId = teamId, Harness = "test-harness", Status = AgentRunStatus.Failed, TaskJson = "{}", + FenceEpoch = 7, CreatedDate = now, CreatedBy = actorId, LastModifiedDate = now, LastModifiedBy = actorId, + }); + await db.SaveChangesAsync(); + _teams.Add(teamId); + + return new World(teamId, actorId, runId); + } + + /// A sealed result row for a pin to hang off — the pin's foreign key needs one, and a whole campaign is a different test's subject. + private async Task SealedResultAsync(World world, Guid groupId) + { + using var scope = _fixture.BeginScope(); + var db = scope.Resolve(); + var digest = Convert.ToHexString(System.Security.Cryptography.SHA256.HashData(groupId.ToByteArray())); + db.PairedQualificationProtocol.Add(new PairedQualificationProtocol + { + ObservationGroupId = groupId, TeamId = world.TeamId, SuiteDigest = "sha256:hidden", SuiteVersion = "sha256/retention:v1", + CodeRevision = new string('c', 40), ControlModelRowId = Guid.NewGuid(), CandidateModelRowId = Guid.NewGuid(), + StatisticsVersion = "paired-cluster-bootstrap/v2", Criterion = "Quality", SessionsPerCell = 1, + MinimumIndependentClusters = 1, MinimumStrata = 1, MinimumRequiredExecutionClusters = 1, MinimumEvaluatorHealth = 1, + MaxCostUsdPerLaunch = 3m, MinimumQualityLift = 0.05, OrderingSeed = "frozen-order", ProtocolDigest = digest, + }); + await db.SaveChangesAsync(); + const string outcome = """{"pairedCells": 1, "qualifiedForCapabilityClaim": true}"""; + const string insert = "INSERT INTO paired_qualification_result (observation_group_id, protocol_digest, evidence_digest, result_digest, statistics_version, " + + "expected_observation_count, observation_count, qualified_for_capability_claim, outcome_json, created_date, created_by, last_modified_date, last_modified_by) " + + "VALUES ({0}, {1}, {2}, {3}, 'paired-cluster-bootstrap/v2', 2, 2, TRUE, {5}, now(), {4}, now(), {4})"; + await db.Database.ExecuteSqlRawAsync(insert, [groupId, digest, Digest(groupId, "evidence"), Digest(groupId, "result"), world.ActorId, outcome]); + + return groupId; + } + + private static string Digest(Guid seed, string salt) => + Convert.ToHexString(System.Security.Cryptography.SHA256.HashData([.. seed.ToByteArray(), .. System.Text.Encoding.UTF8.GetBytes(salt)])); + + public Task InitializeAsync() => Task.CompletedTask; + + /// + /// Removes what this class staged and the sweeps did not. These rows carry no bytes and no artifact placements, so + /// nothing here can change what a bounded global sweep elsewhere in the suite sees — but a class that leaves + /// eligible rows behind still hands the next sweep work it never meant to give it. + /// + public async Task DisposeAsync() + { + foreach (var team in _teams) + { + try + { + using var scope = _fixture.BeginScope(); + await scope.Resolve().Database.ExecuteSqlRawAsync("DELETE FROM agent_run_cleanup_receipt WHERE team_id = {0}", [team]); + } + catch { /* best-effort: a row a sweep already reclaimed has nothing left to give back */ } + } + } + + private sealed record World(Guid TeamId, Guid ActorId, Guid AgentRunId); +} diff --git a/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRetentionReaperFlowTests.cs b/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRetentionReaperFlowTests.cs index f36cc891b..466c1d221 100644 --- a/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRetentionReaperFlowTests.cs +++ b/backend/tests/CodeSpace.IntegrationTests/Workflows/DurableRetentionReaperFlowTests.cs @@ -434,7 +434,7 @@ public async Task A_settlement_whose_revision_moved_writes_nothing() await ElapseQuarantineAsync(stream); // somebody else moved the row after this claim was taken var moved = await StreamAsync(stream); - var settled = await Cursor().SettleAsync(claimed, DurableRetentionDecision.Quarantine(DateTimeOffset.UtcNow.AddDays(1)), CancellationToken.None); + var settled = await Cursor().SettleAsync(Window(), claimed, DurableRetentionDecision.Quarantine(DateTimeOffset.UtcNow.AddDays(1)), CancellationToken.None); settled.ShouldBeFalse("the claim is stale, so it settles nothing"); var after = await StreamAsync(stream); @@ -945,15 +945,19 @@ private LogStreamRetentionCursor Cursor() return new LogStreamRetentionCursor(scope.Resolve>(), scope.Resolve(), NullLogger.Instance); } - /// What the production claim query returns for the window the loop would build right now. - private async Task> CandidatesAsync(int limit) + /// The window the loop would build right now, from this class's own rule. + private static DurableRetentionSweepWindow Window() { var rule = DurableRetentionPolicy.LogStream; var now = DateTimeOffset.UtcNow; - return await Cursor().ClaimAsync(new DurableRetentionSweepWindow(now, now - rule.MinimumAge, now - rule.RecheckInterval), limit, CancellationToken.None); + return new DurableRetentionSweepWindow(now, now - rule.MinimumAge, now - rule.RecheckInterval, now + rule.RecheckInterval); } + /// What the production claim query returns for that window. + private async Task> CandidatesAsync(int limit) => + await Cursor().ClaimAsync(Window(), limit, CancellationToken.None); + private async Task> ClaimAsync(int limit) => (await CandidatesAsync(limit)).Select(candidate => candidate.Id).ToList(); private async Task StreamAsync(Guid streamId) diff --git a/backend/tests/CodeSpace.UnitTests/Workflows/Retention/DurableRetentionPolicyTests.cs b/backend/tests/CodeSpace.UnitTests/Workflows/Retention/DurableRetentionPolicyTests.cs index 8a3be8511..a66c766c2 100644 --- a/backend/tests/CodeSpace.UnitTests/Workflows/Retention/DurableRetentionPolicyTests.cs +++ b/backend/tests/CodeSpace.UnitTests/Workflows/Retention/DurableRetentionPolicyTests.cs @@ -24,25 +24,27 @@ public sealed class DurableRetentionPolicyTests private static readonly DateTimeOffset Now = new(2026, 9, 18, 12, 0, 0, TimeSpan.Zero); private static readonly DurableRetentionRule Rule = DurableRetentionPolicy.LogStream; - [Fact] - public void The_committed_rule_table_is_pinned_to_its_literal_windows() + [Theory] + [InlineData(DurableRecordClass.LogStream, 30)] + [InlineData(DurableRecordClass.CleanupReceipt, 30)] + public void The_committed_rule_table_is_pinned_to_its_literal_windows(DurableRecordClass value, int minimumAgeDays) { - var rule = DurableRetentionPolicy.For(DurableRecordClass.LogStream).ShouldNotBeNull(); + var rule = DurableRetentionPolicy.For(value).ShouldNotBeNull(); - rule.MinimumAge.ShouldBe(TimeSpan.FromDays(30), "a log stream is not a candidate until a month after its capture settled"); + rule.MinimumAge.ShouldBe(TimeSpan.FromDays(minimumAgeDays), "the age floor is measured from the record's own terminal instant"); rule.QuarantineWindow.ShouldBe(TimeSpan.FromHours(24), "a second, independent day passes after the first uncited observation"); - rule.RecheckInterval.ShouldBe(TimeSpan.FromHours(24), "a stream that was looked at and kept is left alone this long, so one unreclaimable row cannot own a batch slot"); + rule.RecheckInterval.ShouldBe(TimeSpan.FromHours(24), "a record that was looked at and kept is left alone this long, so one unreclaimable row cannot own a batch slot"); } [Fact] public void Every_declared_class_has_a_rule_and_the_table_declares_nothing_else() { // The table advertises what is actually reclaimed. A class listed here without a cursor would read as a - // promise the system does not keep; a cursor whose class is missing claims nothing at all. Either way the - // right time to notice is here. + // promise the system does not keep; a cursor whose class is missing claims nothing at all — the loop warns + // and moves on, which is a plane that is silently dead. Either way the right time to notice is here. DurableRetentionPolicy.Rules.Keys.Order().ShouldBe(Enum.GetValues().Order(), customMessage: "a class with no rule is never claimed, so adding one to the enum without a rule silently disables its plane"); - Enum.GetValues().ShouldBe([DurableRecordClass.LogStream], + Enum.GetValues().ShouldBe([DurableRecordClass.LogStream, DurableRecordClass.CleanupReceipt], customMessage: "a class belongs here only together with the cursor that sweeps it — add both in one change, never the rule first"); } diff --git a/backend/tests/CodeSpace.UnitTests/Workflows/Retention/RecordCitationSitesTests.cs b/backend/tests/CodeSpace.UnitTests/Workflows/Retention/RecordCitationSitesTests.cs new file mode 100644 index 000000000..61fcf0c12 --- /dev/null +++ b/backend/tests/CodeSpace.UnitTests/Workflows/Retention/RecordCitationSitesTests.cs @@ -0,0 +1,225 @@ +using CodeSpace.Core.Persistence.Db; +using CodeSpace.Core.Persistence.Entities; +using CodeSpace.Core.Services.Workflows.Retention.Cursors; +using CodeSpace.Messages.Agents.Recovery; +using Microsoft.EntityFrameworkCore; +using Shouldly; + +namespace CodeSpace.UnitTests.Workflows.Retention; + +/// +/// What may still cite a settled cleanup receipt, which of its outcomes a sweep may touch at all, and whether the +/// cursor's raw SQL says the same thing as its C#. The list is the correctness of the cursor, exactly as +/// ArtifactReferenceOracle.ReferenceSites is for artifacts: a citer missing from it makes the cursor delete a +/// row something still points at. +/// +/// Three checks keep the list from being decoration — it is compared against the EF model (a new column that +/// names a receipt reds even though nobody thought about retention while adding it), against the cursor's own +/// classification body (an entry with no probe beside it reds), and against the SQL the claim actually runs (two +/// copies of one rule drift unobserved). +/// +[Trait("Category", "Unit")] +public sealed class RecordCitationSitesTests +{ + private const string UnreachableDatabase = "Host=127.0.0.1;Port=1;Database=unused;Username=unused;Password=unused"; + + [Fact] + public void Every_citer_of_a_settled_cleanup_receipt_is_pinned() + { + CleanupReceiptRetentionCursor.CitationSites.ShouldBe(new[] + { + // A sealed qualification result still cites the receipt as part of the evidence it was computed from. + // Nothing else names a receipt by id: the Room reads them by RUN and folds them, and the orphan sweep + // selects only Orphaned rows, which this cursor never claims. + ("paired_qualification_result_pin", "pinned_id"), + }, ignoreOrder: true); + } + + /// + /// The outcomes a sweep may reclaim, held against the ledger's OWN definition of settled rather than a second + /// opinion about it. Orphaned is addressed to a sweep on the host that owes the work and can still become + /// Compensated. Unknown is rewritable — the ledger's upsert fences only the two settled outcomes, so a live leak + /// can move Orphaned → Unknown — and it is READ by the Room's recovery card with no age bound at all. + /// + [Fact] + public void A_sweep_may_reclaim_exactly_the_outcomes_the_ledger_itself_calls_settled() + { + CleanupReceiptRetentionCursor.SettledOutcomes.ShouldBe([RunResourceOutcome.Completed, RunResourceOutcome.Compensated], ignoreOrder: true); + + foreach (var outcome in Enum.GetValues()) + CleanupReceiptRetentionCursor.SettledOutcomes.Contains(outcome).ShouldBe(Receipt(outcome).IsSettled, + $"'{outcome}' disagrees with RunCleanupReceipt.IsSettled — the plane that WRITES receipts already says which ones are over"); + } + + [Fact] + public void Every_pinned_citation_site_is_actually_probed_by_the_cursor() + { + using var db = BuildContext(); + + // ClassifyAsync's own body, never the whole file: the deleting statement repeats the same probe, so a wider + // match would keep passing with the classification's probe removed — and the classification is what the loop + // asks before anything is deleted. + var body = MethodBody(Source(), "ClassifyAsync"); + var probes = ExistenceQuestions(body); + + probes.ShouldNotBeNullOrWhiteSpace("the classification asks no existence question at all, so this check would pass by reading nothing"); + body.ShouldNotContain("SettleAsync(", customMessage: + "the extract ran past the method it was scoped to; a delimiter that stops only at private members walks through every public one after it, " + + "and a probe living there would satisfy this check without the classification ever asking"); + + var unprobed = CleanupReceiptRetentionCursor.CitationSites + .Select(site => (site.Table, site.Column, Member: MemberOf(db, site.Table, site.Column))) + .Where(site => !probes.Contains($".{site.Member}", StringComparison.Ordinal)) + .Select(site => $"{site.Table}.{site.Column} (no use of .{site.Member} in the citation probe)") + .ToList(); + + unprobed.ShouldBeEmpty("a site the cursor lists but never asks about is a claim it does not keep:\n " + string.Join("\n ", unprobed)); + } + + /// + /// The drift detector, over the EF model rather than the list: a future column that names a cleanup receipt reds + /// here even though nobody thought about retention while adding it. + /// + [Fact] + public void A_mapped_column_that_names_a_cleanup_receipt_and_is_not_probed_fails_this_test() + { + using var db = BuildContext(); + var probed = CleanupReceiptRetentionCursor.CitationSites.Select(site => $"{site.Table}.{site.Column}").ToHashSet(StringComparer.Ordinal); + + var mapped = db.Model.GetEntityTypes() + .SelectMany(entity => entity.GetProperties().Select(property => (Table: entity.GetTableName(), Column: property.GetColumnName()))) + .Where(column => column.Table is not null && column.Column is not null && NamesACleanupReceipt(column.Table!, column.Column!)) + .Select(column => $"{column.Table}.{column.Column}") + .Distinct(StringComparer.Ordinal) + .Order(StringComparer.Ordinal) + .ToList(); + + mapped.Where(column => !probed.Contains(column)).ShouldBeEmpty( + $"a column that names a cleanup receipt which {nameof(CleanupReceiptRetentionCursor)} does not probe would let it delete a row that column still points at — " + + $"add it to {nameof(CleanupReceiptRetentionCursor.CitationSites)} and to the probe beside it"); + } + + /// + /// The claim's raw SQL and the cursor's C# allow-lists are two copies of one rule, and two copies drift + /// unobserved. The claim cannot parameterise an IN list of enum names, so this reads the cursor's own + /// source and compares what its SQL says with what its C# says. + /// + [Fact] + public void The_claims_raw_sql_says_exactly_what_the_allow_lists_say() + { + var expected = string.Join(", ", CleanupReceiptRetentionCursor.SettledOutcomes.Select(outcome => $"'{outcome}'")); + + CleanupReceiptRetentionCursor.SettledOutcomeNames.ShouldBe(expected, "the SQL literal and the enum allow-list are one rule written twice"); + CleanupReceiptRetentionCursor.PinnedKindName.ShouldBe(DurablePinKind.CleanupReceipt.ToString()); + + // Three copies, not two: the claim, the delete, and migration 0237's PARTIAL INDEX, whose predicate has to + // match the claim's or the index quietly drops out of the plan and the sweep starts seq-scanning. + var inLists = OutcomeLists(Source()).Concat(OutcomeLists(Migration())).ToList(); + var pinKinds = System.Text.RegularExpressions.Regex.Matches(Source(), @"pin\.kind = '([^']*)'").Select(match => match.Groups[1].Value).ToList(); + + inLists.Count.ShouldBeGreaterThanOrEqualTo(3, "the claim, the delete and the index predicate all name the settled set; finding fewer means this check read nothing"); + pinKinds.ShouldNotBeEmpty("the claim's pin filter was not found, so this check would pass by reading nothing"); + inLists.ShouldAllBe(list => list == CleanupReceiptRetentionCursor.SettledOutcomeNames, "every outcome list — in the cursor AND in the migration — is the settled set, exactly"); + pinKinds.ShouldAllBe(kind => kind == CleanupReceiptRetentionCursor.PinnedKindName, "the pin kind the SQL names is the one the enum spells"); + } + + private static IEnumerable OutcomeLists(string sql) => + System.Text.RegularExpressions.Regex.Matches(sql, @"outcome IN \(([^)]*)\)").Select(match => match.Groups[1].Value.Trim()); + + /// Migration 0237, read from the scripts that ship with the build rather than copied here. + private static string Migration() => + File.ReadAllText(Path.Combine(AppContext.BaseDirectory, "Persistence", "DbUpFiles", "0237_durable_retention_record_deadlines.sql")); + + /// + /// A column names a cleanup receipt when it ends in cleanup_receipt_id, or when it is the generic pin + /// column the retention plane uses. Near misses, stated rather than left to a reader to wonder about: the + /// receipt's own agent_run_id and team_id point the other way — the receipt names THEM — and + /// workflow_run_model_call_attempt.budget_reservation_id names a different plane's record entirely. + /// + private static bool NamesACleanupReceipt(string table, string column) => + (column.EndsWith("cleanup_receipt_id", StringComparison.Ordinal) && table != "agent_run_cleanup_receipt") + || (table == "paired_qualification_result_pin" && column == "pinned_id"); + + private static RunCleanupReceipt Receipt(RunResourceOutcome outcome) => new() + { + AgentRunId = Guid.NewGuid(), FenceEpoch = 1, Kind = RunResourceKind.Spool, Outcome = outcome, + RecordedByHost = "test", RecordedAt = DateTimeOffset.UtcNow, + }; + + private static string Source() => + File.ReadAllText(Path.Combine(ProductionSourceRoot(), "CodeSpace.Core", "Services", "Workflows", "Retention", "Cursors", $"{nameof(CleanupReceiptRetentionCursor)}.cs")); + + /// The arguments of every Any/AnyAsync in the text — the "does anything still name this" questions, which is what makes a listed table a citation site. + private static string ExistenceQuestions(string source) + { + var questions = new System.Text.StringBuilder(); + + foreach (var call in new[] { "AnyAsync(", ".Any(" }) + { + for (var index = source.IndexOf(call, StringComparison.Ordinal); index >= 0; index = source.IndexOf(call, index + 1, StringComparison.Ordinal)) + { + var end = source.IndexOf('\n', index); + + questions.Append(source[index..(end < 0 ? source.Length : end)]); + } + } + + return questions.ToString(); + } + + /// + /// One method's text, from its declaration to the next member at class indentation — of ANY visibility, because + /// delimiting on private members alone runs a probe's extract on through every public member after it, which is + /// how a scoped check quietly becomes a whole-file one again. + /// + private static string MethodBody(string source, string name) + { + var start = DeclarationOf(source, name); + + if (start < 0) return string.Empty; + + var next = new[] { "\n private ", "\n public ", "\n internal ", "\n protected " } + .Select(member => source.IndexOf(member, start + 1, StringComparison.Ordinal)) + .Where(index => index >= 0) + .DefaultIfEmpty(source.Length) + .Min(); + + return source[start..next]; + } + + /// The DECLARATION of a method, not the first mention of it: a probe is called from elsewhere in the file, and a search that stopped there would read the caller instead. + private static int DeclarationOf(string source, string name) + { + for (var index = source.IndexOf($"{name}(", StringComparison.Ordinal); index >= 0; index = source.IndexOf($"{name}(", index + 1, StringComparison.Ordinal)) + { + var line = source[(source.LastIndexOf('\n', index) + 1)..index]; + + if (line.Contains("public", StringComparison.Ordinal) || line.Contains("private", StringComparison.Ordinal)) return index; + } + + return -1; + } + + private static string MemberOf(CodeSpaceDbContext db, string table, string column) + { + var entity = db.Model.GetEntityTypes().FirstOrDefault(type => type.GetTableName() == table) + ?? throw new InvalidOperationException($"'{table}' is not mapped, so a citation site names a table this build does not have."); + + return (entity.GetProperties().FirstOrDefault(property => property.GetColumnName() == column) + ?? throw new InvalidOperationException($"'{table}.{column}' is not mapped, so a citation site names a column this build does not have.")).Name; + } + + private static string ProductionSourceRoot() + { + for (var directory = new DirectoryInfo(AppContext.BaseDirectory); directory is not null; directory = directory.Parent) + { + var candidate = Path.Combine(directory.FullName, "backend", "src"); + if (Directory.Exists(candidate)) return candidate; + } + + throw new DirectoryNotFoundException($"'backend/src' was not found above '{AppContext.BaseDirectory}', so the probes were never checked. Run the unit suite from the repository checkout."); + } + + private static CodeSpaceDbContext BuildContext() => + new(new DbContextOptionsBuilder().UseNpgsql(UnreachableDatabase).UseSnakeCaseNamingConvention().Options); +}