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 @@ -235,14 +235,20 @@ private static async Task<IReadOnlyList<SegmentPlacement>> ReclaimableAsync(Code

/// <summary>
/// Nothing is left to remove, so either this plane's own lifecycle emptied the stream or something else did. Only
/// the first may be tombstoned: an object whose locations never reached <c>Purged</c> is a loss this cursor did
/// not cause and must not claim. Counting the unaccounted objects over ALL of them, not the batch, is what keeps
/// a stream with more objects than one sweep can hold from being declared drained on a partial view.
/// the first may be tombstoned: EVERY placement of every object has to rest at <c>Purged</c>, which is the state
/// only this lifecycle writes. One <c>Purged</c> placement beside a <c>Deleted</c> one is a deduplicated object
/// half of which left by a path this cursor did not drive — and "some of it was reclaimed by policy" is not a
/// statement the tombstone can make, because a reader takes it to cover the whole stream.
///
/// <para>Counting over ALL objects, not the batch, is what keeps a stream with more objects than one sweep can
/// hold from being declared drained on a partial view.</para>
/// </summary>
private async Task<DrainOutcome> DrainedOrLostAsync(CodeSpaceDbContext db, DurableRetentionCandidate candidate, IQueryable<Guid> objects, CancellationToken cancellationToken)
{
var unaccounted = await objects.CountAsync(objectId => !db.ArtifactLocation
.Any(location => location.TeamId == candidate.TeamId && location.ArtifactObjectId == objectId && location.State == ArtifactLocationState.Purged), cancellationToken).ConfigureAwait(false);
var unaccounted = await objects.CountAsync(objectId =>
!db.ArtifactLocation.Any(location => location.TeamId == candidate.TeamId && location.ArtifactObjectId == objectId)
|| db.ArtifactLocation.Any(location => location.TeamId == candidate.TeamId && location.ArtifactObjectId == objectId && location.State != ArtifactLocationState.Purged),
cancellationToken).ConfigureAwait(false);

return unaccounted == 0 ? DrainOutcome.Drained : BytesAlreadyGone(candidate, unaccounted);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
using CodeSpace.IntegrationTests.Infrastructure;
using CodeSpace.Messages.Agents.Benchmark;
using CodeSpace.Messages.Artifacts;
using CodeSpace.Messages.Constants;
using CodeSpace.Messages.Dtos.Sessions.Room;
using CodeSpace.Messages.Enums;
using CodeSpace.Messages.Retention;
Expand Down Expand Up @@ -465,6 +466,55 @@ public async Task A_stream_whose_bytes_vanished_outside_this_plane_is_never_tomb
"\"purged; retention window elapsed\" about a loss nobody chose inverts the one distinction the tombstone preserves");
}

/// <summary>
/// Who the reclamation is attributed to, and which placement it names — read off the requests the cursor actually
/// made, not off the constant it was supposed to use. Mutation: drop the actor (or the location id) from the
/// request and this reds, where a test that only pinned <c>SystemUsers.SeederId</c>'s value would not notice.
/// </summary>
[Fact]
public async Task Every_purge_this_plane_makes_names_its_placement_and_the_system_actor()
{
var world = await SeedWorldAsync();
var stream = await CaptureAsync(world, "bytes with an audit trail", segments: 2);
await AgeAsync(stream, TimeSpan.FromDays(31));
await SweepAsync();
await ElapseQuarantineAsync(stream);
var placements = await PlacementsOfAsync(world, stream);

var requests = await SweepRecordingPurgesAsync();

requests.Select(request => request.ArtifactLocationId ?? Guid.Empty).ShouldBe(placements, ignoreOrder: true,
customMessage: "an unnamed claim refuses an object with more than one placement, so every request has to say WHICH");
requests.ShouldAllBe(request => request.ActorId == SystemUsers.SeederId);
requests.ShouldAllBe(request => request.TeamId == world.TeamId);
(await StreamAsync(stream)).PurgedAt.ShouldNotBeNull();
}

/// <summary>
/// A deduplicated object half of which left by another path. One placement this plane purged, one another path
/// deleted — and "reclaimed by policy" is not a statement that can cover only half a stream, because a reader
/// takes the tombstone to cover all of it. Mutation: treat ANY Purged placement as drained and this reds with a
/// policy tombstone on a partial loss.
/// </summary>
[Fact]
public async Task A_stream_one_of_whose_placements_left_by_another_path_is_never_tombstoned()
{
var world = await SeedWorldAsync();
var stream = await CaptureAsync(world, "identical", segments: 2, distinctSegments: false);
var placements = await PlacementsOfAsync(world, stream);
placements.Count.ShouldBe(2, "the premise: one object, one placement per write");

await PurgeOnePlacementAsync(world, stream, placements[0]);
await MarkPlacementDeletedAsync(placements[1]);
await AgeAsync(stream, TimeSpan.FromDays(31));
await SweepAsync();
await ElapseQuarantineAsync(stream);
await SweepAsync();

(await StreamAsync(stream)).PurgedAt.ShouldBeNull(
"one placement this plane purged beside one it did not is a partial loss, and the tombstone would report it as a complete policy reclamation");
}

/// <summary>
/// Migration 0236's arm, at the only place that can pin it: the database. A terminal stream admits a retention
/// statement and NOTHING else — a statement that smuggles anything alongside the two columns reads the same
Expand Down Expand Up @@ -842,6 +892,40 @@ private async Task<DurableRetentionSweepSummary> SweepAgainstARefusingDestinatio
return await new DurableRetentionReaper(options, [cursor], NullLogger<DurableRetentionReaper>.Instance).SweepAsync(CancellationToken.None);
}

/// <summary>The real reaper over the real cursor and the real coordinator, with every purge request it makes recorded on the way through.</summary>
private async Task<IReadOnlyList<ArtifactCasPurgeRequest>> SweepRecordingPurgesAsync()
{
using var scope = _fixture.BeginScope();
var options = scope.Resolve<DbContextOptions<CodeSpaceDbContext>>();
var recorder = new RecordingPurgeCoordinator(scope.Resolve<IArtifactCasPurgeCoordinator>());
var cursor = new LogStreamRetentionCursor(options, recorder, NullLogger<LogStreamRetentionCursor>.Instance);
await new DurableRetentionReaper(options, [cursor], NullLogger<DurableRetentionReaper>.Instance).SweepAsync(CancellationToken.None);

return recorder.Requests;
}

/// <summary>Passes every call through to the real coordinator and keeps what was asked, so the audit facts are read off the requests the cursor made rather than off the constant it was meant to use.</summary>
private sealed class RecordingPurgeCoordinator : IArtifactCasPurgeCoordinator
{
private readonly IArtifactCasPurgeCoordinator _inner;

public RecordingPurgeCoordinator(IArtifactCasPurgeCoordinator inner) { _inner = inner; }

public List<ArtifactCasPurgeRequest> Requests { get; } = [];

public Task<ArtifactCasPurgeResult> PurgeAsync(ArtifactCasPurgeRequest request, CancellationToken cancellationToken)
{
Requests.Add(request);

return _inner.PurgeAsync(request, cancellationToken);
}

public Task<ArtifactCasPurgeClaimResult> ClaimAsync(ArtifactCasPurgeRequest request, CancellationToken cancellationToken) => _inner.ClaimAsync(request, cancellationToken);
public Task<ArtifactCasPurgeResult> DeleteAsync(ArtifactCasPurgeClaim claim, CancellationToken cancellationToken) => _inner.DeleteAsync(claim, cancellationToken);
public Task<ArtifactCasReleaseOutcome> ReleaseAsync(ArtifactCasPurgeClaim claim, ArtifactCasReleaseEvidence evidence, CancellationToken cancellationToken) => _inner.ReleaseAsync(claim, evidence, cancellationToken);
public Task<ArtifactCasAbandonResult> AbandonAsync(ArtifactCasPurgeClaim claim, CancellationToken cancellationToken) => _inner.AbandonAsync(claim, cancellationToken);
}

/// <summary>A destination that answers every delete with a refusal and no effect — the shape a revoked key or a read-only bucket produces.</summary>
private sealed class RefusingPurgeCoordinator : IArtifactCasPurgeCoordinator
{
Expand Down Expand Up @@ -897,6 +981,52 @@ private async Task<IReadOnlyList<string>> InsertingTransactionsAsync(Guid result
+ "UNION SELECT DISTINCT xmin::text FROM paired_qualification_result_pin WHERE result_id = {0}", resultId).ToListAsync();
}

/// <summary>Every placement still holding bytes for this stream's objects — the rows a drain has to name one by one.</summary>
private async Task<IReadOnlyList<Guid>> PlacementsOfAsync(World world, Guid streamId)
{
using var scope = _fixture.BeginScope();
var db = scope.Resolve<CodeSpaceDbContext>();
var objects = await ObjectsOfAsync(world, streamId);

return await db.ArtifactLocation.AsNoTracking()
.Where(location => location.TeamId == world.TeamId && objects.Contains(location.ArtifactObjectId)
&& location.State != ArtifactLocationState.Purged && location.State != ArtifactLocationState.Deleted)
.OrderBy(location => location.Id).Select(location => location.Id).ToListAsync();
}

/// <summary>Drives ONE placement through the real purge lifecycle, leaving its siblings alone.</summary>
private async Task PurgeOnePlacementAsync(World world, Guid streamId, Guid locationId)
{
using var scope = _fixture.BeginScope();
var objectId = (await ObjectsOfAsync(world, streamId)).Single();
var outcome = await scope.Resolve<IArtifactCasPurgeCoordinator>().PurgeAsync(new ArtifactCasPurgeRequest
{
TeamId = world.TeamId, ArtifactObjectId = objectId, ArtifactLocationId = locationId, ActorId = SystemUsers.SeederId,
}, CancellationToken.None);

outcome.ShouldBeOfType<ArtifactCasPurgeResult.Purged>();
}

/// <summary>Takes one placement out of the lifecycle by a path this plane never drives. Guards suspended, exactly as the stream's are for time travel.</summary>
private async Task MarkPlacementDeletedAsync(Guid locationId)
{
using var scope = _fixture.BeginScope();
var db = scope.Resolve<CodeSpaceDbContext>();
await db.Database.ExecuteSqlRawAsync("ALTER TABLE artifact_location DISABLE TRIGGER USER");

try
{
await db.ArtifactLocation.Where(location => location.Id == locationId)
.ExecuteUpdateAsync(set => set
.SetProperty(location => location.State, ArtifactLocationState.Deleted)
.SetProperty(location => location.Revision, location => location.Revision + 1));
}
finally
{
await db.Database.ExecuteSqlRawAsync("ALTER TABLE artifact_location ENABLE TRIGGER USER");
}
}

private async Task<IReadOnlyList<Guid>> ObjectsOfAsync(World world, Guid streamId)
{
using var scope = _fixture.BeginScope();
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,10 @@
using CodeSpace.Core.Services.Workflows.Retention;
using CodeSpace.Core.Handlers.QueryHandlers.Agents;
using CodeSpace.Core.Persistence.Entities;
using CodeSpace.Core.Services.Agents.AgentRunLogging;
using CodeSpace.Messages.Artifacts;
using CodeSpace.Messages.Constants;
using CodeSpace.Messages.Dtos.Agents;
using CodeSpace.Messages.Retention;
using Shouldly;

Expand Down Expand Up @@ -41,6 +46,26 @@ public void Every_declared_class_has_a_rule_and_the_table_declares_nothing_else(
customMessage: "a class belongs here only together with the cursor that sweeps it — add both in one change, never the rule first");
}

/// <summary>
/// The exact strings a purged stream puts on the wire. The web client matches both literally — see
/// <c>frontend/src/api/agentRunLogsApi.test.ts</c>, which feeds this availability and this code through its own
/// decoder — and a client that does not recognise them reports the RESPONSE as malformed, which is a worse answer
/// than the wrong one it replaced. Renaming the enum member changes both, so it has to fail here first.
/// </summary>
[Fact]
public void A_purged_read_puts_its_availability_and_code_on_the_wire_verbatim()
{
var read = AgentRunLogWire.Unavailable(Metadata(), 0, new AgentRunLogProblem(AgentRunLogProblemCode.Purged));

read.Availability.ShouldBe(AgentRunLogReadAvailability.Purged);
read.Availability.ToString().ShouldBe("Purged");
read.ProblemCode.ShouldBe("Purged", "the controller sends ProblemCode as the response's `code`, so this literal IS the web client's contract");
read.IsRetryable.ShouldBeFalse("bytes reclaimed on purpose never come back, so no caller should retry the read");
}

private static AgentRunLogMetadata Metadata() => new(Guid.NewGuid(), Guid.NewGuid(), "stdout/v1", "text/plain", "utf-8", "spool/v1",
ArtifactRetention.Run, AgentRunLogStreamState.Completed, 2, 1, 3, 3, null, Now, Now, Now, null);

/// <summary>
/// Who a reclamation is attributed to. The seeder is the established identity for background work — Hangfire
/// workers, scheduled jobs and DbUp all write under it (<c>SystemUsers.cs</c>), and it holds a real seeded row, so
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -69,17 +69,72 @@ public void Every_pinned_citation_site_is_actually_probed_by_the_cursor()
using var db = BuildContext();
var source = File.ReadAllText(Path.Combine(ProductionSourceRoot(), "CodeSpace.Core", "Services", "Workflows", "Retention", "Cursors", $"{nameof(LogStreamRetentionCursor)}.cs"));

// The EXISTENCE QUESTIONS only — not the whole file, and not even the whole probe method. The drain reads
// ArtifactObjectId in three other places, and the probe method itself reads it once more to name the stream's
// own objects, so any wider match keeps passing with the sibling-segment probe deleted: the exact hole this
// test is for. A site is probed when something ASKS whether a row still names it.
var probes = ExistenceQuestions(MethodBody(source, "IsPinnedAsync") + MethodBody(source, "SharesBytesAsync"));
probes.ShouldNotBeNullOrWhiteSpace("no existence question was found in the probe methods, so this test would pass by reading nothing at all");

var unprobed = LogStreamRetentionCursor.CitationSites
.Select(site => (site.Table, site.Column, Member: MemberOf(db, site.Table, site.Column)))
.Where(site => !source.Contains($".{site.Member}", StringComparison.Ordinal))
.Select(site => $"{site.Table}.{site.Column} (no use of .{site.Member})")
.Where(site => !probes.Contains($".{site.Member}", StringComparison.Ordinal))
.Select(site => $"{site.Table}.{site.Column} (no use of .{site.Member} in the citation probes)")
.ToList();

unprobed.ShouldBeEmpty(
$"a site listed in {nameof(LogStreamRetentionCursor.CitationSites)} that nothing reads is a claim the cursor does not keep — "
$"a site listed in {nameof(LogStreamRetentionCursor.CitationSites)} that the citation probes never read is a claim the cursor does not keep — "
+ "add the probe, or take the entry out:\n " + string.Join("\n ", unprobed));
}

/// <summary>
/// The arguments of every <c>AnyAsync</c> in the probes — each one an "does anything still name this" question,
/// which is what makes a listed table a citation SITE rather than a table the cursor happens to read. Each call in
/// these methods ends at its cancellation token, so that is where the span stops.
/// </summary>
private static string ExistenceQuestions(string source)
{
var questions = new System.Text.StringBuilder();

for (var index = source.IndexOf("AnyAsync(", StringComparison.Ordinal); index >= 0; index = source.IndexOf("AnyAsync(", index + 1, StringComparison.Ordinal))
{
var end = source.IndexOf("cancellationToken", index, StringComparison.Ordinal);

questions.Append(source[index..(end < 0 ? source.Length : end)]);
}

return questions.ToString();
}

/// <summary>
/// One method's text, from its declaration to the next member at class indentation. Crude on purpose: it has to
/// hold for an expression-bodied probe and a braced one alike, and anything cleverer would be a parser this test
/// would then depend on being right.
/// </summary>
private static string MethodBody(string source, string name)
{
var start = DeclarationOf(source, name);

if (start < 0) return string.Empty;

var next = source.IndexOf("\n private ", start + 1, StringComparison.Ordinal);

return source[start..(next < 0 ? source.Length : next)];
}

/// <summary>The DECLARATION of a method, not the first mention of it: both probes are called from <c>ClassifyAsync</c> higher up the file, and a search that stopped there would read the caller instead of the probe.</summary>
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 lineStart = source.LastIndexOf('\n', index) + 1;

if (source[lineStart..index].Contains("private", StringComparison.Ordinal)) return index;
}

return -1;
}

/// <summary>The CLR property a column is mapped from — the name the cursor's LINQ has to mention if it reads that column at all.</summary>
private static string MemberOf(CodeSpaceDbContext db, string table, string column)
{
Expand Down
13 changes: 8 additions & 5 deletions frontend/src/api/agentRunLogsApi.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,11 +76,14 @@ describe("Agent Run durable log API", () => {
// wrong-but-well-formed answer, and it hides a healthy deployment doing exactly what its retention policy says.
it("keeps the retention plane's Purged verdict distinct from a missing object", async () => {
vi.stubGlobal("fetch", vi.fn()
.mockResolvedValueOnce(json({ availability: "Purged", code: "log_bytes_purged", isRetryable: false, streamId: "s" }, 410))
.mockResolvedValueOnce(json({ availability: "PhysicalObjectMissing", code: "artifact_missing", isRetryable: false, streamId: "s" }, 410)));

await expect(agentsApi.readRunLogRange("r", "s", 0, 1)).resolves.toEqual({ availability: "Purged", code: "log_bytes_purged", isRetryable: false });
await expect(agentsApi.readRunLogRange("r", "s", 0, 1)).resolves.toEqual({ availability: "PhysicalObjectMissing", code: "artifact_missing", isRetryable: false });
.mockResolvedValueOnce(json({ availability: "Purged", code: "Purged", isRetryable: false, streamId: "s" }, 410))
.mockResolvedValueOnce(json({ availability: "PhysicalObjectMissing", code: "ArtifactMissing", isRetryable: false, streamId: "s" }, 410)));

// The codes are the backend's own AgentRunLogProblemCode names, which the controller sends verbatim
// (`ProblemCode = problem.Code.ToString()`); the pairing is pinned on that side by
// DurableRetentionPolicyTests.A_purged_read_puts_its_availability_and_code_on_the_wire_verbatim.
await expect(agentsApi.readRunLogRange("r", "s", 0, 1)).resolves.toEqual({ availability: "Purged", code: "Purged", isRetryable: false });
await expect(agentsApi.readRunLogRange("r", "s", 0, 1)).resolves.toEqual({ availability: "PhysicalObjectMissing", code: "ArtifactMissing", isRetryable: false });
});

it("fails closed when a success response omits or contradicts its range contract", async () => {
Expand Down
Loading