diff --git a/doc/observability.md b/doc/observability.md index 08dff65ea7..3aaee6c9e2 100644 --- a/doc/observability.md +++ b/doc/observability.md @@ -453,6 +453,19 @@ termination, remove the live execution control, or change the timeout. The context is sent only through the existing opt-in Sentry gate and never leaks into unrelated captures. +The context also samples the pending execution `phase` and `phaseElapsedMs` +when the Stop timer expires. An in-memory tracker belongs to the exact live +execution control; a replaced or missing owner reports `unknown` with a null +elapsed time. Labels come from a closed list and elapsed time uses a monotonic +clock, capped at one day. Nested scopes remain visible while their awaits are +pending, including ACP session close, transport stop, instruction collection, +workspace restore, and host instruction or lease cleanup. `phase_reporting` +means a step is waiting to write timing or teardown error diagnostics. `adapter_execution` and +`host_execution` are coarse labels for work outside those narrower scopes. +The tracker clears when its executor finishes. A phase is diagnostic context, +not evidence that a provider stopped or that files were recovered. No new +run-log or Telemetry event is emitted by this tracker. + The shared reporter also attaches bounded diagnostic contexts for both legacy and native runs: diff --git a/packages/adapter-utils/src/acpx-engine/execute.ts b/packages/adapter-utils/src/acpx-engine/execute.ts index a40d3a36f3..5c575ca8c5 100644 --- a/packages/adapter-utils/src/acpx-engine/execute.ts +++ b/packages/adapter-utils/src/acpx-engine/execute.ts @@ -1,4 +1,5 @@ import { cancellableSandboxStartup } from "./startup-cancellation.js"; +import { withAdapterExecutionPhase, type AdapterExecutionPhase } from "../execution-phase.js"; import fs from "node:fs/promises"; import fsSync from "node:fs"; import os from "node:os"; @@ -4079,9 +4080,9 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { : teardownErr instanceof Error ? teardownErr.message : String(teardownErr); - await ctx + await withAdapterExecutionPhase(ctx, "phase_reporting", () => ctx .onLog("stderr", `[paperclip] ACPX teardown step "${step}" failed: ${reason}\n`) - .catch(() => {}); + .catch(() => {})); }; // Emit one per-phase timing run-log event. It is not an OpenTelemetry // export and it is not a Telemetry event: it carries the phase name @@ -4092,15 +4093,15 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { emitRunPhaseTiming(ctx, phase, now() - startMs, outcome); // Time a settlement step and emit its phase timing on every path. A step // error still emits `failed` before it re-throws to the Phase 3 error policy. - const timedPhase = async (phase: string, run: () => Promise | void): Promise => { + const timedPhase = async (phase: AdapterExecutionPhase, run: () => Promise | void): Promise => { const start = now(); try { - await run(); + await withAdapterExecutionPhase(ctx, phase, run); } catch (error) { - await emitPhase(phase, start, "failed"); + await withAdapterExecutionPhase(ctx, "phase_reporting", () => emitPhase(phase, start, "failed")); throw error; } - await emitPhase(phase, start, "ok"); + await withAdapterExecutionPhase(ctx, "phase_reporting", () => emitPhase(phase, start, "ok")); }; // The turn the run started. It is hoisted to the run scope so the settlement // `endSession` step can cancel a running turn before it closes the runtime @@ -4860,7 +4861,7 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { if (ctx.signal?.aborted) armStopDeadline(); return { cancel: async (reason: string) => { - await turn.cancel({ reason }); + await withAdapterExecutionPhase(ctx, "cancel_turn", () => turn.cancel({ reason })); }, }; }; @@ -5314,7 +5315,7 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { : baseSettlement; // Cancel a running turn before the close (the turn-error path). if (settlement.cancelTurnReason && activeTurn) { - await activeTurn.cancel({ reason: settlement.cancelTurnReason }).catch(() => {}); + await withAdapterExecutionPhase(ctx, "cancel_turn", () => activeTurn!.cancel({ reason: settlement.cancelTurnReason! }).catch(() => {})); } const existing = warmHandles.get(prepared.sessionKey); // Re-read the duplex control-channel disposition here, at the @@ -5353,26 +5354,26 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { ) { // A matching warm entry closes through the warm store, which also // clears its idle timer and flushes its child stderr. - runtimeStopConfirmed = await closeWarmHandle({ + runtimeStopConfirmed = await withAdapterExecutionPhase(ctx, "close_session", () => closeWarmHandle({ handles: warmHandles, key: prepared.sessionKey, entry: existing, reason: settlement.reason, discardPersistentState: settlement.discardPersistentState, - }); + })); return; } const onCloseError = settlement.recordCloseError || ctx.signal?.aborted ? (closeErr: unknown) => recordTeardownError("runtime-close", closeErr) : () => {}; - await runtime + await withAdapterExecutionPhase(ctx, "close_session", () => runtime .close({ handle: settlement.handle, reason: settlement.reason, discardPersistentState: settlement.discardPersistentState, }) .then(() => { runtimeStopConfirmed = true; }) - .catch(onCloseError); + .catch(onCloseError)); if (settlement.dropWarmEntry && warmHandleMatches(existing, runtime, settlement.handle) && existing) { clearWarmHandleTimer(existing); warmHandles.delete(prepared.sessionKey); @@ -5403,18 +5404,18 @@ export function createAcpxEngineExecutor(deps: AcpxEngineExecutorOptions = {}) { // mutable provider files. Require the owned close or local OS handle. const stopped = runtimeStopConfirmed || (!prepared.processSessionBridge && capturedProcessExited(processIdentitySink?.localProcess)); try { - if (stopped) await ctx.onProviderStopped?.(); + if (stopped) await withAdapterExecutionPhase(ctx, "instruction_collection", () => ctx.onProviderStopped?.()); } catch { // Match ACP's fail-soft teardown policy without describing this as // a workspace restore failure or exposing paths from a raw error. await recordTeardownError("instruction-collection", new Error("Instruction collection failed after provider stop. No instruction save is claimed.")); } finally { - await runRuntimeSpan("sandbox.syncBack", async () => { + await withAdapterExecutionPhase(ctx, "workspace_restore", () => runRuntimeSpan("sandbox.syncBack", async () => { const restoreOutcome = await syncBackManagedHome(prepared); if (!restoreOutcome.ok) { workspaceRestoreFailureField = { workspaceRestoreFailure: restoreOutcome.code }; } - }); + })); } }), // The staging lease releases as the run's final act, AFTER the coordinator diff --git a/packages/adapter-utils/src/acpx-engine/settlement-characterization.test.ts b/packages/adapter-utils/src/acpx-engine/settlement-characterization.test.ts index 593fc644d3..ff74f3c8f3 100644 --- a/packages/adapter-utils/src/acpx-engine/settlement-characterization.test.ts +++ b/packages/adapter-utils/src/acpx-engine/settlement-characterization.test.ts @@ -292,6 +292,110 @@ describe("ACP settlement — Layer A: engine teardown orchestration", () => { vi.clearAllMocks(); }); + it.each(["close_session", "stop_transport", "instruction_collection", "workspace_restore", "phase_reporting", "close_error_reporting"])( + "keeps %s visible while its real settlement await is stalled", + async (blockedPhase) => { + const { stateDir, localCwd, executionTarget } = await setupRemoteSandbox(); + stubBridges(); + let resume!: () => void; + let entered!: () => void; + const gate = new Promise((resolve) => { resume = resolve; }); + const reached = new Promise((resolve) => { entered = resolve; }); + const block = () => { entered(); return gate; }; + const scopes = new Map(); + if (blockedPhase === "stop_transport") { + vi.mocked(startAdapterExecutionTargetPaperclipBridge).mockImplementation(async () => ({ env: {}, stop: block }) as never); + } + const execute = createAcpxEngineExecutor({ + createRuntime: () => ({ + ensureSession: async () => okHandle, + startTurn: () => blockedPhase === "close_error_reporting" ? throwingTurn() : completedTurn(), + close: blockedPhase === "close_session" ? block : async () => { + if (blockedPhase === "close_error_reporting") throw new Error("close fixture failed"); + }, + }) as never, + prepareRemoteManagedHome: async (input) => ({ + stagedRuntime: await input.stage([]), + teardown: async () => { + if (blockedPhase === "workspace_restore") await block(); + return { ok: true }; + }, + }), + }); + const execution = execute({ + runId: "pending-settlement-fixture", + ...remoteArgs(stateDir, localCwd, executionTarget), + onProviderStopped: blockedPhase === "instruction_collection" ? block : async () => {}, + onExecutionPhase: (phase: string) => { + const token = Symbol(); + scopes.set(token, phase); + return () => { scopes.delete(token); }; + }, + onEvent: async (event: { eventType: string; payload?: Record }) => { + if (blockedPhase === "phase_reporting" && event.eventType === "run.phase.timing" && event.payload?.phase === "end_session") await block(); + }, + onLog: async (_stream: string, text: string) => { + if (blockedPhase === "close_error_reporting" && text.includes("close fixture failed")) await block(); + }, + } as never); + try { + await reached; + expect([...scopes.values()].at(-1)).toBe(blockedPhase === "close_error_reporting" ? "phase_reporting" : blockedPhase); + } finally { + resume(); + expect((await execution).exitCode).toBe(blockedPhase === "close_error_reporting" ? 1 : 0); + } + expect(scopes.size).toBe(0); + }, + ); + + it("keeps the primary Stop cancellation visible until its await settles", async () => { + const { stateDir, localCwd, executionTarget } = await setupRemoteSandbox(); + stubBridges(); + const controller = new AbortController(); + let started!: () => void; + let cancelled!: () => void; + let resume!: () => void; + const turnStarted = new Promise((resolve) => { started = resolve; }); + const cancellationStarted = new Promise((resolve) => { cancelled = resolve; }); + const gate = new Promise((resolve) => { resume = resolve; }); + const scopes = new Map(); + const execute = createAcpxEngineExecutor({ + createRuntime: () => ({ + ensureSession: async () => okHandle, + startTurn: () => { + started(); + return { + events: (async function* () { await gate; })(), + result: gate.then(() => ({ status: "cancelled", stopReason: "cancelled" })), + cancel: async () => { cancelled(); await gate; }, + }; + }, + close: async () => {}, + }) as never, + }); + const execution = execute({ + runId: "pending-cancel-fixture", + ...remoteArgs(stateDir, localCwd, executionTarget), + signal: controller.signal, + onExecutionPhase: (phase: string) => { + const token = Symbol(); + scopes.set(token, phase); + return () => { scopes.delete(token); }; + }, + } as never); + try { + await turnStarted; + controller.abort(); + await cancellationStarted; + expect([...scopes.values()].at(-1)).toBe("cancel_turn"); + } finally { + resume(); + expect((await execution).exitCode).toBe(1); + } + expect(scopes.size).toBe(0); + }); + it("test_clean_completed_remote_teardown_runs_bridge_stop_then_sync_back_then_lease_release", async () => { // cleanupRemoteBridges (execute.ts:2306-2326) fixes the sub-order for a clean // exit: stop both bridges (allSettled) → run the managed-home sync-back diff --git a/packages/adapter-utils/src/execution-phase.test.ts b/packages/adapter-utils/src/execution-phase.test.ts new file mode 100644 index 0000000000..a90c2c81f1 --- /dev/null +++ b/packages/adapter-utils/src/execution-phase.test.ts @@ -0,0 +1,26 @@ +import { expect, it, vi } from "vitest"; +import { withAdapterExecutionPhase } from "./execution-phase.js"; + +it("enters before a stalled await and releases only when that operation settles", async () => { + const release = vi.fn(); + const enter = vi.fn(() => release); + let resolve!: (value: number) => void; + const operation = new Promise((done) => { resolve = done; }); + const result = withAdapterExecutionPhase({ onExecutionPhase: enter }, "instruction_collection", () => operation); + expect(enter).toHaveBeenCalledWith("instruction_collection"); + await Promise.resolve(); + expect(release).not.toHaveBeenCalled(); + resolve(7); + await expect(result).resolves.toBe(7); + expect(release).toHaveBeenCalledOnce(); +}); + +it.each(["enter", "release"])("preserves values and errors when the %s diagnostic callback throws", async (where) => { + const context = { onExecutionPhase: () => { + if (where === "enter") throw new Error("diagnostic failure"); + return () => { throw new Error("diagnostic failure"); }; + } }; + await expect(withAdapterExecutionPhase(context, "workspace_restore", () => 3)).resolves.toBe(3); + const failure = new Error("operation failure"); + await expect(withAdapterExecutionPhase(context, "workspace_restore", () => { throw failure; })).rejects.toBe(failure); +}); diff --git a/packages/adapter-utils/src/execution-phase.ts b/packages/adapter-utils/src/execution-phase.ts new file mode 100644 index 0000000000..0ebc844763 --- /dev/null +++ b/packages/adapter-utils/src/execution-phase.ts @@ -0,0 +1,28 @@ +/** Content-free, in-memory diagnostic scopes. These never establish stop proof. */ +export const ADAPTER_EXECUTION_PHASES = [ + "host_execution", "adapter_execution", "end_session", "cancel_turn", "close_session", + "settle_reuse", "stop_transport", "sync_back", "instruction_collection", + "workspace_restore", "phase_reporting", "instruction_cleanup", "lease_release", +] as const; + +export type AdapterExecutionPhase = (typeof ADAPTER_EXECUTION_PHASES)[number]; +export type AdapterExecutionPhaseSink = (phase: AdapterExecutionPhase) => (() => void) | void; + +export function isAdapterExecutionPhase(value: unknown): value is AdapterExecutionPhase { + return ADAPTER_EXECUTION_PHASES.some((phase) => phase === value); +} + +/** A broken diagnostic sink must not change execution or teardown outcomes. */ +export async function withAdapterExecutionPhase( + context: { onExecutionPhase?: AdapterExecutionPhaseSink }, + phase: AdapterExecutionPhase, + operation: () => T | Promise, +): Promise { + let release: (() => void) | void = undefined; + try { release = context.onExecutionPhase?.(phase); } catch { /* diagnostics only */ } + try { + return await operation(); + } finally { + try { release?.(); } catch { /* diagnostics only */ } + } +} diff --git a/packages/adapter-utils/src/types.ts b/packages/adapter-utils/src/types.ts index c40da05f91..c11343cae4 100644 --- a/packages/adapter-utils/src/types.ts +++ b/packages/adapter-utils/src/types.ts @@ -5,6 +5,7 @@ import type { SshRemoteExecutionSpec } from "./ssh.js"; import type { AdapterExecutionTarget } from "./execution-target.js"; import type { RuntimeStatusSink } from "./runtime-progress.js"; +import type { AdapterExecutionPhaseSink } from "./execution-phase.js"; import type { ExecutionContinuationEnvelope, NativeFinalizationResult } from "@paperclipai/shared"; export interface AdapterAgent { @@ -196,6 +197,8 @@ export interface AdapterRuntimeEvent { } export interface AdapterExecutionContext { + /** Synchronous, content-free diagnostic scope; never stop or collection authority. */ + onExecutionPhase?: AdapterExecutionPhaseSink; /** Run-scoped operator cancellation; adapters must settle before returning. */ signal?: AbortSignal; /** Opt in to signal-based cancellation before starting provider work. */ diff --git a/server/src/__tests__/run-failure-sentry-real-sdk.test.ts b/server/src/__tests__/run-failure-sentry-real-sdk.test.ts index 583b827f3b..2d0ae8e814 100644 --- a/server/src/__tests__/run-failure-sentry-real-sdk.test.ts +++ b/server/src/__tests__/run-failure-sentry-real-sdk.test.ts @@ -35,7 +35,11 @@ afterEach(async () => { }); describe.skipIf(!sentryPackage)("run failure context with the real Sentry SDK", () => { - it("keeps unconfirmed Stop context off unrelated events and drops arbitrary error fields", async () => { + it.each([ + { phase: "workspace_restore", phaseElapsedMs: 60_123, expectedPhase: "workspace_restore", expectedElapsedMs: 60_123 }, + { phase: "private-provider-phase", phaseElapsedMs: 100, expectedPhase: "unknown", expectedElapsedMs: null }, + { phase: "workspace_restore", phaseElapsedMs: Infinity, expectedPhase: "workspace_restore", expectedElapsedMs: null }, + ])("keeps bounded unconfirmed Stop context off unrelated events ($expectedPhase)", async ({ phase, phaseElapsedMs, expectedPhase, expectedElapsedMs }) => { const Sentry = sentryPackage!; const events: Array> = []; vi.stubEnv("SENTRY_DSN_BACKEND", "https://public@example.invalid/1"); @@ -55,6 +59,7 @@ describe.skipIf(!sentryPackage)("run failure context with the real Sentry SDK", const timeout = Object.assign(new AdapterStopTimeoutError(60_000, { runId: "11111111-1111-4111-8111-111111111111", adapterType: "cursor", runtimeMode: "legacy", abortRequested: true, + phase, phaseElapsedMs, }), { cause: new Error("private-provider-cause"), providerResponse: { headers: "private-provider-headers", body: "private-provider-body" }, @@ -72,6 +77,7 @@ describe.skipIf(!sentryPackage)("run failure context with the real Sentry SDK", contexts: { adapter_stop: { runId: "11111111-1111-4111-8111-111111111111", adapterType: "cursor", runtimeMode: "legacy", abortRequested: true, timeoutMs: 60_000, + phase: expectedPhase, phaseElapsedMs: expectedElapsedMs, } }, }); expect(JSON.stringify(events)).not.toContain("private-provider-"); diff --git a/server/src/__tests__/sentry.test.ts b/server/src/__tests__/sentry.test.ts index 08659e40a1..ba044e8b1b 100644 --- a/server/src/__tests__/sentry.test.ts +++ b/server/src/__tests__/sentry.test.ts @@ -136,6 +136,7 @@ describe("captureException", () => { const error = new AdapterStopTimeoutError(60_000, { runId: "11111111-1111-4111-8111-111111111111", adapterType: "cursor", runtimeMode: "legacy", abortRequested: true, + phase: "instruction_collection", phaseElapsedMs: 60_321, }); Object.assign(error, { providerResponse: "private fixture payload" }); captureException(error); @@ -144,7 +145,8 @@ describe("captureException", () => { expect(sdk.captureException.mock.calls[0]).toEqual([ expect.objectContaining({ message: error.message, stack: error.stack }), { tags: { error_code: "adapter_stop_unconfirmed" }, fingerprint: ["{{ default }}"], contexts: { - adapter_stop: { runId: "11111111-1111-4111-8111-111111111111", adapterType: "cursor", runtimeMode: "legacy", abortRequested: true, timeoutMs: 60_000 }, + adapter_stop: { runId: "11111111-1111-4111-8111-111111111111", adapterType: "cursor", runtimeMode: "legacy", abortRequested: true, timeoutMs: 60_000, + phase: "instruction_collection", phaseElapsedMs: 60_321 }, } }, ]); expect(JSON.stringify(sdk.captureException.mock.calls[0])).not.toContain("private fixture payload"); @@ -159,10 +161,12 @@ describe("captureException", () => { await sentryReady; captureException(new AdapterStopTimeoutError(NaN, { runId: "private fixture payload", adapterType: "private fixture payload", runtimeMode: "private fixture payload", + phase: "private fixture payload", phaseElapsedMs: Infinity, })); expect(JSON.stringify(sdk.captureException.mock.calls)).not.toContain("private fixture payload"); expect(sdk.captureException.mock.calls[0]).toEqual([expect.any(Error), expect.objectContaining({ contexts: { - adapter_stop: { runId: null, adapterType: "unknown", runtimeMode: "unknown", abortRequested: null, timeoutMs: null }, + adapter_stop: { runId: null, adapterType: "unknown", runtimeMode: "unknown", abortRequested: null, timeoutMs: null, + phase: "unknown", phaseElapsedMs: null }, } })]); }); diff --git a/server/src/services/adapter-execution-control.test.ts b/server/src/services/adapter-execution-control.test.ts index 93c7143d39..7167e51f24 100644 --- a/server/src/services/adapter-execution-control.test.ts +++ b/server/src/services/adapter-execution-control.test.ts @@ -10,6 +10,16 @@ import { afterEach(() => vi.useRealTimers()); +async function stopTimeout(pending: Promise): Promise { + try { + await pending; + } catch (error) { + expect(error).toBeInstanceOf(AdapterStopTimeoutError); + return error as AdapterStopTimeoutError; + } + throw new Error("Stop unexpectedly settled"); +} + it("holds readiness for every exact-run no-owner Stop and releases owners idempotently", async () => { const runId = "registration-after-two-stops"; const first = captureAdapterStopOwnership(runId); @@ -139,6 +149,7 @@ it("records the exact unconfirmed run without finishing or removing its control" expect(error).toBeInstanceOf(AdapterStopTimeoutError); expect((error as AdapterStopTimeoutError).diagnostics).toEqual({ runId, adapterType: "claude_local", runtimeMode: "legacy", abortRequested: true, timeoutMs: 1000, + phase: "unknown", phaseElapsedMs: null, }); expect(adapterExecutionControls.get(runId)).toBe(control); expect(finished).toBe(false); @@ -151,3 +162,45 @@ it("records the exact unconfirmed run without finishing or removing its control" adapterExecutionControls.delete(runId); } }); + +it("samples the pending phase at each Stop timeout without settling the executor", async () => { + vi.useFakeTimers(); + const runId = "11111111-1111-4111-8111-111111111111"; + const control = createAdapterExecutionControl(); + await registerAdapterExecutionControl(runId, control); + const read = vi.spyOn(control.phases, "snapshot"); + const pending = stopTimeout(waitForAdapterStop(control.settled, 1000, { runId }, control)); + expect(read).not.toHaveBeenCalled(); + control.phases.enter("instruction_collection"); + await vi.advanceTimersByTimeAsync(1000); + const first = await pending; + expect(first.diagnostics).toMatchObject({ phase: "instruction_collection", phaseElapsedMs: expect.any(Number) }); + const secondWait = stopTimeout(waitForAdapterStop(control.settled, 1000, { runId }, control)); + control.phases.enter("lease_release"); + await vi.advanceTimersByTimeAsync(1000); + expect((await secondWait).diagnostics.phase).toBe("lease_release"); + expect(first.diagnostics.phase).toBe("instruction_collection"); + expect(Object.isFrozen(first.diagnostics)).toBe(true); + expect(adapterExecutionControls.get(runId)).toBe(control); + control.finish(); + adapterExecutionControls.delete(runId); +}); + +it.each(["replaced", "wrong_run", "wrong_promise", "throwing_reader"])("omits pending phase for a %s control", async (scenario) => { + vi.useFakeTimers(); + const runId = "11111111-1111-4111-8111-111111111111"; + const owner = createAdapterExecutionControl(); + const other = createAdapterExecutionControl(); + await registerAdapterExecutionControl(runId, owner); + owner.phases.enter("instruction_collection"); + const pending = stopTimeout(waitForAdapterStop(scenario === "wrong_promise" ? other.settled : owner.settled, 1000, + { runId: scenario === "wrong_run" ? "22222222-2222-4222-8222-222222222222" : runId }, owner)); + if (scenario === "replaced") await registerAdapterExecutionControl(runId, other); + if (scenario === "throwing_reader") vi.spyOn(owner.phases, "snapshot").mockImplementation(() => { throw new Error("private fixture payload"); }); + await vi.advanceTimersByTimeAsync(1000); + expect((await pending).diagnostics).toMatchObject({ phase: "unknown", phaseElapsedMs: null }); + expect(JSON.stringify((await pending).diagnostics)).not.toContain("private fixture payload"); + owner.finish(); + other.finish(); + adapterExecutionControls.delete(runId); +}); diff --git a/server/src/services/adapter-execution-control.ts b/server/src/services/adapter-execution-control.ts index a2964c7e84..06f60cab9f 100644 --- a/server/src/services/adapter-execution-control.ts +++ b/server/src/services/adapter-execution-control.ts @@ -1,13 +1,15 @@ import { AdapterStopTimeoutError, type AdapterStopContext } from "./adapter-stop-timeout.js"; +import { createAdapterExecutionPhaseTracker } from "./adapter-execution-phase.js"; /** Live adapter ownership shared by routes and scheduler service instances. */ export function createAdapterExecutionControl() { const controller = new AbortController(); + const phases = createAdapterExecutionPhaseTracker(); let finish!: () => void; const settled = new Promise((resolve) => { finish = resolve; }); - return { controller, settled, finish }; + return { controller, settled, phases, finish: () => { phases.finish(); finish(); } }; } export const adapterExecutionControls = new Map< @@ -70,7 +72,17 @@ export async function waitForAdapterStop( settled: Promise, timeoutMs = 60_000, diagnostics?: AdapterStopContext, + owner?: ReturnType, ) { + const currentPhase = () => { + try { + return owner && owner.settled === settled && diagnostics?.runId + && adapterExecutionControls.get(diagnostics.runId) === owner + ? owner.phases.snapshot() : null; + } catch { + return null; + } + }; let timer: ReturnType | undefined; try { await Promise.race([ @@ -79,7 +91,12 @@ export async function waitForAdapterStop( timer = setTimeout( () => reject( - new AdapterStopTimeoutError(timeoutMs, diagnostics), + new AdapterStopTimeoutError(timeoutMs, { + ...diagnostics, + phase: undefined, + phaseElapsedMs: undefined, + ...currentPhase(), + }), ), timeoutMs, ); diff --git a/server/src/services/adapter-execution-phase.test.ts b/server/src/services/adapter-execution-phase.test.ts new file mode 100644 index 0000000000..43446f412d --- /dev/null +++ b/server/src/services/adapter-execution-phase.test.ts @@ -0,0 +1,60 @@ +import { expect, it } from "vitest"; +import { createAdapterExecutionPhaseTracker, MAX_EXECUTION_PHASE_ELAPSED_MS } from "./adapter-execution-phase.js"; + +it("returns the latest pending scope and resumes the outer scope's original age", () => { + let now = 10; + const phases = createAdapterExecutionPhaseTracker(() => now); + const outer = phases.enter("sync_back")!; + now = 30; + const inner = phases.enter("instruction_collection")!; + now = 80; + expect(phases.snapshot()).toEqual({ phase: "instruction_collection", phaseElapsedMs: 50 }); + inner(); + expect(phases.snapshot()).toEqual({ phase: "sync_back", phaseElapsedMs: 70 }); + outer(); + expect(phases.snapshot()).toBeNull(); +}); + +it("late and repeated releases never clear a newer scope or revive a completed scope", () => { + const phases = createAdapterExecutionPhaseTracker(() => 20); + const first = phases.enter("end_session")!; + const second = phases.enter("close_session")!; + first(); + first(); + expect(phases.snapshot()?.phase).toBe("close_session"); + second(); + expect(phases.snapshot()).toBeNull(); +}); + +it("isolates trackers and refuses late updates after the executor finishes", () => { + const first = createAdapterExecutionPhaseTracker(); + const second = createAdapterExecutionPhaseTracker(); + first.enter("workspace_restore"); + second.enter("instruction_cleanup"); + first.finish(); + first.enter("lease_release"); + expect(first.snapshot()).toBeNull(); + expect(second.snapshot()?.phase).toBe("instruction_cleanup"); +}); + +it("bounds retained scopes and fails closed on overflow or arbitrary labels", () => { + const phases = createAdapterExecutionPhaseTracker(); + phases.enter("private fixture payload" as "sync_back"); + expect(phases.snapshot()).toBeNull(); + for (let i = 0; i < 100; i++) phases.enter("sync_back"); + expect(phases.snapshot()).toBeNull(); + phases.finish(); + expect(phases.snapshot()).toBeNull(); +}); + +it("uses a finite nonnegative capped elapsed time", () => { + let now = 10; + const phases = createAdapterExecutionPhaseTracker(() => now); + phases.enter("close_session"); + now = Infinity; + expect(phases.snapshot()).toBeNull(); + now = 9; + expect(phases.snapshot()).toBeNull(); + now = MAX_EXECUTION_PHASE_ELAPSED_MS * 2; + expect(phases.snapshot()).toEqual({ phase: "close_session", phaseElapsedMs: MAX_EXECUTION_PHASE_ELAPSED_MS }); +}); diff --git a/server/src/services/adapter-execution-phase.ts b/server/src/services/adapter-execution-phase.ts new file mode 100644 index 0000000000..e70abb6eaa --- /dev/null +++ b/server/src/services/adapter-execution-phase.ts @@ -0,0 +1,44 @@ +import { performance } from "node:perf_hooks"; +import { + isAdapterExecutionPhase, + type AdapterExecutionPhase, +} from "@paperclipai/adapter-utils/execution-phase"; + +export const MAX_EXECUTION_PHASE_ELAPSED_MS = 86_400_000; +const MAX_ACTIVE_PHASES = 16; + +/** Owned by one executor, never looked up by agent, task, or provider identity. */ +export function createAdapterExecutionPhaseTracker(now = () => performance.now()) { + const scopes = new Map(); + let finished = false; + let overflow = false; + return { + enter(phase: AdapterExecutionPhase) { + if (finished || !isAdapterExecutionPhase(phase)) return; + if (scopes.size >= MAX_ACTIVE_PHASES) { + // Refuse further attribution rather than retain an unbounded history or + // describe an older scope as the current one. finish clears the tracker. + overflow = true; + return; + } + const token = Symbol(); + let startedAt: number; + try { startedAt = now(); } catch { return; } + scopes.set(token, { phase, startedAt }); + return () => { scopes.delete(token); }; + }, + snapshot() { + if (finished || overflow) return null; + const current = [...scopes.values()].at(-1); + if (!current) return null; + let elapsed: number; + try { elapsed = now() - current.startedAt; } catch { return null; } + if (!Number.isFinite(elapsed) || elapsed < 0) return null; + return { phase: current.phase, phaseElapsedMs: Math.min(Math.floor(elapsed), MAX_EXECUTION_PHASE_ELAPSED_MS) }; + }, + finish() { + finished = true; + scopes.clear(); + }, + }; +} diff --git a/server/src/services/adapter-stop-timeout.ts b/server/src/services/adapter-stop-timeout.ts index d90ed90bfc..2225a813fe 100644 --- a/server/src/services/adapter-stop-timeout.ts +++ b/server/src/services/adapter-stop-timeout.ts @@ -1,10 +1,14 @@ import { AGENT_ADAPTER_TYPES } from "@paperclipai/shared"; +import { isAdapterExecutionPhase } from "@paperclipai/adapter-utils/execution-phase"; +import { MAX_EXECUTION_PHASE_ELAPSED_MS } from "./adapter-execution-phase.js"; export interface AdapterStopContext { runId?: string; adapterType?: string; runtimeMode?: string; abortRequested?: boolean; + phase?: string; + phaseElapsedMs?: number; } /** Diagnostic identity only. This error never acknowledges termination. */ @@ -22,6 +26,10 @@ export class AdapterStopTimeoutError extends Error { adapterType: AGENT_ADAPTER_TYPES.some((type) => type === context?.adapterType) ? context!.adapterType! : "unknown", runtimeMode: context?.runtimeMode === "native" || context?.runtimeMode === "legacy" ? context.runtimeMode : "unknown", abortRequested: typeof context?.abortRequested === "boolean" ? context.abortRequested : null, + phase: isAdapterExecutionPhase(context?.phase) ? context.phase : "unknown", + phaseElapsedMs: isAdapterExecutionPhase(context?.phase) && Number.isSafeInteger(context?.phaseElapsedMs) + && context!.phaseElapsedMs! >= 0 && context!.phaseElapsedMs! <= MAX_EXECUTION_PHASE_ELAPSED_MS + ? context!.phaseElapsedMs! : null, }); } } diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index a1c57d00f2..1c1e7c64c2 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -24,6 +24,7 @@ import { githubBotConnectionIdsForRun } from "./chat-github-tools.js"; import { isBrowserUseConnection } from "./browser-use-client.js"; import { readQueuedInteractionResponse } from "./queued-interaction-response.js"; import { AGENT_CHAT_DIRECTIVE, isConversation, isConversationExecutionWake, isWaitingConversation, prepareConversationTurn, settleConversationTurn } from "./agent-conversations.js"; +import { withAdapterExecutionPhase } from "@paperclipai/adapter-utils/execution-phase"; import { getConversationConfirmationContext, type ConversationConfirmationContext } from "./conversation-confirmation-context.js"; import { PROCESS_IDENTITY_RECORDED, recordNativeLocalProcessStop } from "./native-local-process-stop.js"; import { hasAcknowledgedNativeReassignmentStopIntent, hasAcknowledgedNativeStopIntent, isAcknowledgedNativeStop, acknowledgedNativeStopExecutionHasStopped } from "./acknowledged-native-stop.js"; @@ -20422,6 +20423,10 @@ export function heartbeatService( run.controllerBootId !== legacyControllerBootId) return; activeRunExecutions.add(run.id); const executionControl = createAdapterExecutionControl(); + // This coarse scope also covers host finalization after the adapter returns. + // Nested scopes refine the label without making any termination claim. + executionControl.phases.enter("host_execution"); + const executionPhaseContext = { onExecutionPhase: executionControl.phases.enter }; const controllerLease = watchLegacyControllerLease(db, run, executionControl.controller); let runScratch: HeartbeatRunScratch | null = null; let githubLauncherLocation: @@ -24908,7 +24913,7 @@ export function heartbeatService( await dispatchResolvedInteractionContinuationWithAtomicGate( (markDispatchStarted) => { legacyAdapterEntered = true; - return adapter.execute({ + return withAdapterExecutionPhase(executionPhaseContext, "adapter_execution", () => adapter.execute({ getFreshSessionHandoff, runId: run.id, agent, @@ -24932,6 +24937,7 @@ export function heartbeatService( onLog, onMeta: onAdapterMeta, onEvent: onAdapterEvent, + onExecutionPhase: executionControl.phases.enter, startupTraceContext: getStartupTraceContext(), onRuntimeProgress: async (progress) => { await recordCurrentHeartbeatRunRuntimeProgress( @@ -24983,7 +24989,7 @@ export function heartbeatService( }); }, authToken: authToken ?? undefined, - }); + })); }, ); if (!guardedDispatch.dispatched) return; @@ -25226,7 +25232,7 @@ export function heartbeatService( ); } await nativeInstructionReservation?.release(); - await releaseInstructionCopy(); + await withAdapterExecutionPhase(executionPhaseContext, "instruction_cleanup", releaseInstructionCopy); } // Reconcile the referenced-project set against the real remote staging outcome. A referenced // project can pass authorization and clone locally at run prep, then fail to stage into the @@ -26632,8 +26638,8 @@ export function heartbeatService( message: "Instruction edits could not be recovered before environment release. No instruction save is claimed.", payload: { state: "unavailable", code: uncapturedInstructions.errorCode } }); } - await releaseInstructionCopy(); - await releaseEnvironmentLeasesForRun({ + await withAdapterExecutionPhase(executionPhaseContext, "instruction_cleanup", releaseInstructionCopy); + await withAdapterExecutionPhase(executionPhaseContext, "lease_release", () => releaseEnvironmentLeasesForRun({ runId: run.id, companyId: run.companyId, agentId: run.agentId, @@ -26641,7 +26647,7 @@ export function heartbeatService( failureReason: latestRun?.error ?? undefined, providerResourceDisposition: providerResourceDispositionForRun, nativeLifecycleTelemetry: nativeLifecycleTelemetryForRun, - }); + })); await releaseRuntimeServicesForRun(run.id).catch(() => undefined); } if ( @@ -29613,7 +29619,7 @@ export function heartbeatService( adapterType: agent?.adapterType, runtimeMode: run.runtimeMode, abortRequested: control.controller.signal.aborted, - }); + }, control); const stopped = await getRun(run.id); if (stopped && isHeartbeatRunTerminalStatus(stopped.status)) { if (