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
20 changes: 10 additions & 10 deletions backend/src/CodeSpace.Core/Services/Agents/AgentRetryContinuity.cs
Original file line number Diff line number Diff line change
Expand Up @@ -11,25 +11,25 @@ namespace CodeSpace.Core.Services.Agents;
public static class AgentRetryContinuity
{
/// <summary>The honest-redo line: fires ONLY when a resumed conversation exists but the workspace was NOT pinned to a prior pushed branch — never on a genuine cold-start retry (no prior attempt at all), which stays byte-identical.</summary>
public const string HonestNoContinuityHint = "Note: your prior attempt's conversation is restored, but its git changes were NOT preserved in this workspace (no pushed branch was found to continue from) — you must redo any relevant file changes from scratch.";
public const string HonestNoContinuityHint = "Note: your prior attempt's conversation is restored, but its git changes were NOT preserved in this workspace (your prior attempt pushed no branch of its own) — you must redo any relevant file changes from scratch.";

/// <summary>Append <see cref="HonestNoContinuityHint"/> to a resumed task's goal. One composition, so the two lanes cannot drift on the separator either.</summary>
public static string WithHonestNoContinuityHint(string goal) => $"{goal}\n\n{HonestNoContinuityHint}";

/// <summary>
/// 3c: what a CROSS-HOST continuation is told, and the reason it needs its own sentence. The two lanes above
/// retry an attempt that FINISHED on a live host, so their only open question is whether a branch was pushed.
/// This lane continues an attempt whose machine was lost mid-run: the conversation comes from a checkpoint taken
/// some time before the loss, and the working tree is simply gone. Both halves have to be said, because the
/// restored transcript will describe edits — possibly edits made after the checkpoint — that the new workspace
/// does not contain, and an agent that is not told will read its own transcript as evidence about files it cannot
/// see.
/// This lane continues an attempt whose machine or process was lost mid-run: the conversation comes from a
/// checkpoint taken some time before the loss, and the working tree is simply gone. Both halves have to be said,
/// because the restored transcript will describe edits — possibly edits made after the checkpoint — that the new
/// workspace does not contain, and an agent that is not told will read its own transcript as evidence about files
/// it cannot see.
/// </summary>
public const string LostHostPreamble = "Note: the machine running your previous attempt was lost mid-run. Your conversation is restored from a checkpoint taken before that, so it may describe work you did after the checkpoint, and it may be missing your last few turns.";
public const string LostHostPreamble = "Note: the machine or the process running your previous attempt was lost mid-run. Your conversation is restored from a checkpoint taken before that, so it may describe work you did after the checkpoint, and it may be missing your last few turns.";

/// <summary>Said when the lost attempt HAD published a branch: the workspace is checked out at it, so the published work is present and only the unpublished remainder is gone. Takes the branch name so the agent can verify rather than take the claim on trust.</summary>
public static string LostHostPublishedBranchHint(string branch) =>
$"Your previous attempt published branch `{branch}`, and this workspace is checked out AT that branch — that work is here. Anything you had NOT published to it died with the machine, so check the files before continuing and redo whatever is missing.";
$"Your previous attempt published branch `{branch}`, and this workspace is checked out AT that branch — that work is here. Anything you had NOT published to it was lost with that attempt, so check the files before continuing and redo whatever is missing.";

/// <summary>
/// Append the cross-host continuation's honesty block to a resumed task's goal: the preamble always, then what
Expand All @@ -49,8 +49,8 @@ public static string WithLostHostHint(string goal, string? publishedBranch, bool
/// <summary>
/// 3c: said when the lost host's checkpoint could not be READ — reaped, or its storage unreachable. The attempt
/// still runs (failing it would spend the retry this whole path exists to improve), but it runs COLD, and an
/// agent that was going to be handed a conversation must be told it is not getting one. Appended to the
/// lost-host block rather than replacing it: the machine really was lost, which is still the reason.
/// agent that was going to be handed a conversation must be told it is not getting one. Appended to the lost-host
/// block rather than replacing it: the machine or the process really was lost, which is still the reason.
/// </summary>
public const string LostHostCheckpointUnreadableHint = "Your previous conversation could not be recovered either — the checkpoint it was stored in is no longer readable — so you are starting this task from the beginning.";

Expand Down
7 changes: 4 additions & 3 deletions backend/src/CodeSpace.Core/Services/Agents/AgentRunService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -777,7 +777,8 @@ public async Task<bool> CancelRunningAsync(Guid runId, string reason, AgentRunAb
// 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.
// the cancel as a host loss. Only an abandon-class ending keeps these columns: the reconciler's abandon, or
// its spool recovery.
var cancelled = await _db.AgentRun
.Where(r => r.Id == runId && r.Status == AgentRunStatus.Running && r.FenceEpoch == snapshot.FenceEpoch)
.ExecuteUpdateAsync(s => s
Expand Down Expand Up @@ -976,8 +977,8 @@ await _db.AgentRun.AsNoTracking().SingleOrDefaultAsync(r => r.Id == runId, cance
// 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.
// abandon-class ending: every clean landing — completion or cancel — releases these two columns in its own
// terminal write; only the reconciler's abandon, or its spool recovery, keeps them.
if (TryResumableFromCheckpoint(candidate.Id, candidate.Status, candidate.SessionId, candidate.SessionTranscriptCheckpointArtifactId, candidate.SessionTranscriptCheckpointAt) is { } continued) return continued;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -267,7 +267,11 @@ private Lease Install(Lease lease, TimeSpan ttl)

if (superseded is not null) CloseQuietly(superseded.Listener);

_ = Task.Run(() => AcceptAsync(lease), CancellationToken.None);
// CALLED, not handed to Task.Run: an async method runs synchronously up to its first await, so the loop's first
// wait is registered on the listener before this returns. The managed HttpListener fails only the waits it
// already holds when it closes; a close that lands while a pool thread is still registering the first one is
// never delivered, and that loop then waits for the life of the worker on a listener that no longer exists.
_ = AcceptAsync(lease);

return lease;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -642,7 +642,7 @@ private async Task<SupervisorPriorDecision> FoldAcceptanceGradeAsync(SupervisorP
// without a pulse the reconciler reads a genuinely-alive resolve grade as abandoned and re-dispatches
// the run mid-grade. Starts only when a real grade fires (every early return above skips it).
using var heartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
var heartbeat = RunGradingHeartbeatLoopAsync(supervisorRunId, nodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, heartbeatCts.Token, "Supervisor resolve acceptance grading is still in progress.");
var heartbeat = RunGradingHeartbeatLoopAsync(supervisorRunId, nodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, heartbeatCts.Token, TimeProvider.System, "Supervisor resolve acceptance grading is still in progress.");

BenchmarkGrade grade;
try
Expand Down Expand Up @@ -762,7 +762,7 @@ private async Task<SupervisorPriorDecision> FoldUnitAcceptanceGradeAsync(Supervi
// reconciler reads this genuinely-alive fold as abandoned and re-dispatches the run mid-grade (the S3
// adversarial scan's M4: the baseline roughly doubled the silent window).
using var heartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
var heartbeat = RunGradingHeartbeatLoopAsync(supervisorRunId, nodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, heartbeatCts.Token, "Supervisor per-unit acceptance grading is still in progress.");
var heartbeat = RunGradingHeartbeatLoopAsync(supervisorRunId, nodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, heartbeatCts.Token, TimeProvider.System, "Supervisor per-unit acceptance grading is still in progress.");

try
{
Expand Down Expand Up @@ -1552,7 +1552,7 @@ private static IReadOnlyList<SupervisorAgentResult> BranchlessUnits(SupervisorTu
private async Task<BenchmarkGrade> GradeStopTargetsWithHeartbeatAsync(Guid supervisorRunId, string nodeId, Guid teamId, IReadOnlyList<(Guid RepositoryId, string Alias, string Branch)> targets, IReadOnlyList<(string Label, SupervisorAcceptanceSpec? Spec)> gates, IReadOnlyDictionary<Guid, string> oracleBaseShas, IReadOnlyList<string> oracleFloorPrograms, CancellationToken cancellationToken)
{
using var heartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
var heartbeat = RunGradingHeartbeatLoopAsync(supervisorRunId, nodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, heartbeatCts.Token);
var heartbeat = RunGradingHeartbeatLoopAsync(supervisorRunId, nodeId, SupervisorLane.AcceptanceGradeHeartbeatInterval, heartbeatCts.Token, TimeProvider.System);

try
{
Expand All @@ -1567,14 +1567,14 @@ private async Task<BenchmarkGrade> GradeStopTargetsWithHeartbeatAsync(Guid super
}
}

/// <summary>The heartbeat loop itself: sleeps, logs, repeats — until <paramref name="cancellationToken"/> fires (grading finished). A cancellation mid-sleep is the expected exit, never propagated as a fault. Internal + interval-parameterized so a unit test can pin the cancellation contract with a millisecond-scale interval instead of waiting out the real 90s production value.</summary>
internal async Task RunGradingHeartbeatLoopAsync(Guid supervisorRunId, string nodeId, TimeSpan interval, CancellationToken cancellationToken, string message = "Supervisor stop acceptance grading is still in progress.")
/// <summary>The heartbeat loop itself: sleeps, logs, repeats — until <paramref name="cancellationToken"/> fires (grading finished). A cancellation mid-sleep is the expected exit, never propagated as a fault. Internal + clock-parameterized so a unit test drives the sleep on a fake <paramref name="timeProvider"/> instead of racing the wall clock; REQUIRED rather than defaulting to the system clock for the reason <see cref="HeartbeatLoop.RunAsync"/> gives — a default is how a call site keeps the wall clock without saying so.</summary>
internal async Task RunGradingHeartbeatLoopAsync(Guid supervisorRunId, string nodeId, TimeSpan interval, CancellationToken cancellationToken, TimeProvider timeProvider, string message = "Supervisor stop acceptance grading is still in progress.")
{
try
{
while (true)
{
await Task.Delay(interval, cancellationToken).ConfigureAwait(false);
await Task.Delay(interval, timeProvider, cancellationToken).ConfigureAwait(false);

await _recordLogger.LogAsync(supervisorRunId, nodeId, Workflows.Lifecycle.LogLevel.Info,
message, cancellationToken).ConfigureAwait(false);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,9 +42,10 @@ public static class ArtifactRetentionPolicy
/// 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 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
/// a single long run would hold every superseded copy of a growing transcript for over a week. An abandon-class
/// ending (the reconciler's abandon, or its spool recovery) 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.
/// </summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -808,7 +808,7 @@ public async Task A_retried_agent_node_restores_the_lost_hosts_checkpoint()
resumed.ResumedFromCheckpointAt.ShouldNotBeNull("the launch stamps this onto the run's permanent confinement record");
resumed.ResumedFromAgentRunId.ShouldBe(lostAgent, "which attempt took over from which is a column, not prose");
retry.ResumedFromAgentRunId.ShouldBe(lostAgent, "and the task's provenance is promoted onto the row, like AgentDefinitionId");
resumed.Goal.ShouldContain("machine running your previous attempt was lost", Case.Sensitive,
resumed.Goal.ShouldContain("the machine or the process running your previous attempt was lost", Case.Sensitive,
"a restored conversation describes a working tree this sandbox does not have, and the agent must be told rather than left to infer it");
}
finally
Expand Down Expand Up @@ -850,7 +850,7 @@ public async Task A_host_loss_with_no_checkpoint_is_retried_cold_and_claims_noth
resumed.ResumedFromCheckpointAt.ShouldBeNull();
resumed.ResumedFromAgentRunId.ShouldBeNull();
retry.ResumedFromAgentRunId.ShouldBeNull();
resumed.Goal.ShouldNotContain("machine running your previous attempt was lost", Case.Sensitive, "nothing may assert a restored conversation this attempt does not have");
resumed.Goal.ShouldNotContain("the machine or the process running your previous attempt was lost", Case.Sensitive, "nothing may assert a restored conversation this attempt does not have");
}
finally
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -520,11 +520,17 @@ public async Task A_lease_whose_listener_died_stops_being_claimed()

broker.HasLease(runId).ShouldBeTrue("precondition: the lease is live and its listener is accepting");

while (logger.Warned.Wait(0)) { } // whatever the open itself warned about is not the signal awaited below

// The platform failed the accept — a listener closed under us, an error nobody enumerated. The lease is still
// in the table and still inside its window, so nothing about TIME will correct it.
broker.BreakListenerForTest(runId);

await WaitUntilAsync(() => !broker.HasLease(runId), TimeSpan.FromSeconds(10),
// Woken by the drop's LAST effect rather than by polling HasLease: the warning is written only after the lease
// has left the table, so once it lands both assertions below read a finished drop instead of racing one.
(await logger.Warned.WaitAsync(TimeSpan.FromSeconds(10))).ShouldBeTrue("the accept loop never reported its dead listener within 10s — the close never reached a loop still registering its first wait, or the failure was swallowed without a word");

broker.HasLease(runId).ShouldBeFalse(
"a lease whose accept loop has stopped went on reporting itself live. That is the worst answer this class can give: the child's connections sit unaccepted in a backlog instead of being refused, and a re-attach reading HasLease true concludes the run still has model access — so it lands no verdict and leaves the run Running, with no model and no explanation, for as long as the worker lives");

logger.Warnings.ShouldContain(line => line.Contains(runId.ToString(), StringComparison.Ordinal),
Expand Down Expand Up @@ -575,7 +581,7 @@ public void A_rebind_is_only_built_for_a_handle_whose_agent_this_host_can_reach(
}

/// <summary>Stands in for this worker's own host identity inside <c>[InlineData]</c>, which cannot carry a runtime value.</summary>
private const string ThisHost = "this-host";
private const string ThisHost = "\0this-host";

/// <summary>The re-bind a later worker would build from what a run's durable handle carries — the point being that every value comes from <paramref name="brokered"/>, because a re-bind restores an address and never mints one.</summary>
private static ModelCredentialRebindRequest RebindOf(BrokeredModelCredential brokered, Guid runId, long epoch) => new()
Expand Down Expand Up @@ -654,12 +660,18 @@ private sealed class CapturingLogger : Microsoft.Extensions.Logging.ILogger<Loop
{
public List<string> Warnings { get; } = [];

/// <summary>Released AFTER each warning is recorded, so a waiter that acquires it reads a <see cref="Warnings"/> that already holds that line.</summary>
public SemaphoreSlim Warned { get; } = new(0);

public IDisposable BeginScope<TState>(TState state) where TState : notnull => NullScope.Instance;
public bool IsEnabled(Microsoft.Extensions.Logging.LogLevel logLevel) => true;

public void Log<TState>(Microsoft.Extensions.Logging.LogLevel logLevel, Microsoft.Extensions.Logging.EventId eventId, TState state, Exception? exception, Func<TState, Exception?, string> formatter)
{
if (logLevel >= Microsoft.Extensions.Logging.LogLevel.Warning) Warnings.Add(formatter(state, exception));
if (logLevel < Microsoft.Extensions.Logging.LogLevel.Warning) return;

Warnings.Add(formatter(state, exception));
Warned.Release();
}

private sealed class NullScope : IDisposable { public static readonly NullScope Instance = new(); public void Dispose() { } }
Expand Down
Loading
Loading