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 @@ -4765,10 +4765,10 @@
/// <summary>The same reconstruction from a payload that came from somewhere other than the row — an offloaded one fetched back out of the artifact store.</summary>
private static AgentEvent ReplayedEvent(AgentEventKind kind, string? text, string? dataJson)
{
if (dataJson is not { Length: > 0 } json) return new AgentEvent { Kind = kind, Text = text };

Check warning on line 4768 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / recurring jobs fire (worker host · Postgres)

Possible null reference assignment.

Check warning on line 4768 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (E2ETests · HTTP · Postgres)

Possible null reference assignment.

Check warning on line 4768 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (UnitTests)

Possible null reference assignment.

Check warning on line 4768 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (UnitTests)

Possible null reference assignment.

Check warning on line 4768 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (LargeCapture · 4 GiB · Postgres)

Possible null reference assignment.

Check warning on line 4768 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (IntegrationTests · Postgres)

Possible null reference assignment.

try { using var doc = JsonDocument.Parse(json); return new AgentEvent { Kind = kind, Text = text, Data = doc.RootElement.Clone() }; }

Check warning on line 4770 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / recurring jobs fire (worker host · Postgres)

Possible null reference assignment.

Check warning on line 4770 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (E2ETests · HTTP · Postgres)

Possible null reference assignment.

Check warning on line 4770 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (UnitTests)

Possible null reference assignment.

Check warning on line 4770 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (UnitTests)

Possible null reference assignment.

Check warning on line 4770 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (LargeCapture · 4 GiB · Postgres)

Possible null reference assignment.

Check warning on line 4770 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (IntegrationTests · Postgres)

Possible null reference assignment.
catch (JsonException) { return new AgentEvent { Kind = kind, Text = text }; }

Check warning on line 4771 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / recurring jobs fire (worker host · Postgres)

Possible null reference assignment.

Check warning on line 4771 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (E2ETests · HTTP · Postgres)

Possible null reference assignment.

Check warning on line 4771 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (UnitTests)

Possible null reference assignment.

Check warning on line 4771 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (UnitTests)

Possible null reference assignment.

Check warning on line 4771 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (LargeCapture · 4 GiB · Postgres)

Possible null reference assignment.

Check warning on line 4771 in backend/src/CodeSpace.Core/Services/Agents/AgentRunExecutor.cs

View workflow job for this annotation

GitHub Actions / dotnet test (IntegrationTests · Postgres)

Possible null reference assignment.
}

/// <summary>Ask the row, on a token of its own, whether the run actually reached a terminal state — the only honest answer to "did the landing take?" once an exception has been raised somewhere after the fenced write.</summary>
Expand Down Expand Up @@ -5408,7 +5408,7 @@
/// The checkpoint coordinates for a run whose envelope OPTED IN, else null — the one place the opt-in is read, so
/// the produce side and the consume side cannot disagree about which runs are checkpointed. A run whose failed
/// attempt nobody can retry writes nothing, which is what keeps the artifact store free of a per-minute
/// transcript copy for every benchmark cell, review child and supervisor unit on the fleet.
/// transcript copy for every benchmark cell, review child and supervisor unit that no retry could resume.
/// </summary>
private static SessionCheckpointTick? CheckpointTickFor(AgentTask task, Guid teamId, IAgentHarness harness, string? workingDirectory, AgentRunFacts facts) =>
task.CheckpointSessionTranscript ? new SessionCheckpointTick(teamId, harness, workingDirectory, facts) : null;
Expand Down
33 changes: 31 additions & 2 deletions backend/src/CodeSpace.Core/Services/Agents/AgentRunService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -774,13 +774,19 @@ public async Task<bool> CancelRunningAsync(Guid runId, string reason, AgentRunAb

if (snapshot is null || snapshot.Status != AgentRunStatus.Running) return false;

// 3c: a deliberate cancel is a CLEAN landing, so it releases the mid-run session checkpoint exactly as
// completion does. Nobody owes this run a continuation, and a kept reference would pin the artifact
// Referenced (terminal in the retention ledger) for good — and make a later retry of the same subtask read
// the cancel as a host loss. Only the reconciler's abandon keeps these columns.
var cancelled = await _db.AgentRun
.Where(r => r.Id == runId && r.Status == AgentRunStatus.Running && r.FenceEpoch == snapshot.FenceEpoch)
.ExecuteUpdateAsync(s => s
.SetProperty(r => r.Status, AgentRunStatus.Cancelled)
.SetProperty(r => r.FenceEpoch, r => r.FenceEpoch + 1)
.SetProperty(r => r.Error, reason)
.SetProperty(r => r.CompletedAt, (DateTimeOffset?)DateTimeOffset.UtcNow), cancellationToken)
.SetProperty(r => r.CompletedAt, (DateTimeOffset?)DateTimeOffset.UtcNow)
.SetProperty(r => r.SessionTranscriptCheckpointArtifactId, (Guid?)null)
.SetProperty(r => r.SessionTranscriptCheckpointAt, (DateTimeOffset?)null), cancellationToken)
.ConfigureAwait(false);

if (cancelled == 0) return false;
Expand Down Expand Up @@ -956,7 +962,7 @@ await _db.AgentRun.AsNoTracking().SingleOrDefaultAsync(r => r.Id == runId, cance
var candidates = await _db.AgentRun.AsNoTracking()
.Where(a => a.TeamId == teamId && a.WorkflowRunId == supervisorRunId && a.SessionId != null)
.OrderByDescending(a => a.CreatedDate).ThenByDescending(a => a.Id)
.Select(a => new { a.Id, a.SessionId, a.ResultJson, a.TaskJson })
.Select(a => new { a.Id, a.Status, a.SessionId, a.ResultJson, a.TaskJson, a.SessionTranscriptCheckpointArtifactId, a.SessionTranscriptCheckpointAt })
.ToListAsync(cancellationToken).ConfigureAwait(false);

// The most-recent RESUMABLE prior attempt of THIS subtask — skip a captured-but-transcript-less attempt so it
Expand All @@ -966,11 +972,34 @@ await _db.AgentRun.AsNoTracking().SingleOrDefaultAsync(r => r.Id == runId, cance
if (SubtaskIdOf(candidate.TaskJson) != subtaskId) continue;

if (TryResumable(candidate.Id, candidate.SessionId, candidate.ResultJson) is { } resumable) return resumable;

// 3c: a prior attempt whose HOST died has no result at all — the reconciler's abandon writes none — so
// the captured-transcript read above finds nothing for exactly the population a warm retry helps most.
// Its mid-run checkpoint is what survives, and a TERMINAL row still naming one is the signature of an
// abandon: every clean landing — completion or cancel — releases these two columns in its own terminal
// write; only an abandon keeps them.
if (TryResumableFromCheckpoint(candidate.Id, candidate.Status, candidate.SessionId, candidate.SessionTranscriptCheckpointArtifactId, candidate.SessionTranscriptCheckpointAt) is { } continued) return continued;
}

return null;
}

/// <summary>
/// A prior attempt's MID-RUN checkpoint as a resumable session — the host-loss counterpart of
/// <see cref="TryResumable"/>. Both halves are required for the same reason that one demands both: a session id
/// with no transcript resumes into "No conversation found", and a transcript nothing can address is not
/// resumable at all.
///
/// <para>And only from a TERMINAL row. The checkpoint columns are written while the attempt is Running, so a
/// row still Running holds the checkpoint of a session that may be live — a kill-wave is best-effort, and a
/// revived parent stops the orphan sweep from selecting it — and resuming it would fork a conversation that is
/// still being written.</para>
/// </summary>
private static ResumableSession? TryResumableFromCheckpoint(Guid agentRunId, AgentRunStatus status, string? sessionId, Guid? checkpointArtifactId, DateTimeOffset? checkpointAt) =>
AgentRunStateMachine.IsTerminal(status) && sessionId is { Length: > 0 } sid && checkpointArtifactId is { } artifactId && checkpointAt is { } at
? new ResumableSession(agentRunId, sid, null, artifactId, at)
: null;

/// <summary>
/// A prior agent run's (id, session id, result json) → its RESUMABLE session, or null when not resumable: no
/// session id, no transcript (both-or-neither — a resume without one fails "No conversation found"), or a
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,19 @@ internal static IReadOnlyList<string> DependsOnFor(SupervisorPlannedSubtask? pla
internal static DependencyStagingResult PreferPriorAttemptStaging(DependencyStagingResult priorAttemptStaging, DependencyStagingResult dependencyStaging) =>
priorAttemptStaging.Ref is not null ? priorAttemptStaging : dependencyStaging;

/// <summary>
/// 3c: whether a staged unit opts into the mid-run session checkpoint — only when a later retry could actually
/// consume it. A checkpoint costs a whole-file read and an artifact write a minute for as long as the unit runs,
/// and a host loss keeps the survivor for good, so a unit nobody can resume must not pay for one.
///
/// <para>Two exclusions. A unit with no subtask id can never be FOUND by the retry lookup, which matches on it:
/// that is every resolver unit, whose task is built without one, and a re-resolve stages a fresh resolver rather
/// than resuming the old one. And a unit the run can no longer afford to respawn has nobody to hand a restored
/// conversation to (<see cref="SupervisorBounds.CanRespawnAfterWave"/>).</para>
/// </summary>
internal static bool CheckpointsSessionTranscript(AgentTask task, SupervisorTurnContext context, int waveSize) =>
task.SubtaskId is { Length: > 0 } && SupervisorBounds.CanRespawnAfterWave(context, waveSize);

/// <summary>The subtask's target repository, resolved the SAME way <see cref="BuildTaskWithGoal"/> will resolve it — a pure pre-computation so dependency staging can look up the right repo's manifest before the task itself is built.</summary>
private static Guid? ResolveTargetRepositoryId(SupervisorAgentDispatch? spec, SupervisorTurnContext context)
{
Expand Down Expand Up @@ -433,7 +446,11 @@ private async Task<SupervisorExecution> ExecuteRetryAsync(SupervisorDecision dec

var (escalatedTask, escalation) = await ApplyRetryEscalationAsync(builtTask, priorResult, context, cancellationToken).ConfigureAwait(false);

var task = ApplyRetryDisposition(escalatedTask, prior, priorResult, workspaceHasPriorWork: effectiveStaging.Ref is not null);
// The continuity sentence describes the resumed attempt's OWN git state, so it reads that attempt's own
// pushed branch — never the effective clone ref. For a dependent unit whose attempt pushed nothing, the
// effective ref is the PRODUCER's handoff branch: it says nothing about whether this attempt's work survived,
// and naming it would tell the agent another unit's branch is its own published work.
var task = ApplyRetryDisposition(escalatedTask, prior, priorResult, workspaceRef: priorAttemptStaging.Ref);

if (AgentRetryCauses.Classify(priorResult?.Error) == AgentRetryCauses.GatewayFormatFault)
_logger.LogWarning("Supervisor retry of subtask {SubtaskId}: the prior attempt died on a gateway FORMAT fault — retrying FRESH (a conversation replay re-triggers the fault) with extended thinking disabled ({EnvVar}=0)", retry.SubtaskId, AgentRetryCauses.MaxThinkingTokensEnvVar);
Expand Down Expand Up @@ -513,20 +530,44 @@ private async Task<SupervisorExecution> ExecuteRetryAsync(SupervisorDecision dec
/// World-state continuity (the prior branch tip staging) is decided elsewhere and stays UNCHANGED either way —
/// the degrade drops the broken conversation, never the preserved work.
/// </summary>
internal static AgentTask ApplyRetryDisposition(AgentTask task, ResumableSession? prior, SupervisorAgentResult? priorResult, bool workspaceHasPriorWork)
internal static AgentTask ApplyRetryDisposition(AgentTask task, ResumableSession? prior, SupervisorAgentResult? priorResult, string? workspaceRef)
{
if (AgentRetryCauses.Classify(priorResult?.Error) == AgentRetryCauses.GatewayFormatFault)
return AgentRetryCauses.ApplyFormatFaultMitigation(task);

return prior is null ? task : ApplyResumeRecord(task, prior, workspaceHasPriorWork);
return prior is null ? task : ApplyResumeRecord(task, prior, workspaceRef);
}

/// <summary>The pure fold of a resumable prior attempt onto the task: always stamps the session/transcript, and — ONLY when <paramref name="workspaceHasPriorWork"/> is false — appends the honest-redo line so the hint's truth value always matches the actual git state. Internal + static so the honesty branch is unit-pinned directly.</summary>
internal static AgentTask ApplyResumeRecord(AgentTask task, ResumableSession prior, bool workspaceHasPriorWork)
/// <summary>
/// The pure fold of a resumable prior attempt onto the task: always stamps the session + transcript, then says
/// what is true about the WORLD that conversation refers to.
///
/// <para>Two shapes, because the prior attempt ended two different ways. An attempt that FINISHED left its
/// workspace behind, so the only open question is whether it pushed a branch —
/// <paramref name="workspaceRef"/> null means it did not, and the honest-redo line says so. An attempt whose
/// HOST died (<see cref="ResumableSession.CheckpointAt"/>) left nothing but the conversation: its clone is
/// unreachable, so it owes the lost-host block instead, it carries its provenance onto the new run
/// (<c>ResumedFromCheckpointAt</c> for the permanent confinement record, <c>ResumedFromAgentRunId</c> for the
/// column), and its transcript ref is marked a CHECKPOINT so the executor degrades to a cold start rather than
/// failing the attempt when those bytes cannot be read.</para>
///
/// <para>Internal + static so both honesty branches are unit-pinned directly.</para>
/// </summary>
/// <param name="workspaceRef">The branch the RESUMED attempt itself pushed, or null when it pushed none — never the effective clone ref, which for a dependent unit is its producer's handoff branch and says nothing about this attempt's work.</param>
internal static AgentTask ApplyResumeRecord(AgentTask task, ResumableSession prior, string? workspaceRef)
{
var resumed = task with { ResumeFromSessionId = prior.SessionId, RestoredTranscript = prior.InlineTranscript, RestoredTranscriptArtifactId = prior.TranscriptArtifactId };

return workspaceHasPriorWork ? resumed : resumed with { Goal = AgentRetryContinuity.WithHonestNoContinuityHint(resumed.Goal) };
if (prior.CheckpointAt is { } checkpointAt)
return resumed with
{
RestoredTranscriptIsCheckpoint = true,
ResumedFromCheckpointAt = checkpointAt,
ResumedFromAgentRunId = prior.AgentRunId,
Goal = AgentRetryContinuity.WithLostHostHint(resumed.Goal, workspaceRef, treeOwed: task.RepositoryId is not null),
};

return workspaceRef is not null ? resumed : resumed with { Goal = AgentRetryContinuity.WithHonestNoContinuityHint(resumed.Goal) };
}

/// <summary>
Expand Down Expand Up @@ -775,9 +816,12 @@ private async Task<SupervisorExecution> StageAgentsAndParkAsync(IReadOnlyList<(A
var reclaimed = k < orphans.Count;
reclaimedAny |= reclaimed;

// 3c: the ONE staging seam every spawn wave, retry and resolve passes through, so the checkpoint opt-in is
// decided once here rather than at each verb's own task build — see CheckpointsSessionTranscript for who
// is excluded and why.
var agentRunId = reclaimed
? orphans[k]
: await CreateResolvedAgentRunAsync(tasks[k].Task, tasks[k].Spec, context, cancellationToken).ConfigureAwait(false);
: await CreateResolvedAgentRunAsync(tasks[k].Task with { CheckpointSessionTranscript = CheckpointsSessionTranscript(tasks[k].Task, context, tasks.Count) }, tasks[k].Spec, context, cancellationToken).ConfigureAwait(false);

StageAgentWait(context, k, agentRunId);
agentRunIds.Add(agentRunId);
Expand Down
22 changes: 22 additions & 0 deletions backend/src/CodeSpace.Core/Services/Supervisor/SupervisorBounds.cs
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,28 @@ private static bool IsReauthorableRejection(SupervisorPriorDecision decision) =>
return null;
}

/// <summary>
/// Whether the run could still afford to RESPAWN a unit this wave is about to stage — read from the same
/// total-spawn bound <see cref="PostDecision"/> enforces, so the two can never disagree about what the run can
/// still do.
///
/// <para>A later <c>retry</c> costs exactly one spawn, and <see cref="PostDecision"/> refuses it when
/// <c>TotalSpawnedAgents + 1</c> would exceed the cap. By then this wave's own <paramref name="waveSize"/> agents
/// are on the tape, so the room for it has to exist NOW: strictly less than the cap, not at it.</para>
///
/// <para>The spawn cap and nothing else, deliberately. It is the one bound that says "no further agent can ever
/// be created" — whereas a run near its no-progress cap can still retry, because a wave that makes progress
/// resets that cadence. Gating on no-progress would leave the units most likely to be retried un-checkpointed,
/// which is the opposite of the point. The 3c consumer of this predicate pays a whole-file read and an artifact
/// write per minute, so it is worth asking whether anyone can consume the result.</para>
/// </summary>
public static bool CanRespawnAfterWave(SupervisorTurnContext context, int waveSize)
{
ArgumentNullException.ThrowIfNull(context);

return context.TotalSpawnedAgents + waveSize < (context.MaxTotalSpawns ?? SupervisorLane.DefaultMaxTotalSpawns);
}

/// <summary>How many agents the decision would spawn: a spawn fans out its <c>subtaskIds</c>; a retry is exactly one. Best-effort read — a malformed payload reads 0 (it stages nothing, so it can't breach a count bound).</summary>
internal static int SpawnCount(SupervisorDecision decision)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,10 @@ public static class ArtifactRetentionPolicy

/// <summary>
/// A mid-run session-transcript checkpoint. TWO HOURS, not seven days, and the short floor is the whole reason
/// the class exists: a run writes one of these per minute, each supersedes the last, and the run's terminal write
/// clears the column that references the survivor — so on the seven-day floor a single long run would hold every
/// superseded copy of a growing transcript for over a week. Two hours still sits far outside the window in which
/// the class exists: a run writes one of these per minute, each supersedes the last, and a clean landing
/// (completion or a deliberate cancel) clears the column that references the survivor — so on the seven-day floor
/// a single long run would hold every superseded copy of a growing transcript for over a week. An abandon keeps
/// the column on purpose, and that survivor is then Referenced for good; this floor does not collect it. Two hours still sits far outside the window in which
/// the reference lands (the stamp is the next statement after the write) and far outside the window in which a
/// continuation reads it (an abandon follows the host's death within one liveness window), so the floor costs
/// nothing it protects. The quarantine stays the standard 24 h: the second, independent wait is unchanged.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -506,7 +506,8 @@ private static AgentTask ApplyRespawnResumeHint(AgentTask task, JsonElement? pri
/// base ref/pin byte-identical.</para>
///
/// <para>The returned flag is whether the honest-redo line is OWED, and it follows the PRIMARY repo alone —
/// the same <c>workspaceHasPriorWork: effectiveStaging.Ref is not null</c> read the supervisor's retry uses. A
/// the same read the supervisor's retry makes of its prior attempt's own pushed branch
/// (<c>workspaceRef: priorAttemptStaging.Ref</c>), which is keyed on the target repository too. A
/// multi-repo attempt whose primary push FAILED while a sibling's succeeded still repins that sibling, but its
/// primary re-clones the default branch, so the resumed conversation must still be told its changes are not
/// there: an OR across repos would suppress the line precisely where the agent's own repo lost its work. Nothing
Expand Down
7 changes: 4 additions & 3 deletions backend/src/CodeSpace.Messages/Agents/AgentTask.cs
Original file line number Diff line number Diff line change
Expand Up @@ -109,9 +109,10 @@ public sealed record AgentTask
///
/// <para>An opt-in rather than a default, because a checkpoint nobody will consume is pure waste: it costs a
/// whole-file read and an artifact write per minute, per running agent, per worker. Only a producer whose failed
/// attempt can actually be RETRIED sets it — today that is <c>agent.run</c> for a node whose own retry policy
/// allows more than one attempt. The benchmark lanes (one attempt per cell by protocol), review children and
/// supervisor units leave it false.</para>
/// attempt can actually be RETRIED sets it — today <c>agent.run</c> for a node whose own retry policy allows more
/// than one attempt, and a supervisor unit that carries a subtask id while its run can still afford to respawn
/// it. The benchmark lanes (one attempt per cell by protocol), review children and supervisor resolver units
/// leave it false.</para>
///
/// <para><c>[JsonIgnore(WhenWritingDefault)]</c> so an envelope that did not opt in adds nothing to task_json.</para>
/// </summary>
Expand Down
Loading
Loading