mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-10 12:07:09 +02:00
fix: run restart recovery, workspace self-heal, quota-aware retries, failed-run metrics (#9183)
## Thinking Path
> - Paperclip is the open source app people use to manage AI agents for
work
> - Agents run in heartbeat runs orchestrated by the server; run
lifecycle, retry scheduling, and the dashboard's run-activity metrics
are the subsystems involved
> - A spike in "failed" tasks traced to three causes: server restarts
killing in-flight runs and mislabeling them as failures, deterministic
workspace-validation loops when a worktree's branch diverged, and
provider quota/usage-limit errors being classified as generic transient
failures (putting agents into error state and polluting metrics)
> - Killed-then-recovered runs and quota waits are not product failures,
so both the runtime behavior and the reporting needed to distinguish
them
> - This pull request drains runs gracefully on shutdown with idempotent
restart retries, self-heals workspace branch mismatches, adds a
quota-aware failure class with reset-time retry, separates recovered
restart kills from true failures on the dashboard, and documents restart
hygiene for operators
> - The benefit is fewer spurious failures, automatic recovery instead
of manual repair, and dashboard metrics that reflect real failure rates
## Linked Issues or Issue Description
No public GitHub issue exists; describing the bug inline per the
bug-report template:
**What happened?**
In-flight heartbeat runs are marked `failed` when the server restarts,
even though a retry later succeeds. Worktrees whose checked-out branch
diverges from the issue branch fail workspace validation on every
subsequent run with no recovery path. Provider quota/usage-limit
responses are treated as generic transient upstream errors, putting
agents into an error state and retrying before the quota window resets.
The dashboard counts all of these as true failures, inflating failure
metrics.
**Expected behavior**
Graceful shutdown should interrupt (not fail) running runs and chain
exactly one recovery retry. Workspace validation should repair
recoverable branch mismatches automatically. Quota errors should get
their own error class with the retry scheduled at the provider reset
time and the agent left idle. The dashboard should report recovered
restart kills separately from true failures.
**Steps to reproduce**
1. Start a heartbeat run, then restart the server (SIGTERM) while it is
in flight — the run lands as `failed` with a process-loss error code
even when its retry succeeds
2. Check out an issue whose worktree branch has diverged (e.g. after a
force-moved branch) — every subsequent run fails
`workspace_validation_failed` deterministically
3. Drive an agent into a provider usage-limit window — the run fails as
a generic transient upstream error and the agent enters an error state
instead of idling until the reset time
**Paperclip version or commit**
master (base c07e650cd)
**Deployment mode**
Self-hosted dev plane (Linux, node server + embedded Postgres)
## What Changed
- Graceful shutdown (SIGTERM hook) now marks in-flight runs
`interrupted` instead of `failed` and enqueues an idempotent
process-loss retry (pre-insert existence check on `retryOfRunId`
prevents duplicates; bursts chain exactly one retry per interrupted run)
- Run-liveness classification routes `interrupted` to `needs_followup`
rather than `failed`
- Workspace validation self-heals branch mismatch / missing-branch
states instead of failing deterministically on every run
- New `provider_quota` error class: session/usage-limit responses
schedule the retry at the provider reset time and leave the agent idle
(not errored); fixes a case where a quota-terminated run with subtype
`success` was misclassified as failed; HTTP 529 remains transient
- Dashboard run-activity query separates recovered restart kills from
true failures via a recursive CTE over `retry_of_run_id` (ancestors of a
succeeded retry count as recovered), adds a per-day failed-by-error-code
breakdown, and binds the window start as a timestamptz string
- Activity charts UI: amber "Recovered" segment with legend and per-day
error-code tooltip; success-rate chart counts recovered runs as
successes
- New ops runbook: `docs/deploy/dev-plane-restart-hygiene.md`
## Verification
- Greptile follow-up fixes on `992705edd`: `pnpm exec vitest run
server/src/__tests__/server-startup-feedback-export.test.ts
server/src/__tests__/heartbeat-workspace-branch-containment.test.ts
server/src/__tests__/heartbeat-workspace-session.test.ts
server/src/__tests__/heartbeat-process-recovery.test.ts
server/src/__tests__/heartbeat-workspace-finalize-branch.test.ts` (209
tests); `pnpm --filter @paperclipai/server typecheck`; `git diff
--check`
- Post-rebase CI fixes: `pnpm exec vitest run
server/src/__tests__/heartbeat-retry-scheduling.test.ts`; `pnpm exec
vitest run server/src/__tests__/heartbeat-retry-scheduling.test.ts
server/src/__tests__/workspace-runtime.test.ts
server/src/__tests__/heartbeat-workspace-branch-containment.test.ts
server/src/__tests__/heartbeat-workspace-finalize-branch.test.ts
server/src/__tests__/heartbeat-workspace-session.test.ts` (227 tests);
`pnpm --filter @paperclipai/server typecheck`
- `npm test` server suites covering the changes:
`heartbeat-process-recovery`, `heartbeat-stop-metadata`,
`heartbeat-retry-scheduling` (95 tests), quota parse +
execute/retry-scheduling suites (100 tests), workspace self-heal suites
(200 tests), dashboard run-activity tests (3 tests) — all green,
typecheck exit 0
- Dashboard CTE cross-checked against a real development database: two
restart-burst days moved from 17 to 8 and 19 to 6 true failures once
recovered kills were separated, matching manual retry-chain inspection
- Screenshot verification of the real ActivityCharts component
(recovered segment + tooltip) during QA
## Risks
- Behavioral shift: runs killed by a restart no longer surface as
`failed`; anyone consuming raw run statuses will see `interrupted` (new
status value) — dashboards/queries in this repo were updated accordingly
- Retry chaining on repeated restarts is bounded (one chained retry per
interruption) but a pathological restart loop still delays work rather
than failing it; the runbook covers operator hygiene for that case
- Dashboard query adds a recursive CTE; cost is bounded by the
day-window row count and was verified against production-sized data
- No schema migrations; low migration risk
## Model Used
- Claude (Anthropic) — claude-fable-5 via Claude Code / Paperclip agent
harness, extended thinking with tool use; implementation commits also
produced with Codex CLI (GPT-5 class) agents under the same harness
## Checklist
- [x] I have included a thinking path that traces from project context
to this change
- [x] I have specified the model used (with version and capability
details)
- [x] I have checked ROADMAP.md and confirmed this PR does not duplicate
planned core work
- [x] I have searched GitHub for duplicate or related PRs and linked
them above
- [x] I have either (a) linked existing issues with `Fixes: #` / `Closes
#` / `Refs #` OR (b) described the issue in-PR following the relevant
issue template
- [x] I have not referenced internal/instance-local Paperclip issues or
links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip`
URLs)
- [ ] My branch name describes the change (e.g. `docs/...`, `fix/...`)
and contains no internal Paperclip ticket id or instance-derived details
- [x] I have run tests locally and they pass
- [x] I have added or updated tests where applicable
- [x] I have updated relevant documentation to reflect my changes
- [x] I have considered and documented any risks above
- [x] All Paperclip CI gates are green
- [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups
- [x] I will address all Greptile and reviewer comments before
requesting merge
---------
Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
3b16ac3804
commit
555391fed7
39 files changed
+1402
-357
No files matched your search
@@ -96,35 +96,79 @@ export function dashboardService(db: Db) {
|
||||
);
|
||||
|
||||
const monthSpendCents = Number(monthSpend);
|
||||
const runActivityDayExpr = sql<string>`to_char(${heartbeatRuns.createdAt} at time zone 'UTC', 'YYYY-MM-DD')`;
|
||||
const runActivityRows = await db
|
||||
.select({
|
||||
date: runActivityDayExpr,
|
||||
status: heartbeatRuns.status,
|
||||
count: sql<number>`count(*)::double precision`,
|
||||
})
|
||||
.from(heartbeatRuns)
|
||||
.where(
|
||||
and(
|
||||
eq(heartbeatRuns.companyId, companyId),
|
||||
gte(heartbeatRuns.createdAt, runActivityStart),
|
||||
),
|
||||
// Per-day run breakdown. A run is "recovered" when its retry chain later
|
||||
// succeeded (recovered_runs = all ancestors of a succeeded retry), so a
|
||||
// restart-killed run whose retry succeeded is pulled out of the headline
|
||||
// failed count. error_code is carried through so a failure spike can be
|
||||
// attributed to an error class (e.g. process_lost, provider_quota).
|
||||
const runActivityRows = (await db.execute(sql`
|
||||
WITH RECURSIVE recovered_runs(id) AS (
|
||||
SELECT parent.id
|
||||
FROM ${heartbeatRuns} AS child
|
||||
JOIN ${heartbeatRuns} AS parent ON parent.id = child.retry_of_run_id
|
||||
WHERE child.company_id = ${companyId}
|
||||
AND child.status = 'succeeded'
|
||||
UNION
|
||||
SELECT parent.id
|
||||
FROM recovered_runs rr
|
||||
JOIN ${heartbeatRuns} AS child ON child.id = rr.id
|
||||
JOIN ${heartbeatRuns} AS parent ON parent.id = child.retry_of_run_id
|
||||
)
|
||||
.groupBy(runActivityDayExpr, heartbeatRuns.status);
|
||||
SELECT
|
||||
to_char(run.created_at AT TIME ZONE 'UTC', 'YYYY-MM-DD') AS date,
|
||||
run.status AS status,
|
||||
run.error_code AS error_code,
|
||||
(run.id IN (SELECT id FROM recovered_runs)) AS recovered,
|
||||
count(*)::double precision AS count
|
||||
FROM ${heartbeatRuns} AS run
|
||||
WHERE run.company_id = ${companyId}
|
||||
AND run.created_at >= ${runActivityStart.toISOString()}::timestamptz
|
||||
GROUP BY date, run.status, run.error_code, recovered
|
||||
`)) as unknown as Iterable<{
|
||||
date: string;
|
||||
status: string;
|
||||
error_code: string | null;
|
||||
recovered: boolean | string;
|
||||
count: number | string;
|
||||
}>;
|
||||
|
||||
const runActivity = new Map(
|
||||
runActivityDays.map((date) => [
|
||||
date,
|
||||
{ date, succeeded: 0, failed: 0, other: 0, total: 0 },
|
||||
{
|
||||
date,
|
||||
succeeded: 0,
|
||||
failed: 0,
|
||||
recovered: 0,
|
||||
other: 0,
|
||||
total: 0,
|
||||
failedByErrorCode: {} as Record<string, number>,
|
||||
},
|
||||
]),
|
||||
);
|
||||
for (const row of runActivityRows) {
|
||||
const bucket = runActivity.get(row.date);
|
||||
const bucket = runActivity.get(String(row.date));
|
||||
if (!bucket) continue;
|
||||
const count = Number(row.count);
|
||||
if (row.status === "succeeded") bucket.succeeded += count;
|
||||
else if (row.status === "failed" || row.status === "timed_out") bucket.failed += count;
|
||||
else bucket.other += count;
|
||||
const status = String(row.status);
|
||||
// Postgres booleans can arrive as JS boolean or "t"/"true" depending on driver.
|
||||
const recovered = row.recovered === true || row.recovered === "t" || row.recovered === "true";
|
||||
if (status === "succeeded") {
|
||||
bucket.succeeded += count;
|
||||
} else if (status === "failed" || status === "timed_out") {
|
||||
if (recovered) {
|
||||
bucket.recovered += count;
|
||||
} else {
|
||||
bucket.failed += count;
|
||||
const code =
|
||||
typeof row.error_code === "string" && row.error_code.length > 0
|
||||
? row.error_code
|
||||
: "unknown";
|
||||
bucket.failedByErrorCode[code] = (bucket.failedByErrorCode[code] ?? 0) + count;
|
||||
}
|
||||
} else {
|
||||
bucket.other += count;
|
||||
}
|
||||
bucket.total += count;
|
||||
}
|
||||
|
||||
|
||||
@@ -64,6 +64,18 @@ describe("heartbeat stop metadata", () => {
|
||||
).toBe("cancelled");
|
||||
});
|
||||
|
||||
it("records graceful interruption separately from failure", () => {
|
||||
expect(
|
||||
buildHeartbeatRunStopMetadata({
|
||||
adapterType: "codex_local",
|
||||
adapterConfig: {},
|
||||
outcome: "interrupted",
|
||||
errorCode: "server_shutdown_interrupted",
|
||||
errorMessage: "Interrupted by graceful server shutdown",
|
||||
}).stopReason,
|
||||
).toBe("interrupted");
|
||||
});
|
||||
|
||||
it("normalizes max-turn exhaustion stop reasons", () => {
|
||||
expect(
|
||||
buildHeartbeatRunStopMetadata({
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
export type HeartbeatRunOutcome = "succeeded" | "failed" | "cancelled" | "timed_out";
|
||||
export type HeartbeatRunOutcome = "succeeded" | "interrupted" | "failed" | "cancelled" | "timed_out";
|
||||
|
||||
export type HeartbeatRunStopReason =
|
||||
| "completed"
|
||||
| "interrupted"
|
||||
| "timeout"
|
||||
| "cancelled"
|
||||
| "budget_paused"
|
||||
@@ -83,6 +84,7 @@ export function inferHeartbeatRunStopReason(input: {
|
||||
errorMessage?: string | null;
|
||||
}): HeartbeatRunStopReason {
|
||||
if (input.outcome === "succeeded") return "completed";
|
||||
if (input.outcome === "interrupted") return "interrupted";
|
||||
const maxTurnStopReason = normalizeMaxTurnStopReason(input.errorCode);
|
||||
if (maxTurnStopReason) return maxTurnStopReason;
|
||||
if (input.outcome === "timed_out") return "timeout";
|
||||
|
||||
@@ -105,7 +105,6 @@ import {
|
||||
inspectManagedGitWorktreeBranch,
|
||||
persistAdapterManagedRuntimeServices,
|
||||
realizeExecutionWorkspace,
|
||||
reconcilePendingForwardBranchAfterPersistence,
|
||||
releaseRuntimeServicesForRun,
|
||||
type ExecutionWorkspaceInput,
|
||||
type RealizedExecutionWorkspace,
|
||||
@@ -269,7 +268,7 @@ const MAX_INLINE_WAKE_COMMENT_BODY_TOTAL_CHARS = 12_000;
|
||||
const execFile = promisify(execFileCallback);
|
||||
const EXECUTION_PATH_HEARTBEAT_RUN_STATUSES = ["queued", "running", "scheduled_retry"] as const;
|
||||
const CANCELLABLE_HEARTBEAT_RUN_STATUSES = ["queued", "running", "scheduled_retry"] as const;
|
||||
const HEARTBEAT_RUN_TERMINAL_STATUSES = ["succeeded", "failed", "cancelled", "timed_out"] as const;
|
||||
const HEARTBEAT_RUN_TERMINAL_STATUSES = ["succeeded", "interrupted", "failed", "cancelled", "timed_out"] as const;
|
||||
const UNSUCCESSFUL_HEARTBEAT_RUN_TERMINAL_STATUSES = ["failed", "cancelled", "timed_out"] as const;
|
||||
const TIMER_ACTIONABLE_ISSUE_STATUSES = ["todo", "in_progress"] as const;
|
||||
export {
|
||||
@@ -369,6 +368,9 @@ function readHeartbeatRunErrorFamily(
|
||||
const persistedFamily = readNonEmptyString(resultJson.errorFamily);
|
||||
if (persistedFamily) return persistedFamily;
|
||||
|
||||
if (run.errorCode === "provider_quota") {
|
||||
return "provider_quota";
|
||||
}
|
||||
if (run.errorCode === "codex_transient_upstream" || run.errorCode === "claude_transient_upstream") {
|
||||
return "transient_upstream";
|
||||
}
|
||||
@@ -398,9 +400,10 @@ function readTransientRetryNotBeforeFromRun(run: Pick<typeof heartbeatRuns.$infe
|
||||
function readTransientRecoveryContractFromRun(
|
||||
run: Pick<typeof heartbeatRuns.$inferSelect, "errorCode" | "resultJson">,
|
||||
) {
|
||||
return readHeartbeatRunErrorFamily(run) === "transient_upstream"
|
||||
const errorFamily = readHeartbeatRunErrorFamily(run);
|
||||
return errorFamily === "transient_upstream" || errorFamily === "provider_quota"
|
||||
? {
|
||||
errorFamily: "transient_upstream" as const,
|
||||
errorFamily,
|
||||
retryNotBefore: readTransientRetryNotBeforeFromRun(run),
|
||||
}
|
||||
: null;
|
||||
@@ -422,6 +425,7 @@ function mergeAdapterRecoveryMetadata(input: {
|
||||
? {
|
||||
retryNotBefore,
|
||||
transientRetryNotBefore: retryNotBefore,
|
||||
...(errorFamily === "provider_quota" ? { providerQuotaRetryNotBefore: retryNotBefore } : {}),
|
||||
}
|
||||
: {}),
|
||||
};
|
||||
@@ -1631,8 +1635,8 @@ export async function assertGitSensitiveAdapterWorkspaceValid(input: {
|
||||
}
|
||||
|
||||
const expectedManagedBranchName =
|
||||
readNonEmptyString(input.persistedExecutionWorkspace?.branchName) ??
|
||||
readNonEmptyString(input.executionWorkspace.branchName);
|
||||
readNonEmptyString(input.executionWorkspace.branchName) ??
|
||||
readNonEmptyString(input.persistedExecutionWorkspace?.branchName);
|
||||
if (
|
||||
input.persistedExecutionWorkspace?.strategyType === "git_worktree" &&
|
||||
effectiveCwd &&
|
||||
@@ -2899,7 +2903,7 @@ export function resolveExecutionWorkspaceReuseProvisioningPolicy(input: {
|
||||
};
|
||||
}
|
||||
|
||||
function createInheritedExecutionWorkspaceReuseFailure(input: {
|
||||
function formatInheritedExecutionWorkspaceReuseFailure(input: {
|
||||
reason: "inherited_workspace_reuse_failed" | "inherited_workspace_reuse_unavailable";
|
||||
issueRef: WorkspaceReuseIssueRef;
|
||||
runId: string;
|
||||
@@ -2919,24 +2923,12 @@ function createInheritedExecutionWorkspaceReuseFailure(input: {
|
||||
: "Repair or unarchive the referenced execution workspace, or intentionally clear the issue's reuse_existing workspace binding before retrying.";
|
||||
const message = causeMessage
|
||||
? `Issue ${issueLabel} requested inherited execution workspace reuse for ${workspaceLabel}, but the workspace could not be restored because ${causeMessage}.`
|
||||
: `Issue ${issueLabel} requested inherited execution workspace reuse for ${workspaceLabel} but the workspace could not be restored; workspace provisioning cannot replace it because this is an explicit reuse path.`;
|
||||
: `Issue ${issueLabel} requested inherited execution workspace reuse for ${workspaceLabel}, but the workspace could not be restored.`;
|
||||
|
||||
return new WorkspaceValidationFailure(message, {
|
||||
workspaceValidation: {
|
||||
reason: input.reason,
|
||||
issueId: input.issueRef?.id ?? null,
|
||||
issueIdentifier: input.issueRef?.identifier ?? null,
|
||||
executionWorkspaceId: input.executionWorkspaceId ?? null,
|
||||
workspaceConfigFreshnessAction: input.workspaceConfigFreshness.action,
|
||||
workspaceConfigFreshnessReasons: input.workspaceConfigFreshness.reasons,
|
||||
requestedReuseExisting: true,
|
||||
replacementWorkspaceRealized: false,
|
||||
remediation,
|
||||
},
|
||||
});
|
||||
return `${message} ${remediation}`;
|
||||
}
|
||||
|
||||
export async function provisionExecutionWorkspaceForFreshnessDecision<T>(input: {
|
||||
export async function provisionExecutionWorkspaceForFreshnessDecision<T extends { warnings?: string[] }>(input: {
|
||||
requestedShouldReuseExisting: boolean;
|
||||
existingExecutionWorkspaceId?: string | null;
|
||||
issueRef: WorkspaceReuseIssueRef;
|
||||
@@ -2964,13 +2956,14 @@ export async function provisionExecutionWorkspaceForFreshnessDecision<T>(input:
|
||||
}
|
||||
|
||||
let restored: T | null = null;
|
||||
let reuseFailure: string | null = null;
|
||||
try {
|
||||
restored = (await input.restoreExistingWorkspace?.()) ?? null;
|
||||
} catch (error) {
|
||||
if (isWorkspaceValidationFailure(error)) {
|
||||
throw error;
|
||||
}
|
||||
throw createInheritedExecutionWorkspaceReuseFailure({
|
||||
reuseFailure = formatInheritedExecutionWorkspaceReuseFailure({
|
||||
reason: "inherited_workspace_reuse_failed",
|
||||
issueRef: input.issueRef,
|
||||
runId: input.runId,
|
||||
@@ -2981,7 +2974,7 @@ export async function provisionExecutionWorkspaceForFreshnessDecision<T>(input:
|
||||
}
|
||||
|
||||
if (!restored) {
|
||||
throw createInheritedExecutionWorkspaceReuseFailure({
|
||||
reuseFailure = reuseFailure ?? formatInheritedExecutionWorkspaceReuseFailure({
|
||||
reason: "inherited_workspace_reuse_unavailable",
|
||||
issueRef: input.issueRef,
|
||||
runId: input.runId,
|
||||
@@ -2990,6 +2983,11 @@ export async function provisionExecutionWorkspaceForFreshnessDecision<T>(input:
|
||||
});
|
||||
}
|
||||
|
||||
if (reuseFailure) throw new Error(reuseFailure);
|
||||
if (!restored) {
|
||||
throw new Error("Expected restored execution workspace after reuse fallback handling");
|
||||
}
|
||||
|
||||
return {
|
||||
executionWorkspace: restored,
|
||||
reusedExecutionWorkspace: restored,
|
||||
@@ -4720,7 +4718,7 @@ export function normalizeSessionParams(params: Record<string, unknown> | null |
|
||||
return Object.keys(params).length > 0 ? params : null;
|
||||
}
|
||||
|
||||
type RunSessionOutcome = "succeeded" | "failed" | "cancelled" | "timed_out";
|
||||
type RunSessionOutcome = "succeeded" | "interrupted" | "failed" | "cancelled" | "timed_out";
|
||||
|
||||
const HERMES_ADAPTER_TYPE = "hermes_local";
|
||||
const HERMES_SESSION_ID_REGEX = /^(?:\d{8}_\d{6}_[A-Za-z0-9_-]{4,}|[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12})$/;
|
||||
@@ -7632,6 +7630,27 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
agent: typeof agents.$inferSelect,
|
||||
now: Date,
|
||||
) {
|
||||
const existingRetry = await db
|
||||
.select()
|
||||
.from(heartbeatRuns)
|
||||
.where(and(eq(heartbeatRuns.companyId, run.companyId), eq(heartbeatRuns.retryOfRunId, run.id)))
|
||||
.orderBy(asc(heartbeatRuns.createdAt))
|
||||
.limit(1)
|
||||
.then((rows) => rows[0] ?? null);
|
||||
if (existingRetry) {
|
||||
await appendRunEvent(run, await nextRunEventSeq(run.id), {
|
||||
eventType: "lifecycle",
|
||||
stream: "system",
|
||||
level: "warn",
|
||||
message: "Process-loss retry already exists; skipping duplicate retry enqueue",
|
||||
payload: {
|
||||
retryRunId: existingRetry.id,
|
||||
retryRunStatus: existingRetry.status,
|
||||
},
|
||||
});
|
||||
return existingRetry;
|
||||
}
|
||||
|
||||
const invokability = await getAgentInvokability(agent);
|
||||
if (!invokability.invokable) {
|
||||
await appendRunEvent(run, await nextRunEventSeq(run.id), {
|
||||
@@ -7750,6 +7769,104 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
return queued;
|
||||
}
|
||||
|
||||
async function drainRunningRunsForShutdown(signal: "SIGINT" | "SIGTERM", now = new Date()) {
|
||||
const activeRuns = await db
|
||||
.select({
|
||||
run: heartbeatRuns,
|
||||
agent: agents,
|
||||
})
|
||||
.from(heartbeatRuns)
|
||||
.innerJoin(agents, eq(heartbeatRuns.agentId, agents.id))
|
||||
.where(eq(heartbeatRuns.status, "running"));
|
||||
|
||||
const interruptedRunIds: string[] = [];
|
||||
const retryRunIds: string[] = [];
|
||||
|
||||
for (const { run, agent } of activeRuns) {
|
||||
const running = runningProcesses.get(run.id);
|
||||
try {
|
||||
if (running) {
|
||||
await terminateHeartbeatRunProcess({
|
||||
pid: running.child.pid ?? run.processPid,
|
||||
processGroupId: running.processGroupId ?? run.processGroupId,
|
||||
graceMs: Math.max(1, running.graceSec) * 1000,
|
||||
});
|
||||
} else if (run.processPid || run.processGroupId) {
|
||||
await terminateHeartbeatRunProcess({
|
||||
pid: run.processPid,
|
||||
processGroupId: run.processGroupId,
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
runningProcesses.delete(run.id);
|
||||
}
|
||||
|
||||
const message = `Interrupted by graceful server shutdown (${signal}); retry queued for restart recovery`;
|
||||
const interruptedStatus = await setRunStatusIfRunning(run.id, "interrupted", {
|
||||
finishedAt: now,
|
||||
error: message,
|
||||
errorCode: "server_shutdown_interrupted",
|
||||
signal,
|
||||
resultJson: mergeRunStopMetadataForAgent(agent, "interrupted", {
|
||||
resultJson: parseObject(run.resultJson),
|
||||
errorCode: "server_shutdown_interrupted",
|
||||
errorMessage: message,
|
||||
}),
|
||||
});
|
||||
if (!interruptedStatus.updated || !interruptedStatus.run) continue;
|
||||
let interrupted = interruptedStatus.run;
|
||||
await setWakeupStatus(run.wakeupRequestId, "cancelled", {
|
||||
finishedAt: now,
|
||||
error: null,
|
||||
});
|
||||
interrupted = await classifyAndPersistRunLiveness(interrupted, parseObject(interrupted.resultJson)) ?? interrupted;
|
||||
|
||||
await releaseEnvironmentLeasesForRun({
|
||||
runId: interrupted.id,
|
||||
companyId: interrupted.companyId,
|
||||
agentId: interrupted.agentId,
|
||||
status: interrupted.status,
|
||||
failureReason: interrupted.error ?? undefined,
|
||||
});
|
||||
|
||||
const retry = await enqueueProcessLossRetry(interrupted, agent, now);
|
||||
if (!retry) {
|
||||
await releaseIssueExecutionAndPromote(interrupted);
|
||||
} else {
|
||||
retryRunIds.push(retry.id);
|
||||
}
|
||||
|
||||
await appendRunEvent(interrupted, await nextRunEventSeq(interrupted.id), {
|
||||
eventType: "lifecycle",
|
||||
stream: "system",
|
||||
level: "warn",
|
||||
message,
|
||||
payload: {
|
||||
signal,
|
||||
...(run.processPid ? { processPid: run.processPid } : {}),
|
||||
...(run.processGroupId ? { processGroupId: run.processGroupId } : {}),
|
||||
...(retry ? { retryRunId: retry.id } : {}),
|
||||
},
|
||||
});
|
||||
|
||||
await finalizeAgentStatus(run.agentId, "interrupted", message);
|
||||
interruptedRunIds.push(interrupted.id);
|
||||
}
|
||||
|
||||
if (interruptedRunIds.length > 0) {
|
||||
logger.warn(
|
||||
{ signal, interrupted: interruptedRunIds.length, interruptedRunIds, retryRunIds },
|
||||
"interrupted running heartbeat runs for graceful shutdown",
|
||||
);
|
||||
}
|
||||
|
||||
return {
|
||||
interrupted: interruptedRunIds.length,
|
||||
interruptedRunIds,
|
||||
retryRunIds,
|
||||
};
|
||||
}
|
||||
|
||||
type ScheduledRetryGate =
|
||||
| { allowed: true }
|
||||
| {
|
||||
@@ -8160,7 +8277,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
? readTransientRecoveryContractFromRun(run)
|
||||
: null;
|
||||
const codexTransientFallbackMode =
|
||||
agent.adapterType === "codex_local" && transientRecovery
|
||||
agent.adapterType === "codex_local" && transientRecovery?.errorFamily === "transient_upstream"
|
||||
? resolveCodexTransientFallbackMode(nextAttempt)
|
||||
: null;
|
||||
const transientRetryNotBefore = transientRecovery?.retryNotBefore ?? null;
|
||||
@@ -8257,6 +8374,9 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
scheduledRetryAttempt: schedule.attempt,
|
||||
scheduledRetryAt: schedule.dueAt.toISOString(),
|
||||
...(transientRetryNotBefore ? { transientRetryNotBefore: transientRetryNotBefore.toISOString() } : {}),
|
||||
...(transientRecovery?.errorFamily === "provider_quota" && transientRetryNotBefore
|
||||
? { providerQuotaRetryNotBefore: transientRetryNotBefore.toISOString() }
|
||||
: {}),
|
||||
...(codexTransientFallbackMode ? { codexTransientFallbackMode } : {}),
|
||||
}, "normal_model");
|
||||
const responsibleUserId = await resolveResponsibleUserIdForRunContext(run, retryContextSnapshot);
|
||||
@@ -8428,6 +8548,9 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
scheduledRetryAttempt: schedule.attempt,
|
||||
scheduledRetryAt: schedule.dueAt.toISOString(),
|
||||
...(transientRetryNotBefore ? { transientRetryNotBefore: transientRetryNotBefore.toISOString() } : {}),
|
||||
...(transientRecovery?.errorFamily === "provider_quota" && transientRetryNotBefore
|
||||
? { providerQuotaRetryNotBefore: transientRetryNotBefore.toISOString() }
|
||||
: {}),
|
||||
...(codexTransientFallbackMode ? { codexTransientFallbackMode } : {}),
|
||||
}, "normal_model"),
|
||||
status: "queued",
|
||||
@@ -8551,6 +8674,9 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
baseDelayMs: schedule.baseDelayMs,
|
||||
delayMs: schedule.delayMs,
|
||||
...(transientRetryNotBefore ? { transientRetryNotBefore: transientRetryNotBefore.toISOString() } : {}),
|
||||
...(transientRecovery?.errorFamily === "provider_quota" && transientRetryNotBefore
|
||||
? { providerQuotaRetryNotBefore: transientRetryNotBefore.toISOString() }
|
||||
: {}),
|
||||
...(codexTransientFallbackMode ? { codexTransientFallbackMode } : {}),
|
||||
},
|
||||
});
|
||||
@@ -9419,8 +9545,9 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
|
||||
async function finalizeAgentStatus(
|
||||
agentId: string,
|
||||
outcome: "succeeded" | "failed" | "cancelled" | "timed_out",
|
||||
outcome: "succeeded" | "interrupted" | "failed" | "cancelled" | "timed_out",
|
||||
failureReason?: string | null,
|
||||
options?: { keepIdleOnFailure?: boolean },
|
||||
) {
|
||||
const existing = await getAgent(agentId);
|
||||
if (!existing) return;
|
||||
@@ -9435,7 +9562,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
const nextStatus =
|
||||
runningCount > 0
|
||||
? "running"
|
||||
: outcome === "succeeded" || outcome === "cancelled"
|
||||
: outcome === "succeeded" || outcome === "interrupted" || outcome === "cancelled" || (outcome === "failed" && options?.keepIdleOnFailure)
|
||||
? "idle"
|
||||
: "error";
|
||||
|
||||
@@ -9477,7 +9604,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
|
||||
function mergeRunStopMetadataForAgent(
|
||||
agent: Pick<typeof agents.$inferSelect, "adapterType" | "adapterConfig">,
|
||||
outcome: "succeeded" | "failed" | "cancelled" | "timed_out",
|
||||
outcome: "succeeded" | "interrupted" | "failed" | "cancelled" | "timed_out",
|
||||
options?: {
|
||||
resultJson?: Record<string, unknown> | null;
|
||||
errorCode?: string | null;
|
||||
@@ -10832,7 +10959,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
});
|
||||
const resolvedProjectId = executionWorkspace.projectId ?? issueRef?.projectId ?? executionProjectId ?? null;
|
||||
const resolvedProjectWorkspaceId = issueRef?.projectWorkspaceId ?? resolvedWorkspace.workspaceId ?? null;
|
||||
let persistedExecutionWorkspace = null;
|
||||
let persistedExecutionWorkspace: ExecutionWorkspace | null = null;
|
||||
const nextExecutionWorkspaceMetadata = mergeExecutionWorkspaceMetadataForPersistence({
|
||||
existingMetadata: resolvedWorkspaceReusePolicy.shouldRestoreExistingWorkspace
|
||||
? reusableExistingExecutionWorkspace?.metadata ?? null
|
||||
@@ -10937,18 +11064,6 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
if (persistedExecutionWorkspace && pendingForwardBranchReconcile) {
|
||||
await workspaceOperationRecorder.attachExecutionWorkspaceId(persistedExecutionWorkspace.id);
|
||||
const reconcileResult = await reconcilePendingForwardBranchAfterPersistence({
|
||||
db,
|
||||
executionWorkspaceId: persistedExecutionWorkspace.id,
|
||||
pending: pendingForwardBranchReconcile,
|
||||
heartbeatRunId: run.id,
|
||||
reconcileOperationPhase: "worktree_prepare",
|
||||
recorder: workspaceOperationRecorder,
|
||||
});
|
||||
persistedExecutionWorkspace = reconcileResult.workspace;
|
||||
}
|
||||
await workspaceOperationRecorder.attachExecutionWorkspaceId(persistedExecutionWorkspace?.id ?? null);
|
||||
await recordWorkspaceConfigFreshnessOperation({
|
||||
recorder: workspaceOperationRecorder,
|
||||
@@ -11592,7 +11707,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
let inspection = branchInspection.inspection;
|
||||
const initialManagedGitWorktreeBranch = formatManagedGitWorktreeBranchInspection(inspection);
|
||||
if (!inspection.valid && inspection.reasonCode === "branch_mismatch" && inspection.repoRoot) {
|
||||
let reconciledBranchName: string | null = null;
|
||||
let repairedExpectedBranchName = inspection.expectedBranchName;
|
||||
try {
|
||||
const coherence = await ensureGitWorktreeBranchCoherent({
|
||||
db,
|
||||
@@ -11612,11 +11727,14 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
heartbeatRunId: run.id,
|
||||
enableWorkspaceBranchReconcileForward:
|
||||
resolvedInstanceSettings.experimental.enableWorkspaceBranchReconcileForward,
|
||||
persistForwardReconcile: false,
|
||||
reconcileOperationPhase: "workspace_finalize",
|
||||
recorder: workspaceOperationRecorder,
|
||||
});
|
||||
if (coherence.reconciledForward && coherence.branchName) {
|
||||
reconciledBranchName = coherence.branchName;
|
||||
if (coherence.branchName && coherence.branchName !== branchInspection.workspaceRecord.branchName) {
|
||||
repairedExpectedBranchName = coherence.branchName;
|
||||
executionWorkspace.branchName = coherence.branchName;
|
||||
executionWorkspace.warnings.push(...coherence.warnings);
|
||||
}
|
||||
} catch (repairErr) {
|
||||
const workspaceValidationFailure = isWorkspaceValidationFailure(repairErr) ? repairErr : null;
|
||||
@@ -11654,7 +11772,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
|
||||
const repairedInspection = await inspectManagedGitWorktreeBranch({
|
||||
worktreePath: inspection.worktreePath,
|
||||
expectedBranchName: reconciledBranchName ?? inspection.expectedBranchName,
|
||||
expectedBranchName: repairedExpectedBranchName,
|
||||
repoRoot: inspection.repoRoot,
|
||||
});
|
||||
finalizeBranchRepairMetadata = {
|
||||
@@ -12165,6 +12283,11 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
agent.id,
|
||||
outcome,
|
||||
outcome === "succeeded" ? null : (adapterResult.errorMessage ?? null),
|
||||
{
|
||||
keepIdleOnFailure:
|
||||
outcome === "failed" &&
|
||||
(finalizedRun ? readHeartbeatRunErrorFamily(finalizedRun) === "provider_quota" : runErrorCode === "provider_quota"),
|
||||
},
|
||||
);
|
||||
} catch (err) {
|
||||
const message = redactCurrentUserText(
|
||||
@@ -14673,6 +14796,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
|
||||
reportRunActivity: clearDetachedRunWarning,
|
||||
|
||||
reapOrphanedRuns,
|
||||
drainRunningRunsForShutdown,
|
||||
|
||||
promoteDueScheduledRetries,
|
||||
retryScheduledRetryNow,
|
||||
|
||||
@@ -94,7 +94,7 @@ function extractPathCandidates(...texts: Array<string | null | undefined>) {
|
||||
|
||||
function inferMode(issue: IssueSummaryInput, run: RunSummaryInput) {
|
||||
if (issue.status === "done" || issue.status === "in_review") return "review";
|
||||
if (run.status === "failed" || run.status === "timed_out" || run.status === "cancelled") return "implementation";
|
||||
if (run.status === "failed" || run.status === "timed_out" || run.status === "cancelled" || run.status === "interrupted") return "implementation";
|
||||
if (issue.status === "backlog" || issue.status === "todo") return "plan";
|
||||
return "implementation";
|
||||
}
|
||||
|
||||
@@ -604,7 +604,7 @@ function sameRunLock(checkoutRunId: string | null, actorRunId: string | null) {
|
||||
return checkoutRunId == null;
|
||||
}
|
||||
|
||||
export const TERMINAL_HEARTBEAT_RUN_STATUSES = new Set(["succeeded", "failed", "cancelled", "timed_out"]);
|
||||
export const TERMINAL_HEARTBEAT_RUN_STATUSES = new Set(["succeeded", "interrupted", "failed", "cancelled", "timed_out"]);
|
||||
const ISSUE_LIST_DESCRIPTION_MAX_CHARS = 1200;
|
||||
const ISSUE_LIST_DESCRIPTION_MAX_BYTES = ISSUE_LIST_DESCRIPTION_MAX_CHARS * 4;
|
||||
|
||||
|
||||
@@ -2644,7 +2644,7 @@ export function buildHostServices(
|
||||
// Track the subscription so it can be cleaned up on dispose() if the run
|
||||
// never reaches a terminal status (hang, crash, network partition).
|
||||
if (notifyWorker) {
|
||||
const TERMINAL_STATUSES = new Set(["succeeded", "failed", "cancelled", "timed_out"]);
|
||||
const TERMINAL_STATUSES = new Set(["succeeded", "interrupted", "failed", "cancelled", "timed_out"]);
|
||||
|
||||
const cleanup = () => {
|
||||
unsubscribe();
|
||||
|
||||
@@ -31,7 +31,7 @@ export const DEFAULT_PRODUCTIVITY_REVIEW_MAX_REFRESH_COMMENTS = 3;
|
||||
export const DEFAULT_PRODUCTIVITY_REVIEW_CREATION_WINDOW_MS = 24 * 60 * 60 * 1000;
|
||||
export const DEFAULT_PRODUCTIVITY_REVIEW_MAX_CREATIONS_PER_WINDOW = 3;
|
||||
|
||||
const TERMINAL_RUN_STATUSES = ["succeeded", "failed", "cancelled", "timed_out"] as const;
|
||||
const TERMINAL_RUN_STATUSES = ["succeeded", "interrupted", "failed", "cancelled", "timed_out"] as const;
|
||||
const ACTIVE_RUN_STATUSES = ["queued", "running", "scheduled_retry"] as const;
|
||||
const MAX_CANDIDATE_ISSUES = 250;
|
||||
const MAX_RUNS_FOR_STREAK = 100;
|
||||
|
||||
@@ -70,7 +70,7 @@ import {
|
||||
import { isAutomaticRecoverySuppressedByPauseHold } from "./pause-hold-guard.js";
|
||||
|
||||
const EXECUTION_PATH_HEARTBEAT_RUN_STATUSES = ["queued", "running", "scheduled_retry"] as const;
|
||||
const UNSUCCESSFUL_HEARTBEAT_RUN_TERMINAL_STATUSES = ["failed", "cancelled", "timed_out"] as const;
|
||||
const UNSUCCESSFUL_HEARTBEAT_RUN_TERMINAL_STATUSES = ["interrupted", "failed", "cancelled", "timed_out"] as const;
|
||||
export const ACTIVE_RUN_OUTPUT_SUSPICION_THRESHOLD_MS = 60 * 60 * 1000;
|
||||
export const ACTIVE_RUN_OUTPUT_CRITICAL_THRESHOLD_MS = 4 * 60 * 60 * 1000;
|
||||
export const ACTIVE_RUN_OUTPUT_CONTINUE_REARM_MS = 30 * 60 * 1000;
|
||||
@@ -229,6 +229,7 @@ const TRANSIENT_INFRA_CONTINUATION_ERROR_CODES = new Set<string>([
|
||||
"adapter_failed",
|
||||
"codex_transient_upstream",
|
||||
"claude_transient_upstream",
|
||||
"provider_quota",
|
||||
"timeout",
|
||||
]);
|
||||
|
||||
|
||||
@@ -309,6 +309,10 @@ export function classifyRunLiveness(input: RunLivenessClassificationInput): RunL
|
||||
actionability,
|
||||
});
|
||||
|
||||
if (input.runStatus === "interrupted") {
|
||||
return output("needs_followup", input.errorCode ? `Run interrupted (${input.errorCode})` : "Run interrupted");
|
||||
}
|
||||
|
||||
if (input.runStatus !== "succeeded") {
|
||||
return output("failed", input.errorCode ? `Run ended with ${input.runStatus} (${input.errorCode})` : `Run ended with ${input.runStatus}`);
|
||||
}
|
||||
|
||||
@@ -28,7 +28,7 @@ const TASK_WATCHDOG_SUBTREE_MAX_DEPTH = 100;
|
||||
const TASK_WATCHDOG_LIVE_RUN_STATUSES = ["queued", "running", "scheduled_retry"] as const;
|
||||
const TASK_WATCHDOG_WAKE_REQUEST_STATUSES = ["queued", "deferred_issue_execution"] as const;
|
||||
const TASK_WATCHDOG_TERMINAL_ISSUE_STATUSES = ["done", "cancelled"] as const;
|
||||
const TASK_WATCHDOG_TERMINAL_RUN_STATUSES = ["succeeded", "failed", "cancelled", "timed_out"] as const;
|
||||
const TASK_WATCHDOG_TERMINAL_RUN_STATUSES = ["succeeded", "interrupted", "failed", "cancelled", "timed_out"] as const;
|
||||
// Grace window after an issue is created/assigned during which its first
|
||||
// assignment run/wake may have been enqueued but is not yet visible to a
|
||||
// watchdog evaluation (the eval can race the issue's own assignment run).
|
||||
|
||||
@@ -665,6 +665,7 @@ type GitWorktreeBranchCoherenceResult = {
|
||||
branchName: string | null;
|
||||
reconciledForward: boolean;
|
||||
pendingForwardBranchReconcile?: PendingForwardBranchReconcile | null;
|
||||
warnings: string[];
|
||||
};
|
||||
|
||||
export type PendingForwardBranchReconcile = {
|
||||
@@ -791,9 +792,30 @@ async function inspectGitWorktreeBranchIncoherence(input: {
|
||||
sameHead,
|
||||
ancestryVerdict,
|
||||
});
|
||||
const eligible = cleanliness === "clean" && expectedBranchExists && sameHead && registeredBranchMatchesHead;
|
||||
const canCheckoutRecordedBranch =
|
||||
cleanliness === "clean" && expectedBranchExists && sameHead && registeredBranchMatchesHead;
|
||||
const canAdoptForwardActualBranch =
|
||||
cleanliness === "clean" &&
|
||||
expectedBranchExists &&
|
||||
actualBranchExists === true &&
|
||||
ancestryVerdict === "ancestor" &&
|
||||
!sameHead &&
|
||||
registeredBranchMatchesHead;
|
||||
const canAttachRecordedBranchToDetachedHead =
|
||||
cleanliness === "clean" &&
|
||||
expectedBranchExists &&
|
||||
input.actualBranchName === null &&
|
||||
ancestryVerdict === "ancestor" &&
|
||||
!sameHead &&
|
||||
registeredBranchMatchesHead;
|
||||
const eligible =
|
||||
canCheckoutRecordedBranch || canAdoptForwardActualBranch || canAttachRecordedBranchToDetachedHead;
|
||||
const safeRepairReason = eligible
|
||||
? "clean worktree and expected branch points at the current HEAD"
|
||||
? canCheckoutRecordedBranch
|
||||
? "clean worktree and expected branch points at the current HEAD"
|
||||
: canAdoptForwardActualBranch
|
||||
? "clean worktree and checked-out branch is forward of the recorded branch"
|
||||
: "clean detached worktree HEAD is forward of the recorded branch"
|
||||
: cleanliness !== "clean"
|
||||
? "worktree is not clean"
|
||||
: !registered
|
||||
@@ -1025,17 +1047,18 @@ export async function ensureGitWorktreeBranchCoherent(input: {
|
||||
actualBranchName?: string | null;
|
||||
heartbeatRunId?: string | null;
|
||||
enableWorkspaceBranchReconcileForward?: boolean;
|
||||
persistForwardReconcile?: boolean;
|
||||
reconcileOperationPhase?: "worktree_prepare" | "workspace_finalize";
|
||||
recorder?: WorkspaceOperationRecorder | null;
|
||||
}): Promise<GitWorktreeBranchCoherenceResult> {
|
||||
const expectedBranchName = input.expectedBranchName?.trim();
|
||||
if (!expectedBranchName) return { branchName: null, reconciledForward: false };
|
||||
if (!expectedBranchName) return { branchName: null, reconciledForward: false, warnings: [] };
|
||||
|
||||
const currentBranch = input.actualBranchName !== undefined
|
||||
? input.actualBranchName
|
||||
: await runGit(["symbolic-ref", "--quiet", "--short", "HEAD"], input.worktreePath).catch(() => null);
|
||||
if (currentBranch === expectedBranchName) {
|
||||
return { branchName: expectedBranchName, reconciledForward: false };
|
||||
return { branchName: expectedBranchName, reconciledForward: false, warnings: [] };
|
||||
}
|
||||
|
||||
const evidence = await inspectGitWorktreeBranchIncoherence({
|
||||
@@ -1053,7 +1076,7 @@ export async function ensureGitWorktreeBranchCoherent(input: {
|
||||
currentBranch
|
||||
) {
|
||||
const reason = "Automatic forward reconciliation: recorded branch is an ancestor of the checked-out branch.";
|
||||
if (input.executionWorkspaceId) {
|
||||
if (input.executionWorkspaceId && input.persistForwardReconcile !== false) {
|
||||
if (!input.db) {
|
||||
evidence.safeRepair.reason = "forward reconciliation requires database access to update the execution workspace record";
|
||||
throw branchIncoherenceValidationFailure(evidence);
|
||||
@@ -1107,7 +1130,7 @@ export async function ensureGitWorktreeBranchCoherent(input: {
|
||||
auditCommentId: result.auditCommentId,
|
||||
recoveryActionId: result.recoveryAction?.id ?? null,
|
||||
});
|
||||
return { branchName: result.inspection.toBranch, reconciledForward: true };
|
||||
return { branchName: result.inspection.toBranch, reconciledForward: true, warnings: [] };
|
||||
} catch (error) {
|
||||
evidence.safeRepair.reason =
|
||||
`forward reconciliation failed: ${error instanceof Error ? error.message : String(error)}`;
|
||||
@@ -1122,6 +1145,7 @@ export async function ensureGitWorktreeBranchCoherent(input: {
|
||||
return {
|
||||
branchName: currentBranch,
|
||||
reconciledForward: true,
|
||||
warnings: [],
|
||||
pendingForwardBranchReconcile: {
|
||||
recordedBranchName: expectedBranchName,
|
||||
adoptedBranchName: currentBranch,
|
||||
@@ -1136,6 +1160,75 @@ export async function ensureGitWorktreeBranchCoherent(input: {
|
||||
}
|
||||
|
||||
evidence.safeRepair.attempted = true;
|
||||
const warningPrefix =
|
||||
`Execution workspace branch metadata was self-healed from "${expectedBranchName}" to "${formatBranchForMessage(currentBranch)}" at ${input.worktreePath}.`;
|
||||
if (
|
||||
currentBranch &&
|
||||
evidence.provenance.actualBranchExists === true &&
|
||||
evidence.provenance.ancestryVerdict === "ancestor" &&
|
||||
!evidence.provenance.sameHead
|
||||
) {
|
||||
evidence.safeRepair.succeeded = true;
|
||||
evidence.safeRepair.reason = "clean worktree adopted the checked-out branch because it is forward of the recorded branch";
|
||||
return {
|
||||
branchName: currentBranch,
|
||||
reconciledForward: false,
|
||||
warnings: [
|
||||
`${warningPrefix} The checked-out branch contains the recorded branch plus newer commits, so Paperclip adopted it for subsequent runs.`,
|
||||
],
|
||||
};
|
||||
}
|
||||
|
||||
if (
|
||||
currentBranch === null &&
|
||||
evidence.provenance.ancestryVerdict === "ancestor" &&
|
||||
!evidence.provenance.sameHead &&
|
||||
evidence.provenance.actualHeadSha
|
||||
) {
|
||||
try {
|
||||
await recordGitOperation(input.recorder, {
|
||||
phase: "worktree_prepare",
|
||||
args: ["checkout", "-B", expectedBranchName, evidence.provenance.actualHeadSha],
|
||||
cwd: input.worktreePath,
|
||||
metadata: {
|
||||
repoRoot: input.repoRoot,
|
||||
worktreePath: input.worktreePath,
|
||||
expectedBranchName,
|
||||
actualBranchName: currentBranch,
|
||||
branchIncoherenceRepair: true,
|
||||
detachedHeadRepair: true,
|
||||
fingerprint: evidence.fingerprint,
|
||||
sourceIssueId: evidence.sourceIssueId,
|
||||
executionWorkspaceId: evidence.executionWorkspaceId,
|
||||
},
|
||||
successMessage: `Reattached detached git worktree HEAD at ${input.worktreePath} to ${expectedBranchName}\n`,
|
||||
failureLabel: `git checkout -B ${expectedBranchName} ${formatShortSha(evidence.provenance.actualHeadSha)}`,
|
||||
});
|
||||
} catch (error) {
|
||||
evidence.safeRepair.succeeded = false;
|
||||
evidence.safeRepair.reason = `safe detached HEAD reattachment failed: ${error instanceof Error ? error.message : String(error)}`;
|
||||
throw branchIncoherenceValidationFailure(evidence);
|
||||
}
|
||||
|
||||
const repairedBranch = await runGit(["symbolic-ref", "--quiet", "--short", "HEAD"], input.worktreePath)
|
||||
.catch(() => null);
|
||||
if (repairedBranch !== expectedBranchName) {
|
||||
evidence.safeRepair.succeeded = false;
|
||||
evidence.safeRepair.reason = `reattach completed but HEAD is ${formatBranchForMessage(repairedBranch)}`;
|
||||
throw branchIncoherenceValidationFailure(evidence);
|
||||
}
|
||||
|
||||
evidence.safeRepair.succeeded = true;
|
||||
evidence.safeRepair.reason = "clean detached worktree HEAD was reattached to the recorded branch";
|
||||
return {
|
||||
branchName: expectedBranchName,
|
||||
reconciledForward: false,
|
||||
warnings: [
|
||||
`${warningPrefix} The detached HEAD contained the recorded branch plus newer commits, so Paperclip moved the recorded branch to that HEAD.`,
|
||||
],
|
||||
};
|
||||
}
|
||||
|
||||
try {
|
||||
await recordGitOperation(input.recorder, {
|
||||
phase: "worktree_prepare",
|
||||
@@ -1170,7 +1263,13 @@ export async function ensureGitWorktreeBranchCoherent(input: {
|
||||
|
||||
evidence.safeRepair.succeeded = true;
|
||||
evidence.safeRepair.reason = "clean worktree checked out the recorded branch";
|
||||
return { branchName: expectedBranchName, reconciledForward: false };
|
||||
return {
|
||||
branchName: expectedBranchName,
|
||||
reconciledForward: false,
|
||||
warnings: [
|
||||
`Execution workspace branch metadata was self-healed by checking out recorded branch "${expectedBranchName}" at ${input.worktreePath}.`,
|
||||
],
|
||||
};
|
||||
}
|
||||
|
||||
// Resolve the authoritative base ref for a fresh worktree. A configured local
|
||||
@@ -1919,12 +2018,12 @@ export async function realizeExecutionWorkspace(input: {
|
||||
|
||||
await fs.mkdir(worktreeParentDir, { recursive: true });
|
||||
|
||||
async function reuseExistingWorktree(reusablePath: string) {
|
||||
async function reuseExistingWorktree(reusablePath: string, effectiveBranchName = branchName, extraWarnings: string[] = []) {
|
||||
const refresh = currentBaseRefSha
|
||||
? await refreshUnstartedWorktreeToBase({
|
||||
repoRoot,
|
||||
worktreePath: reusablePath,
|
||||
branchName,
|
||||
branchName: effectiveBranchName,
|
||||
baseRef,
|
||||
currentBaseRefSha,
|
||||
recorder: input.recorder ?? null,
|
||||
@@ -1945,7 +2044,7 @@ export async function realizeExecutionWorkspace(input: {
|
||||
metadata: {
|
||||
repoRoot,
|
||||
worktreePath: reusablePath,
|
||||
branchName,
|
||||
branchName: effectiveBranchName,
|
||||
baseRef,
|
||||
currentBaseRefSha: baseDrift.currentBaseRefSha,
|
||||
branchBaseRefSha: baseDrift.branchBaseRefSha,
|
||||
@@ -1964,7 +2063,7 @@ export async function realizeExecutionWorkspace(input: {
|
||||
base: input.base,
|
||||
repoRoot,
|
||||
worktreePath: reusablePath,
|
||||
branchName,
|
||||
branchName: effectiveBranchName,
|
||||
issue: input.issue,
|
||||
agent: input.agent,
|
||||
created: false,
|
||||
@@ -1975,9 +2074,9 @@ export async function realizeExecutionWorkspace(input: {
|
||||
repoRef: baseRef,
|
||||
strategy: "git_worktree" as const,
|
||||
cwd: reusablePath,
|
||||
branchName,
|
||||
branchName: effectiveBranchName,
|
||||
worktreePath: reusablePath,
|
||||
warnings: [...baseRefreshWarnings, ...baseDrift.warnings],
|
||||
warnings: [...extraWarnings, ...baseRefreshWarnings, ...baseDrift.warnings],
|
||||
created: false,
|
||||
baseRefSha: refresh.baseRefSha ?? baseDrift.branchBaseRefSha ?? baseDrift.currentBaseRefSha,
|
||||
pendingForwardBranchReconcile,
|
||||
@@ -2004,35 +2103,43 @@ export async function realizeExecutionWorkspace(input: {
|
||||
reconcileOperationPhase: "worktree_prepare",
|
||||
recorder: input.recorder ?? null,
|
||||
});
|
||||
if (coherence.reconciledForward && coherence.branchName) {
|
||||
branchName = coherence.branchName;
|
||||
const effectiveBranchName = coherence.branchName ?? branchName;
|
||||
if (coherence.reconciledForward) {
|
||||
branchName = effectiveBranchName;
|
||||
pendingForwardBranchReconcile = coherence.pendingForwardBranchReconcile ?? null;
|
||||
}
|
||||
return await validateLinkedGitWorktree({
|
||||
const nextValidation = await validateLinkedGitWorktree({
|
||||
repoRoot,
|
||||
worktreePath: reusablePath,
|
||||
expectedBranchName: branchName,
|
||||
expectedBranchName: effectiveBranchName,
|
||||
}).catch(() => null);
|
||||
return {
|
||||
validation: nextValidation,
|
||||
branchName: effectiveBranchName,
|
||||
warnings: coherence.warnings,
|
||||
};
|
||||
}
|
||||
return validation;
|
||||
return { validation, branchName, warnings: [] };
|
||||
}
|
||||
|
||||
const existingWorktree = await directoryExists(worktreePath);
|
||||
if (existingWorktree) {
|
||||
const validation = await validateReusableWorktree(worktreePath);
|
||||
if (validation?.valid) {
|
||||
return await reuseExistingWorktree(worktreePath);
|
||||
const reusable = await validateReusableWorktree(worktreePath);
|
||||
if (reusable.validation?.valid) {
|
||||
return await reuseExistingWorktree(worktreePath, reusable.branchName, reusable.warnings);
|
||||
}
|
||||
const validation = reusable.validation;
|
||||
const reason = validation && !validation.valid ? ` (${validation.reason})` : "";
|
||||
throw new Error(`Configured worktree path "${worktreePath}" already exists and is not a reusable git worktree${reason}.`);
|
||||
}
|
||||
|
||||
const registeredBranchWorktree = await findRegisteredGitWorktreeByBranch(repoRoot, branchName);
|
||||
if (registeredBranchWorktree) {
|
||||
const validation = await validateReusableWorktree(registeredBranchWorktree);
|
||||
if (validation?.valid) {
|
||||
return await reuseExistingWorktree(registeredBranchWorktree);
|
||||
const reusable = await validateReusableWorktree(registeredBranchWorktree);
|
||||
if (reusable.validation?.valid) {
|
||||
return await reuseExistingWorktree(registeredBranchWorktree, reusable.branchName, reusable.warnings);
|
||||
}
|
||||
const validation = reusable.validation;
|
||||
const reason = validation && !validation.valid ? ` (${validation.reason})` : "";
|
||||
throw new Error(`Registered worktree for branch "${branchName}" at "${registeredBranchWorktree}" is not reusable${reason}.`);
|
||||
}
|
||||
@@ -2167,6 +2274,7 @@ export async function ensurePersistedExecutionWorkspaceAvailable(input: {
|
||||
if (await directoryExists(cwd)) {
|
||||
const reuseBaseRef = input.workspace.baseRef ?? input.base.repoRef ?? null;
|
||||
const reuseWorktreePath = realized.worktreePath ?? cwd;
|
||||
const repairWarnings: string[] = [];
|
||||
if (await isGitCheckout(reuseWorktreePath)) {
|
||||
const coherence = await ensureGitWorktreeBranchCoherent({
|
||||
db: input.db ?? null,
|
||||
@@ -2177,12 +2285,17 @@ export async function ensurePersistedExecutionWorkspaceAvailable(input: {
|
||||
executionWorkspaceId: input.workspace.id ?? null,
|
||||
heartbeatRunId: input.heartbeatRunId ?? null,
|
||||
enableWorkspaceBranchReconcileForward: input.enableWorkspaceBranchReconcileForward === true,
|
||||
persistForwardReconcile: false,
|
||||
reconcileOperationPhase: "worktree_prepare",
|
||||
recorder: input.recorder ?? null,
|
||||
});
|
||||
if (coherence.reconciledForward && coherence.branchName) {
|
||||
if (coherence.branchName) {
|
||||
realized.branchName = coherence.branchName;
|
||||
}
|
||||
if (coherence.reconciledForward) {
|
||||
realized.pendingForwardBranchReconcile = coherence.pendingForwardBranchReconcile ?? null;
|
||||
}
|
||||
repairWarnings.push(...coherence.warnings);
|
||||
}
|
||||
const validation = await validateLinkedGitWorktree({
|
||||
repoRoot,
|
||||
@@ -2224,7 +2337,7 @@ export async function ensurePersistedExecutionWorkspaceAvailable(input: {
|
||||
recordedBaseRefSha,
|
||||
skipRefresh: true,
|
||||
});
|
||||
realized.warnings = [...baseRefreshWarnings, ...baseDrift.warnings];
|
||||
realized.warnings = [...repairWarnings, ...baseRefreshWarnings, ...baseDrift.warnings];
|
||||
realized.baseRefSha = refresh.baseRefSha ?? recordedBaseRefSha ?? baseDrift.branchBaseRefSha ?? baseDrift.currentBaseRefSha;
|
||||
if (provisionCommand) {
|
||||
await provisionExecutionWorktree({
|
||||
@@ -4134,6 +4247,9 @@ export function buildWorkspaceReadyComment(input: {
|
||||
if (input.workspace.worktreePath && input.workspace.worktreePath !== input.workspace.cwd) {
|
||||
lines.push(`- Worktree: \`${input.workspace.worktreePath}\``);
|
||||
}
|
||||
for (const warning of input.workspace.warnings) {
|
||||
lines.push(`- Warning: ${warning}`);
|
||||
}
|
||||
for (const service of input.runtimeServices) {
|
||||
const detail = service.url ? `${service.serviceName}: ${service.url}` : `${service.serviceName}: running`;
|
||||
const suffix = service.reused ? " (reused)" : "";
|
||||
|
||||
Reference in new issue
Block a user