From e2b6070339351f017d2a971a464a7ec658b21dbe Mon Sep 17 00:00:00 2001 From: Daniel Saldarriaga Date: Mon, 5 Oct 2026 10:27:14 +0200 Subject: [PATCH 1/2] fix: persist explicit goal cancellation --- README.md | 9 ++- dist/server.js | 120 +++++++++++++++++++++++++++++++-------- src/server.ts | 125 ++++++++++++++++++++++++++++++++--------- test/server-v2.test.ts | 75 ++++++++++++++++++++++++- test/server.test.ts | 58 +++++++++++++++++++ 5 files changed, 335 insertions(+), 52 deletions(-) diff --git a/README.md b/README.md index 648110e..906ca9b 100644 --- a/README.md +++ b/README.md @@ -315,6 +315,13 @@ OpenCode plugin modules are target-specific. This package exports separate modul } ``` -Codex goal mode has deeper runtime integration for thread lifecycle control. This plugin implements the same workflow using OpenCode plugin hooks. Token usage is read from OpenCode step-finish usage when available and falls back to message token metadata or text estimation when exact usage is unavailable. Continuation is driven by V2 `session.execution.succeeded` events and legacy `session.idle` / `session.status` idle notifications, never by intermediate model-step completion. V2 execution starts arm busy tracking; native `session.retry.scheduled` events cancel plugin recovery while OpenCode retries. Terminal execution transport failures use bounded recovery, while interruptions (user, shutdown, or superseded) cancel local timers without starting another turn or charging a prompt failure. Interrupted or non-transport-failed executions remain suppressed until the host starts a new execution; this does not change the persisted goal status. Use `/pause_goal` for a durable pause. Each V2 plugin instance handles goal events only for its own location while still observing cross-location child task lifecycles. The optional `max_turn_time` watchdog can retry one goal continuation prompt when a model turn remains busy, without consuming the goal's auto-turn or no-progress budgets; recognized transport failures do count toward the prompt-failure ceiling. By default, continuation is deferred while OpenCode Task child sessions are active or their terminal result still needs an orchestrator turn, bounded by the `max_task_block_seconds` ceiling so an unobservable child cannot stall a goal indefinitely. During compaction on OpenCode 1, the plugin disables OpenCode's generic synthetic auto-continue while an active goal exists so the goal-specific continuation prompt remains authoritative; on OpenCode 2, compaction runs inside the session's execution flow, so the plugin instead injects the goal snapshot through the `session.compaction` hook where the host provides it. +Codex goal mode has deeper runtime integration for thread lifecycle control. This plugin implements the same workflow using OpenCode plugin hooks. Token usage is read from OpenCode step-finish usage when available and falls back to message token metadata or text estimation when exact usage is unavailable. Continuation is driven by V2 `session.execution.succeeded` events and legacy `session.idle` / `session.status` idle notifications, never by intermediate model-step completion. V2 execution starts arm busy tracking; native `session.retry.scheduled` events cancel plugin recovery while OpenCode retries. Terminal execution transport failures use bounded recovery. A V2 `session.execution.interrupted` event with reason `user`, or a V1 session/assistant `MessageAbortedError`, persists a terminal `cancelled` goal and invalidates timers and outstanding continuation preparation without charging a prompt failure. A later idle, unrelated prompt, or plugin reload cannot resume that goal; start an explicitly requested new goal with `/goal ` or `/goal replace `. Shutdown, superseded, and other interruptions only suppress local continuation until the host starts another execution; they do not cancel the persisted goal. Use `/pause_goal` for a durable, resumable pause. Each V2 plugin instance handles goal events only for its own location while still observing cross-location child task lifecycles. The optional `max_turn_time` watchdog can retry one goal continuation prompt when a model turn remains busy, without consuming the goal's auto-turn or no-progress budgets; recognized transport failures do count toward the prompt-failure ceiling. By default, continuation is deferred while OpenCode Task child sessions are active or their terminal result still needs an orchestrator turn, bounded by the `max_task_block_seconds` ceiling so an unobservable child cannot stall a goal indefinitely. During compaction on OpenCode 1, the plugin disables OpenCode's generic synthetic auto-continue while an active goal exists so the goal-specific continuation prompt remains authoritative; on OpenCode 2, compaction runs inside the session's execution flow, so the plugin instead injects the goal snapshot through the `session.compaction` hook where the host provides it. The goal sidebar shows the current status, elapsed time, token usage, auto-continue count, latest checkpoint, latest status message, stop reason, and objective when a goal is active, paused, or safety-limited. It checks the shared goal state file every second so usage and checkpoints stay current during a long run. Closed goals remain visible briefly through the latest tool state as achieved or unmet. + + +### Zed and ACP lifecycle boundaries + +Cancelling an active goal from an ACP client is durable when OpenCode emits the user-cancellation events described above. The plugin prevents subsequent goal continuations, including callbacks still preparing a prompt when cancellation arrives. OpenCode owns cancellation of in-flight model requests, tools, subprocesses, and already submitted prompts; the plugin cannot guarantee process termination if the host does not abort them or does not publish a cancellation signal. It does not interpret ordinary idle events, provider-error text, or completed tool calls as user cancellation. + +Goal lifetime and ACP prompt-turn lifetime are separate. The plugin continues an active goal across successful executions, but does not control the ACP adapter's response to `session/prompt`. A host that returns `end_turn` after the initial execution can therefore show an idle client while subsequent goal work runs. In an isolated OpenCode 2.0.21 ACP test, registered `/goal` commands returned `end_turn` while their submitted execution was still running; a normal `session/prompt` remained open and returned `cancelled` after Cancel. Keeping that ACP request open across goal continuations, and exposing structured goal plans as ACP `plan` updates, requires host integration. Request or phase completion alone must not mark a goal complete. Long-running goals remain supported within the configured budgets; this cancellation handling adds no timeout. diff --git a/dist/server.js b/dist/server.js index b5f3ac8..0a38667 100644 --- a/dist/server.js +++ b/dist/server.js @@ -1937,6 +1937,27 @@ var NON_TRANSPORT_TERMINAL_PATTERN = /\b(?:abort(?:ed)?|interrupt(?:ed|ion)?)\b/ var NON_PROGRESS_TOOLS = new Set(["get_goal", "get_goal_history", "list_all_goals"]); var TASK_TERMINAL_STATES = new Set(["completed", "error", "cancelled"]); var activeContinuations = new Set; + +class ContinuationEpochs { + values = new Map; + current(sessionID) { + return this.values.get(sessionID) ?? 0; + } + invalidate(sessionID) { + this.values.set(sessionID, this.current(sessionID) + 1); + } +} +function isUserAbortEvent(event) { + const properties = event.properties; + if (event.type === "session.error") { + return isRecord(properties?.error) && properties.error.name === "MessageAbortedError"; + } + const message = properties?.info; + return event.type === "message.updated" && isRecord(message) && message.role === "assistant" && isRecord(message.error) && message.error.name === "MessageAbortedError"; +} +function continuationStillReserved(goal, current) { + return current?.id === goal.id && current.status === goal.status && (goal.status !== "active" || current.pendingAttempt?.id === goal.pendingAttempt?.id); +} function restrictedAgentSet(options) { if (options?.allow_goal_execution_from_plan === true) return new Set; @@ -2963,6 +2984,7 @@ var server = async ({ client }, options) => { const toolAttempts = new Map; const explicitResumeRequests = new Set; const restartAfterContinuation = new Set; + const continuationEpochs = new ContinuationEpochs; const watchdogRescuedSessions = new Set; const planAgents = restrictedAgentSet(options); const isPlanAgent = (agent) => typeof agent === "string" && planAgents.has(agent.trim().toLowerCase()); @@ -2974,6 +2996,7 @@ var server = async ({ client }, options) => { maxObjectiveChars: objectiveChars, consumeAutoTurnReset: (sessionID) => explicitResumeRequests.delete(sessionID), stopAutonomy: (sessionID, mode = "stop") => { + continuationEpochs.invalidate(sessionID); cancelScheduledContinuation(sessionID); if (mode === "stop") clearTurnWatchdog(sessionID); @@ -3027,6 +3050,8 @@ var server = async ({ client }, options) => { turnWatchdogs.set(sessionID, watchdog); } async function runTurnWatchdog(sessionID, watchdog) { + const epoch = continuationEpochs.current(sessionID); + const isCurrent = () => !disposed && epoch === continuationEpochs.current(sessionID); let claimedContinuation = false; let claimedGoalID; try { @@ -3062,17 +3087,23 @@ var server = async ({ client }, options) => { claimedContinuation = true; claimedGoalID = current.id; watchdogRescuedSessions.add(sessionID); + if (!isCurrent()) + return; await sendContinuation(client, sessionID, continuationPrompt(current, locale), current.lastPromptAgent ?? latestTurnAgent ?? null); - await recordContinuationResult(sessionID, "success", maxPromptFailures, { + if (!isCurrent()) + return; + const delivered = await recordContinuationResult(sessionID, "success", maxPromptFailures, { armNoProgress: false, started: true, expectedGoalID: claimedGoalID }); - locallyDeliveredPendingSessions.add(sessionID); - clearTurnWatchdog(sessionID); + if (isCurrent() && delivered?.pendingAttempt?.delivered) { + locallyDeliveredPendingSessions.add(sessionID); + clearTurnWatchdog(sessionID); + } } catch (error) { try { - if (claimedContinuation && isTransportError(error)) { + if (claimedContinuation && isCurrent() && isTransportError(error)) { await recordContinuationResult(sessionID, "failure", maxPromptFailures, { expectedGoalID: claimedGoalID }); } await client.app?.log?.({ @@ -3142,16 +3173,24 @@ var server = async ({ client }, options) => { return; if (activeContinuations.has(sessionID)) return; + const epoch = continuationEpochs.current(sessionID); + const isCurrent = () => !disposed && epoch === continuationEpochs.current(sessionID); activeContinuations.add(sessionID); let attemptReservedAt = Date.now(); let attemptGoalID; let attemptID; try { const latestAssistant = await fetchLatestAssistant(client, sessionID); + if (!isCurrent()) + return; taskTracker.observeAssistantMessage(sessionID, latestAssistant); const taskStatus = await taskBlockStatus(sessionID); + if (!isCurrent()) + return; if (taskStatus && taskStatus.blocked) { const deferralGoal = await getGoalInternal(sessionID); + if (!isCurrent()) + return; if (!taskDeferralGoalContinuable(deferralGoal)) { taskDeferredSessions.delete(sessionID); cancelScheduledContinuation(sessionID); @@ -3161,7 +3200,7 @@ var server = async ({ client }, options) => { scheduleSettledContinuation(sessionID, taskStatus.retryAt != null ? taskStatus.retryAt - Date.now() : TASK_BLOCK_RETRY_MS, scheduled != null || taskStatus.retryAt != null); return; } - if (busySessions.has(sessionID)) + if (!isCurrent() || busySessions.has(sessionID)) return; const observed = await recordAssistantMessage(sessionID, latestAssistant, options ?? {}, true); await reconcileLocalMarkerAfterProgress(locallyDeliveredPendingSessions, sessionID, observed.goal); @@ -3171,7 +3210,7 @@ var server = async ({ client }, options) => { if (scheduled && scheduledContinuations.get(sessionID) !== scheduled) return; const current = await getGoalInternal(sessionID); - if (!current) + if (!isCurrent() || !current) return; const latestTurnAgent = agentFromMessage(latestAssistant); if (isPlanAgent(current.lastPromptAgent) || isPlanAgent(latestTurnAgent)) { @@ -3209,7 +3248,7 @@ var server = async ({ client }, options) => { return; if (!autoContinue) return; - if (nativeRetrySessions.has(sessionID)) + if (!isCurrent() || nativeRetrySessions.has(sessionID)) return; const goal = await reserveContinuation(sessionID, maxAutoTurns, minInterval); if (!goal) @@ -3217,7 +3256,8 @@ var server = async ({ client }, options) => { attemptReservedAt = goal.pendingAttempt?.reservedAt ?? Date.now(); attemptGoalID = goal.id; attemptID = goal.pendingAttempt?.id; - if (nativeRetrySessions.has(sessionID)) { + const beforeDelivery = await getGoalInternal(sessionID); + if (!isCurrent() || !continuationStillReserved(goal, beforeDelivery) || busySessions.has(sessionID) || nativeRetrySessions.has(sessionID)) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); return; } @@ -3226,19 +3266,20 @@ var server = async ({ client }, options) => { return; } await sendContinuation(client, sessionID, goal.status === "active" ? continuationPrompt(goal, locale) : limitPrompt(goal, locale), goal.lastPromptAgent ?? latestTurnAgent ?? null); - if (disposed) { + if (!isCurrent()) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); return; } const delivered = await recordContinuationResult(sessionID, "success", maxPromptFailures, { expectedGoalID: attemptGoalID }); - locallyDeliveredPendingSessions.add(sessionID); + if (isCurrent() && delivered?.pendingAttempt?.delivered) + locallyDeliveredPendingSessions.add(sessionID); if (!delivered?.pendingAttempt?.delivered) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); } } catch (error) { - if (disposed) { + if (!isCurrent()) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); return; } @@ -3503,6 +3544,17 @@ var server = async ({ client }, options) => { async event({ event }) { const sessionID = sessionIDFromEvent(event); const eventType = event.type; + if (sessionID && isUserAbortEvent(event)) { + explicitResumeRequests.delete(sessionID); + goalServices.stopAutonomy?.(sessionID); + busySessions.delete(sessionID); + nativeRetrySessions.delete(sessionID); + watchdogRescuedSessions.delete(sessionID); + clearToolAttemptsForSession(toolAttempts, sessionID); + taskTracker.observeSessionStatus(sessionID, "idle"); + await cancelGoal(sessionID); + return; + } if (eventType === "session.created") { taskTracker.observeSessionCreated(event); } @@ -3571,6 +3623,7 @@ var server = async ({ client }, options) => { } } if (sessionID && eventType === "session.deleted") { + continuationEpochs.invalidate(sessionID); explicitResumeRequests.delete(sessionID); busySessions.delete(sessionID); clearTurnWatchdog(sessionID); @@ -3635,6 +3688,7 @@ async function setupV2(context) { const isPlanAgent = (agent) => typeof agent === "string" && planAgents.has(agent.trim().toLowerCase()); const activeContinuationsV2 = new Set; const restartAfterContinuation = new Set; + const continuationEpochs = new ContinuationEpochs; const stoppedExecutions = new Set; const latestStepBySession = new Map; const stepTextBuffers = new Map; @@ -3654,6 +3708,7 @@ async function setupV2(context) { } }, stopAutonomy: (sessionID, mode = "stop") => { + continuationEpochs.invalidate(sessionID); cancelScheduledContinuation(sessionID); if (mode === "stop") clearTurnWatchdog(sessionID); @@ -3712,6 +3767,8 @@ async function setupV2(context) { turnWatchdogs.set(sessionID, watchdog); } async function runTurnWatchdog(sessionID, watchdog) { + const epoch = continuationEpochs.current(sessionID); + const isCurrent = () => !disposed && epoch === continuationEpochs.current(sessionID); let claimedContinuation = false; let claimedGoalID; try { @@ -3743,17 +3800,23 @@ async function setupV2(context) { claimedContinuation = true; claimedGoalID = current.id; watchdogRescuedSessions.add(sessionID); + if (!isCurrent()) + return; await sendContinuation2(sessionID, continuationPrompt(current, locale), current.lastPromptAgent ?? latestStep?.agent ?? null); - await recordContinuationResult(sessionID, "success", maxPromptFailures, { + if (!isCurrent()) + return; + const delivered = await recordContinuationResult(sessionID, "success", maxPromptFailures, { armNoProgress: false, started: true, expectedGoalID: claimedGoalID }); - locallyDeliveredPendingSessions.add(sessionID); - clearTurnWatchdog(sessionID); + if (isCurrent() && delivered?.pendingAttempt?.delivered) { + locallyDeliveredPendingSessions.add(sessionID); + clearTurnWatchdog(sessionID); + } } catch (error) { try { - if (claimedContinuation && isTransportError(error)) { + if (claimedContinuation && isCurrent() && isTransportError(error)) { await recordContinuationResult(sessionID, "failure", maxPromptFailures, { expectedGoalID: claimedGoalID }); } v2ErrorLog("Turn watchdog retry failed", error); @@ -3818,8 +3881,10 @@ async function setupV2(context) { return; if (activeContinuationsV2.has(sessionID)) return; + const epoch = continuationEpochs.current(sessionID); + const isCurrent = () => !disposed && epoch === continuationEpochs.current(sessionID); await taskRecoveryComplete; - if (disposed || stoppedExecutions.has(sessionID) || busySessions.has(sessionID)) + if (!isCurrent() || stoppedExecutions.has(sessionID) || busySessions.has(sessionID)) return; activeContinuationsV2.add(sessionID); let attemptReservedAt = Date.now(); @@ -3833,6 +3898,8 @@ async function setupV2(context) { const taskStatus = taskBlockStatus(sessionID); if (taskStatus && taskStatus.blocked) { const deferralGoal = await getGoalInternal(sessionID); + if (!isCurrent()) + return; if (!taskDeferralGoalContinuable(deferralGoal)) { taskDeferredSessions.delete(sessionID); cancelScheduledContinuation(sessionID); @@ -3864,7 +3931,7 @@ async function setupV2(context) { if (scheduled && scheduledContinuations.get(sessionID) !== scheduled) return; const current = await getGoalInternal(sessionID); - if (!current) + if (!isCurrent() || !current) return; const latestTurnAgent = latestStep?.agent; if (isPlanAgent(current.lastPromptAgent) || isPlanAgent(latestTurnAgent)) { @@ -3902,7 +3969,7 @@ async function setupV2(context) { return; if (!autoContinue) return; - if (nativeRetrySessions.has(sessionID)) + if (!isCurrent() || nativeRetrySessions.has(sessionID)) return; const goal = await reserveContinuation(sessionID, maxAutoTurns, minInterval); if (!goal) { @@ -3915,7 +3982,8 @@ async function setupV2(context) { attemptReservedAt = goal.pendingAttempt?.reservedAt ?? Date.now(); attemptGoalID = goal.id; attemptID = goal.pendingAttempt?.id; - if (nativeRetrySessions.has(sessionID)) { + const beforeDelivery = await getGoalInternal(sessionID); + if (!isCurrent() || !continuationStillReserved(goal, beforeDelivery) || busySessions.has(sessionID) || nativeRetrySessions.has(sessionID)) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); return; } @@ -3924,19 +3992,20 @@ async function setupV2(context) { return; } await sendContinuation2(sessionID, goal.status === "active" ? continuationPrompt(goal, locale) : limitPrompt(goal, locale), goal.lastPromptAgent ?? latestTurnAgent ?? null); - if (disposed) { + if (!isCurrent()) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); return; } const delivered = await recordContinuationResult(sessionID, "success", maxPromptFailures, { expectedGoalID: attemptGoalID }); - locallyDeliveredPendingSessions.add(sessionID); + if (isCurrent() && delivered?.pendingAttempt?.delivered) + locallyDeliveredPendingSessions.add(sessionID); if (!delivered?.pendingAttempt?.delivered) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); } } catch (error) { - if (disposed) { + if (!isCurrent()) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }); return; } @@ -4101,9 +4170,11 @@ async function setupV2(context) { nativeRetrySessions.delete(sessionID); clearTurnWatchdog(sessionID); watchdogRescuedSessions.delete(sessionID); - cancelScheduledContinuation(sessionID); - taskDeferredSessions.delete(sessionID); + goalServices.stopAutonomy?.(sessionID); + clearToolAttemptsForSession(toolAttempts, sessionID); taskTracker.observeSessionStatus(sessionID, "idle"); + if (data.reason === "user") + await cancelGoal(sessionID); return; } case "session.execution.failed": { @@ -4144,6 +4215,7 @@ async function setupV2(context) { case "session.deleted": { if (!sessionID) return; + continuationEpochs.invalidate(sessionID); explicitResumeRequests.delete(sessionID); stoppedExecutions.delete(sessionID); sessionOwnership.delete(sessionID); diff --git a/src/server.ts b/src/server.ts index e659e79..f9e7f94 100644 --- a/src/server.ts +++ b/src/server.ts @@ -134,6 +134,39 @@ type ScheduledContinuation = { purpose: "settle" | "recovery" | "retry" } +// Invalidate work already awaiting transcript/state IO when its goal is stopped +// or replaced. A fresh busy event must not make an older callback current again. +class ContinuationEpochs { + private readonly values = new Map() + + current(sessionID: string) { + return this.values.get(sessionID) ?? 0 + } + + invalidate(sessionID: string) { + this.values.set(sessionID, this.current(sessionID) + 1) + } +} + +function isUserAbortEvent(event: { type?: string; properties?: Record }) { + const properties = event.properties + if (event.type === "session.error") { + return isRecord(properties?.error) && properties.error.name === "MessageAbortedError" + } + const message = properties?.info + return ( + event.type === "message.updated" && isRecord(message) && message.role === "assistant" && + isRecord(message.error) && message.error.name === "MessageAbortedError" + ) +} + +function continuationStillReserved(goal: InternalGoalSnapshot, current: InternalGoalSnapshot | null) { + return ( + current?.id === goal.id && current.status === goal.status && + (goal.status !== "active" || current.pendingAttempt?.id === goal.pendingAttempt?.id) + ) +} + function restrictedAgentSet(options?: Options) { if (options?.allow_goal_execution_from_plan === true) return new Set() const names = Array.isArray(options?.restricted_agents) ? options.restricted_agents : DEFAULT_RESTRICTED_AGENTS @@ -1334,6 +1367,7 @@ const server: Plugin = async ({ client }, options?: Options) => { const toolAttempts = new Map() const explicitResumeRequests = new Set() const restartAfterContinuation = new Set() + const continuationEpochs = new ContinuationEpochs() // Sessions whose busy episode already received a watchdog rescue. Cleared // when the episode ends (idle/deleted), so each busy episode rescues at most // once and a rescue prompt cannot recursively re-arm the watchdog. @@ -1348,6 +1382,7 @@ const server: Plugin = async ({ client }, options?: Options) => { maxObjectiveChars: objectiveChars, consumeAutoTurnReset: (sessionID) => explicitResumeRequests.delete(sessionID), stopAutonomy: (sessionID, mode = "stop") => { + continuationEpochs.invalidate(sessionID) cancelScheduledContinuation(sessionID) if (mode === "stop") clearTurnWatchdog(sessionID) taskDeferredSessions.delete(sessionID) @@ -1404,6 +1439,8 @@ const server: Plugin = async ({ client }, options?: Options) => { } async function runTurnWatchdog(sessionID: string, watchdog: TurnWatchdog) { + const epoch = continuationEpochs.current(sessionID) + const isCurrent = () => !disposed && epoch === continuationEpochs.current(sessionID) let claimedContinuation = false let claimedGoalID: string | undefined try { @@ -1437,6 +1474,7 @@ const server: Plugin = async ({ client }, options?: Options) => { claimedContinuation = true claimedGoalID = current.id watchdogRescuedSessions.add(sessionID) + if (!isCurrent()) return await sendContinuation( client, sessionID, @@ -1448,19 +1486,22 @@ const server: Plugin = async ({ client }, options?: Options) => { // never arms the no-progress evaluation. The rescue delivers while the // session is already inside a busy episode, so the pending attempt is // marked started immediately, and this busy episode rescues only once. - await recordContinuationResult(sessionID, "success", maxPromptFailures, { + if (!isCurrent()) return + const delivered = await recordContinuationResult(sessionID, "success", maxPromptFailures, { armNoProgress: false, started: true, expectedGoalID: claimedGoalID, }) - locallyDeliveredPendingSessions.add(sessionID) - clearTurnWatchdog(sessionID) + if (isCurrent() && delivered?.pendingAttempt?.delivered) { + locallyDeliveredPendingSessions.add(sessionID) + clearTurnWatchdog(sessionID) + } } catch (error) { try { // Watchdog rescues share the same prompt-failure ceiling: recognized // transport errors accumulate toward max_prompt_failures without // consuming auto-turn budgets. - if (claimedContinuation && isTransportError(error)) { + if (claimedContinuation && isCurrent() && isTransportError(error)) { await recordContinuationResult(sessionID, "failure", maxPromptFailures, { expectedGoalID: claimedGoalID }) } await client.app?.log?.({ @@ -1529,6 +1570,8 @@ const server: Plugin = async ({ client }, options?: Options) => { if (disposed) return if (busySessions.has(sessionID)) return if (activeContinuations.has(sessionID)) return + const epoch = continuationEpochs.current(sessionID) + const isCurrent = () => !disposed && epoch === continuationEpochs.current(sessionID) activeContinuations.add(sessionID) // Anchor for bounded-retry scheduling, declared at function scope so the // catch block can use it. Initialized to "now" as a safe default. @@ -1537,14 +1580,17 @@ const server: Plugin = async ({ client }, options?: Options) => { let attemptID: string | undefined try { const latestAssistant = await fetchLatestAssistant(client, sessionID) + if (!isCurrent()) return taskTracker.observeAssistantMessage(sessionID, latestAssistant) const taskStatus = await taskBlockStatus(sessionID) + if (!isCurrent()) return if (taskStatus && taskStatus.blocked) { // Validate the goal before re-arming. The re-arm below runs at TASK_BLOCK_RETRY_MS // and writes nothing to the goal, so a goal completed, cleared, or paused while a // child still blocks would otherwise keep a 1 Hz poll alive until the ceiling - and // forever when max_task_block_seconds is 0. const deferralGoal = await getGoalInternal(sessionID) + if (!isCurrent()) return if (!taskDeferralGoalContinuable(deferralGoal)) { taskDeferredSessions.delete(sessionID) cancelScheduledContinuation(sessionID) @@ -1564,14 +1610,14 @@ const server: Plugin = async ({ client }, options?: Options) => { ) return } - if (busySessions.has(sessionID)) return + if (!isCurrent() || busySessions.has(sessionID)) return const observed = await recordAssistantMessage(sessionID, latestAssistant, options ?? {}, true) await reconcileLocalMarkerAfterProgress(locallyDeliveredPendingSessions, sessionID, observed.goal) const queued = scheduledContinuations.get(sessionID) if (observed.progressed && queued?.purpose !== "settle") cancelScheduledContinuation(sessionID) if (scheduled && scheduledContinuations.get(sessionID) !== scheduled) return const current = await getGoalInternal(sessionID) - if (!current) return + if (!isCurrent() || !current) return const latestTurnAgent = agentFromMessage(latestAssistant) if (isPlanAgent(current.lastPromptAgent) || isPlanAgent(latestTurnAgent)) { if (current.status === "active") await pauseGoalForPlanMode(sessionID) @@ -1620,7 +1666,7 @@ const server: Plugin = async ({ client }, options?: Options) => { const queuedBeforeReserve = scheduledContinuations.get(sessionID) if (queuedBeforeReserve && queuedBeforeReserve !== scheduled) return if (!autoContinue) return - if (nativeRetrySessions.has(sessionID)) return + if (!isCurrent() || nativeRetrySessions.has(sessionID)) return // Reserve (and persist) the attempt BEFORE delivery so a racing busy can // correlate to it. The attempt stays reserved until delivery or rollback. @@ -1629,7 +1675,8 @@ const server: Plugin = async ({ client }, options?: Options) => { attemptReservedAt = goal.pendingAttempt?.reservedAt ?? Date.now() attemptGoalID = goal.id attemptID = goal.pendingAttempt?.id - if (nativeRetrySessions.has(sessionID)) { + const beforeDelivery = await getGoalInternal(sessionID) + if (!isCurrent() || !continuationStillReserved(goal, beforeDelivery) || busySessions.has(sessionID) || nativeRetrySessions.has(sessionID)) { await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }) return } @@ -1643,8 +1690,8 @@ const server: Plugin = async ({ client }, options?: Options) => { goal.status === "active" ? continuationPrompt(goal, locale) : limitPrompt(goal, locale), goal.lastPromptAgent ?? latestTurnAgent ?? null, ) - if (disposed) { - // The plugin was torn down while the prompt was in flight: roll the + if (!isCurrent()) { + // The goal was stopped/replaced or the plugin disposed in flight: roll the // reserved turn back instead of committing a continuation afterward. await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }) return @@ -1654,15 +1701,15 @@ const server: Plugin = async ({ client }, options?: Options) => { const delivered = await recordContinuationResult(sessionID, "success", maxPromptFailures, { expectedGoalID: attemptGoalID, }) - locallyDeliveredPendingSessions.add(sessionID) + if (isCurrent() && delivered?.pendingAttempt?.delivered) locallyDeliveredPendingSessions.add(sessionID) if (!delivered?.pendingAttempt?.delivered) { // The attempt was not present at delivery time (e.g. disposed mid-send): // do not leave a phantom reserved turn. await rollbackContinuationAttempt(sessionID, { goalID: attemptGoalID, attemptID }) } } catch (error) { - if (disposed) { - // The plugin was torn down while the prompt was in flight and the + if (!isCurrent()) { + // The goal was stopped/replaced or the plugin disposed in flight and the // prompt then failed: the reserved attempt was never delivered, so // roll it back instead of counting a transport failure or consuming an // auto-turn. @@ -1961,6 +2008,17 @@ const server: Plugin = async ({ client }, options?: Options) => { async event({ event }) { const sessionID = sessionIDFromEvent(event as never) const eventType = (event as { type?: string }).type + if (sessionID && isUserAbortEvent(event as never)) { + explicitResumeRequests.delete(sessionID) + goalServices.stopAutonomy?.(sessionID) + busySessions.delete(sessionID) + nativeRetrySessions.delete(sessionID) + watchdogRescuedSessions.delete(sessionID) + clearToolAttemptsForSession(toolAttempts, sessionID) + taskTracker.observeSessionStatus(sessionID, "idle") + await cancelGoal(sessionID) + return + } if (eventType === "session.created") { taskTracker.observeSessionCreated(event as { properties?: Record }) } @@ -2044,6 +2102,7 @@ const server: Plugin = async ({ client }, options?: Options) => { } } if (sessionID && eventType === "session.deleted") { + continuationEpochs.invalidate(sessionID) explicitResumeRequests.delete(sessionID) busySessions.delete(sessionID) clearTurnWatchdog(sessionID) @@ -2115,6 +2174,7 @@ async function setupV2(context: PluginV2.Plugin.Context): Promise typeof agent === "string" && planAgents.has(agent.trim().toLowerCase()) const activeContinuationsV2 = new Set() const restartAfterContinuation = new Set() + const continuationEpochs = new ContinuationEpochs() // Interruptions and terminal failures are not successful idle boundaries. // Keep legacy idle notifications and queued recovery from restarting them; // only a new execution started by the host may lift this local suppression. @@ -2137,6 +2197,7 @@ async function setupV2(context: PluginV2.Plugin.Context): Promise { + continuationEpochs.invalidate(sessionID) cancelScheduledContinuation(sessionID) if (mode === "stop") clearTurnWatchdog(sessionID) taskDeferredSessions.delete(sessionID) @@ -2192,6 +2253,8 @@ async function setupV2(context: PluginV2.Plugin.Context): Promise !disposed && epoch === continuationEpochs.current(sessionID) let claimedContinuation = false let claimedGoalID: string | undefined try { @@ -2216,21 +2279,25 @@ async function setupV2(context: PluginV2.Plugin.Context): Promise !disposed && epoch === continuationEpochs.current(sessionID) // Transcript recovery must settle before any continuation decision; // otherwise the first lifecycle event after a restart defers to a task // state that has not been rebuilt yet. await taskRecoveryComplete - if (disposed || stoppedExecutions.has(sessionID) || busySessions.has(sessionID)) return + if (!isCurrent() || stoppedExecutions.has(sessionID) || busySessions.has(sessionID)) return activeContinuationsV2.add(sessionID) let attemptReservedAt = Date.now() let attemptGoalID: string | undefined @@ -2313,6 +2382,7 @@ async function setupV2(context: PluginV2.Plugin.Context): Promise { const mock = makeMockContext({ max_turn_time: 0.02, min_continue_interval_seconds: 0 }) const cleanup = await setupPlugin(mock as never) @@ -1786,6 +1786,79 @@ for (const reason of ["user", "shutdown", "superseded"]) { }) } +test("V2 user cancellation persists across reload and unrelated executions", async () => { + const mock = makeMockContext({ min_continue_interval_seconds: 0, max_turn_time: 0.02 }) + const cleanup = await setupPlugin(mock as never) + await createGoalViaV2Tool(mock, "respect user cancellation") + await mock.stream.push({ type: "session.execution.started", created: 1, data: { sessionID: "ses_v2" } }) + await mock.stream.push({ type: "session.execution.interrupted", created: 2, data: { sessionID: "ses_v2", reason: "user" } }) + await mock.stream.push({ type: "session.execution.interrupted", created: 3, data: { sessionID: "ses_v2", reason: "user" } }) + await mock.stream.push({ type: "session.idle", created: 4, data: { sessionID: "ses_v2" } }) + expect(await getGoalInternal("ses_v2")).toMatchObject({ status: "cancelled", pendingAttempt: null, continuationFailures: 0 }) + const persisted = JSON.parse(await readFile(process.env.OPENCODE_GOAL_STATE_PATH!, "utf8")) + expect(persisted.goals.ses_v2.status).toBe("cancelled") + expect(persisted.goals.ses_v2.history.filter((entry: { type: string }) => entry.type === "cancelled")).toHaveLength(1) + mock.stream.end() + await cleanup() + + const reloaded = makeMockContext({ min_continue_interval_seconds: 0, max_turn_time: 0.02 }) + await setupPlugin(reloaded as never) + await reloaded.stream.push({ type: "session.execution.started", created: 5, data: { sessionID: "ses_v2" } }) + await reloaded.stream.push({ type: "session.execution.succeeded", created: 6, data: { sessionID: "ses_v2" } }) + await new Promise((resolve) => setTimeout(resolve, 100)) + expect(reloaded.promptCalls).toHaveLength(0) + expect((await getGoal("ses_v2"))?.status).toBe("cancelled") + // An explicit new goal in the same session must still work. + await createGoalViaV2Tool(reloaded, "a new user-requested goal") + await reloaded.stream.push({ type: "session.execution.succeeded", created: 7, data: { sessionID: "ses_v2" } }) + await waitFor(() => reloaded.promptCalls.length === 1) + expect(reloaded.promptCalls[0]?.text).toContain("a new user-requested goal") +}) + +test("V2 user cancellation persists even with auto-continue disabled", async () => { + const mock = makeMockContext({ auto_continue: false }) + await setupPlugin(mock as never) + await createGoalViaV2Tool(mock, "an explicitly cancelled manual goal") + await mock.stream.push({ type: "session.execution.interrupted", created: 1, data: { sessionID: "ses_v2", reason: "user" } }) + expect((await getGoal("ses_v2"))?.status).toBe("cancelled") + expect(mock.promptCalls).toHaveLength(0) +}) + +test("V2 another location's cancellation does not cancel the owner's goal", async () => { + await createGoal("ses_other", "a goal owned by another location") + const mock = makeMockContext({}, [], {}, { directory: "/own" }, { + ses_other: { location: { directory: "/other" } }, + }) + await setupPlugin(mock as never) + await mock.stream.push({ type: "session.execution.interrupted", created: 1, data: { sessionID: "ses_other", reason: "user" } }) + expect((await getGoal("ses_other"))?.status).toBe("active") + expect(mock.promptCalls).toHaveLength(0) +}) + +test("V2 cancellation rejects late recovery results without affecting a replacement", async () => { + let rejectOldPrompt: (() => void) | undefined + const mock = makeMockContext({ min_continue_interval_seconds: 0 }) + mock.session.prompt = async (input) => { + mock.promptCalls.push(input) + if (mock.promptCalls.length !== 1) return + await new Promise((_resolve, reject) => { + rejectOldPrompt = () => reject(new Error("network connection failed")) + }) + } + await setupPlugin(mock as never) + await createGoalViaV2Tool(mock, "the old goal") + // Recovery runs on a timer, allowing the host event consumer to observe Cancel. + await mock.stream.push({ type: "session.execution.failed", created: 1, data: { sessionID: "ses_v2", error: { message: "network connection failed" } } }) + await waitFor(() => rejectOldPrompt != null) + await mock.stream.push({ type: "session.execution.interrupted", created: 2, data: { sessionID: "ses_v2", reason: "user" } }) + expect((await getGoal("ses_v2"))?.status).toBe("cancelled") + await createGoalViaV2Tool(mock, "the replacement goal") + rejectOldPrompt?.() + await waitFor(() => mock.promptCalls.length === 2) + expect(await getGoal("ses_v2")).toMatchObject({ objective: "the replacement goal", status: "active", continuationFailures: 0, autoTurns: 1 }) + expect(mock.promptCalls[1]?.text).toContain("the replacement goal") +}) + test("V2 terminal execution transport failure recovers after native retries and respects the failure ceiling", async () => { const mock = makeMockContext({ min_continue_interval_seconds: 0, max_prompt_failures: 1 }) const cleanup = await setupPlugin(mock as never) diff --git a/test/server.test.ts b/test/server.test.ts index d24a837..d687651 100644 --- a/test/server.test.ts +++ b/test/server.test.ts @@ -77,12 +77,70 @@ beforeEach(async () => { process.env.OPENCODE_GOAL_STATE_PATH = join(dir, "goals.json") }) +for (const signal of ["session.error", "message.updated"]) { + test(`V1 ${signal} user abort persists cancellation and prevents later continuations`, async () => { + const calls: unknown[] = [] + const client = { session: { promptAsync: async (input: unknown) => { calls.push(input) } } } + const hooks = await setupServer({ client } as never, { min_continue_interval_seconds: 0, max_turn_time: 0.02 }) + await requireTool(hooks.tool?.create_goal, "create_goal").execute({ objective: "respect cancellation" }, { sessionID: "ses_cancel" } as never) + await hooks.event!({ event: { type: "session.status", properties: { sessionID: "ses_cancel", status: { type: "busy" } } } } as never) + const error = { name: "MessageAbortedError", data: { message: "The operation was aborted." } } + const properties = signal === "session.error" + ? { sessionID: "ses_cancel", error } + : { info: { id: "msg_cancel", sessionID: "ses_cancel", role: "assistant", error, time: { completed: Date.now() } } } + await hooks.event!({ event: { type: signal, properties } } as never) + await hooks.event!({ event: { type: "session.status", properties: { sessionID: "ses_cancel", status: { type: "idle" } } } } as never) + await new Promise((resolve) => setTimeout(resolve, 100)) + expect(calls).toHaveLength(0) + expect(await getGoalInternal("ses_cancel")).toMatchObject({ status: "cancelled", pendingAttempt: null, continuationFailures: 0 }) + const persisted = JSON.parse(await readFile(process.env.OPENCODE_GOAL_STATE_PATH!, "utf8")) + expect(persisted.goals.ses_cancel.status).toBe("cancelled") + await hooks.dispose?.() + const reloaded = await setupServer({ client } as never, { min_continue_interval_seconds: 0 }) + await reloaded.event!({ event: { type: "session.status", properties: { sessionID: "ses_cancel", status: { type: "busy" } } } } as never) + await reloaded.event!({ event: { type: "session.idle", properties: { sessionID: "ses_cancel" } } } as never) + expect(calls).toHaveLength(0) + }) +} + afterEach(async () => { for (const dispose of serverDisposers.splice(0).reverse()) await dispose() delete process.env.OPENCODE_GOAL_STATE_PATH await rm(dir, { recursive: true, force: true }) }) +test("V1 cancellation invalidates a continuation still reading the transcript", async () => { + let releaseTranscript: (() => void) | undefined + const calls: unknown[] = [] + const hooks = await setupServer({ client: { session: { + messages: async () => { + await new Promise((resolve) => { releaseTranscript = resolve }) + return { data: [] } + }, + promptAsync: async (input: unknown) => { calls.push(input) }, + } } } as never, { min_continue_interval_seconds: 0 }) + await requireTool(hooks.tool?.create_goal, "create_goal").execute({ objective: "do not restart after cancellation" }, { sessionID: "ses_cancel_read" } as never) + const idle = hooks.event!({ event: { type: "session.idle", properties: { sessionID: "ses_cancel_read" } } } as never) + await waitFor(() => releaseTranscript != null) + await hooks.event!({ event: { type: "session.error", properties: { sessionID: "ses_cancel_read", error: { name: "MessageAbortedError" } } } } as never) + releaseTranscript?.() + await idle + expect(calls).toHaveLength(0) + expect((await getGoal("ses_cancel_read"))?.status).toBe("cancelled") +}) + +test("V1 only a named session abort cancels the goal", async () => { + const calls: unknown[] = [] + const hooks = await setupServer({ client: { session: { + promptAsync: async (input: unknown) => { calls.push(input) }, + } } } as never, { min_continue_interval_seconds: 0 }) + await requireTool(hooks.tool?.create_goal, "create_goal").execute({ objective: "recover ordinary errors" }, { sessionID: "ses_non_cancel" } as never) + await hooks.event!({ event: { type: "session.error", properties: { sessionID: "ses_non_cancel", error: { name: "APIError", data: { message: "Upstream aborted request" } } } } } as never) + expect((await getGoal("ses_non_cancel"))?.status).toBe("active") + await hooks.event!({ event: { type: "session.idle", properties: { sessionID: "ses_non_cancel" } } } as never) + expect(calls).toHaveLength(1) +}) + test("server plugin exposes Codex-style goal tools", async () => { const calls: unknown[] = [] const hooks = await setupServer( From aa99843c8916abe9815de33c3cc16fbd2b4da558 Mon Sep 17 00:00:00 2001 From: Daniel Saldarriaga Date: Mon, 5 Oct 2026 15:16:33 +0200 Subject: [PATCH 2/2] fix: preserve paused and limited goals on host cancellation --- README.md | 2 +- dist/server.js | 14 ++++++++++++-- src/server.ts | 5 +++-- src/state.ts | 10 ++++++++++ test/server-v2.test.ts | 27 ++++++++++++++++++++++++++- test/server.test.ts | 31 +++++++++++++++++++++++++++++++ test/state.test.ts | 12 ++++++++++++ 7 files changed, 95 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index 906ca9b..288be8b 100644 --- a/README.md +++ b/README.md @@ -315,7 +315,7 @@ OpenCode plugin modules are target-specific. This package exports separate modul } ``` -Codex goal mode has deeper runtime integration for thread lifecycle control. This plugin implements the same workflow using OpenCode plugin hooks. Token usage is read from OpenCode step-finish usage when available and falls back to message token metadata or text estimation when exact usage is unavailable. Continuation is driven by V2 `session.execution.succeeded` events and legacy `session.idle` / `session.status` idle notifications, never by intermediate model-step completion. V2 execution starts arm busy tracking; native `session.retry.scheduled` events cancel plugin recovery while OpenCode retries. Terminal execution transport failures use bounded recovery. A V2 `session.execution.interrupted` event with reason `user`, or a V1 session/assistant `MessageAbortedError`, persists a terminal `cancelled` goal and invalidates timers and outstanding continuation preparation without charging a prompt failure. A later idle, unrelated prompt, or plugin reload cannot resume that goal; start an explicitly requested new goal with `/goal ` or `/goal replace `. Shutdown, superseded, and other interruptions only suppress local continuation until the host starts another execution; they do not cancel the persisted goal. Use `/pause_goal` for a durable, resumable pause. Each V2 plugin instance handles goal events only for its own location while still observing cross-location child task lifecycles. The optional `max_turn_time` watchdog can retry one goal continuation prompt when a model turn remains busy, without consuming the goal's auto-turn or no-progress budgets; recognized transport failures do count toward the prompt-failure ceiling. By default, continuation is deferred while OpenCode Task child sessions are active or their terminal result still needs an orchestrator turn, bounded by the `max_task_block_seconds` ceiling so an unobservable child cannot stall a goal indefinitely. During compaction on OpenCode 1, the plugin disables OpenCode's generic synthetic auto-continue while an active goal exists so the goal-specific continuation prompt remains authoritative; on OpenCode 2, compaction runs inside the session's execution flow, so the plugin instead injects the goal snapshot through the `session.compaction` hook where the host provides it. +Codex goal mode has deeper runtime integration for thread lifecycle control. This plugin implements the same workflow using OpenCode plugin hooks. Token usage is read from OpenCode step-finish usage when available and falls back to message token metadata or text estimation when exact usage is unavailable. Continuation is driven by V2 `session.execution.succeeded` events and legacy `session.idle` / `session.status` idle notifications, never by intermediate model-step completion. V2 execution starts arm busy tracking; native `session.retry.scheduled` events cancel plugin recovery while OpenCode retries. Terminal execution transport failures use bounded recovery. A V2 `session.execution.interrupted` event with reason `user`, or a V1 session/assistant `MessageAbortedError`, persists a terminal `cancelled` state for an active goal and invalidates timers and outstanding continuation preparation without charging a prompt failure. A later idle, unrelated prompt, or plugin reload cannot resume that goal; start an explicitly requested new goal with `/goal ` or `/goal replace `. Cancelling a manual turn preserves paused or limited goals. V2 shutdown, superseded, and other interruptions only suppress local continuation until the host starts another execution; they do not cancel the persisted goal. V1 exposes only `MessageAbortedError`, so it cannot distinguish user cancellation from other host aborts of an active goal. Use `/pause_goal` for a durable, resumable pause. Each V2 plugin instance handles goal events only for its own location while still observing cross-location child task lifecycles. The optional `max_turn_time` watchdog can retry one goal continuation prompt when a model turn remains busy, without consuming the goal's auto-turn or no-progress budgets; recognized transport failures do count toward the prompt-failure ceiling. By default, continuation is deferred while OpenCode Task child sessions are active or their terminal result still needs an orchestrator turn, bounded by the `max_task_block_seconds` ceiling so an unobservable child cannot stall a goal indefinitely. During compaction on OpenCode 1, the plugin disables OpenCode's generic synthetic auto-continue while an active goal exists so the goal-specific continuation prompt remains authoritative; on OpenCode 2, compaction runs inside the session's execution flow, so the plugin instead injects the goal snapshot through the `session.compaction` hook where the host provides it. The goal sidebar shows the current status, elapsed time, token usage, auto-continue count, latest checkpoint, latest status message, stop reason, and objective when a goal is active, paused, or safety-limited. It checks the shared goal state file every second so usage and checkpoints stay current during a long run. Closed goals remain visible briefly through the latest tool state as achieved or unmet. diff --git a/dist/server.js b/dist/server.js index 0a38667..8c4b3c0 100644 --- a/dist/server.js +++ b/dist/server.js @@ -888,6 +888,16 @@ async function cancelGoal(sessionID, reason = "cancelled") { return snapshot(goal); }); } +async function cancelActiveGoal(sessionID) { + return mutate((state) => { + const goal = state.goals[sessionID]; + if (!goal) + return null; + if (goal.status === "active") + cancelGoalRecord(goal, "cancelled"); + return snapshot(goal); + }); +} async function clearGoal(sessionID) { return mutate((state) => { const goal = state.goals[sessionID]; @@ -3552,7 +3562,7 @@ var server = async ({ client }, options) => { watchdogRescuedSessions.delete(sessionID); clearToolAttemptsForSession(toolAttempts, sessionID); taskTracker.observeSessionStatus(sessionID, "idle"); - await cancelGoal(sessionID); + await cancelActiveGoal(sessionID); return; } if (eventType === "session.created") { @@ -4174,7 +4184,7 @@ async function setupV2(context) { clearToolAttemptsForSession(toolAttempts, sessionID); taskTracker.observeSessionStatus(sessionID, "idle"); if (data.reason === "user") - await cancelGoal(sessionID); + await cancelActiveGoal(sessionID); return; } case "session.execution.failed": { diff --git a/src/server.ts b/src/server.ts index f9e7f94..9843fcc 100644 --- a/src/server.ts +++ b/src/server.ts @@ -8,6 +8,7 @@ import type { GoalSnapshot, InternalGoalSnapshot, PendingAttempt } from "./state import { accountUsage, cancelGoal, + cancelActiveGoal, clearGoal, completeGoal, createGoal, @@ -2016,7 +2017,7 @@ const server: Plugin = async ({ client }, options?: Options) => { watchdogRescuedSessions.delete(sessionID) clearToolAttemptsForSession(toolAttempts, sessionID) taskTracker.observeSessionStatus(sessionID, "idle") - await cancelGoal(sessionID) + await cancelActiveGoal(sessionID) return } if (eventType === "session.created") { @@ -2732,7 +2733,7 @@ async function setupV2(context: PluginV2.Plugin.Context): Promise { + const goal = state.goals[sessionID] + if (!goal) return null + if (goal.status === "active") cancelGoalRecord(goal, "cancelled") + return snapshot(goal) + }) +} + export async function clearGoal(sessionID: string) { return mutate((state) => { const goal = state.goals[sessionID] diff --git a/test/server-v2.test.ts b/test/server-v2.test.ts index a3d1d39..755747c 100644 --- a/test/server-v2.test.ts +++ b/test/server-v2.test.ts @@ -3,7 +3,7 @@ import { mkdtemp, readFile, readdir, rm, writeFile } from "node:fs/promises" import { join } from "node:path" import { tmpdir } from "node:os" import plugin from "../src/server" -import { cancelGoal, createGoal, getGoal, getGoalInternal, recordContinuationResult, reserveContinuation } from "../src/state" +import { accountUsage, pauseGoalForPlanMode, setGoalStatus, cancelGoal, createGoal, getGoal, getGoalInternal, recordContinuationResult, reserveContinuation } from "../src/state" const TOOL_NAMES = [ "clear_goal", @@ -1815,6 +1815,31 @@ test("V2 user cancellation persists across reload and unrelated executions", asy expect(reloaded.promptCalls[0]?.text).toContain("a new user-requested goal") }) +for (const status of ["paused", "plan", "budgetLimited", "usageLimited"] as const) { + test(`V2 preserves ${status} goals when a manual execution is cancelled`, async () => { + const mock = makeMockContext({ min_continue_interval_seconds: 0 }) + await setupPlugin(mock as never) + await createGoal("ses_v2", "remain available after cancelling a manual turn", { tokenBudget: status === "budgetLimited" ? 1 : null }) + if (status === "paused") await setGoalStatus("ses_v2", "paused") + if (status === "plan") await pauseGoalForPlanMode("ses_v2") + if (status === "budgetLimited") await accountUsage("ses_v2", 2) + if (status === "usageLimited") { + await reserveContinuation("ses_v2", 1, 0) + await reserveContinuation("ses_v2", 1, 0) + } + const before = await getGoalInternal("ses_v2") + expect(before?.status).toBe(status === "plan" ? "paused" : status) + await mock.stream.push({ type: "session.execution.started", created: 1, data: { sessionID: "ses_v2" } }) + await mock.stream.push({ type: "session.execution.interrupted", created: 2, data: { sessionID: "ses_v2", reason: "user" } }) + await mock.stream.push({ type: "session.idle", created: 3, data: { sessionID: "ses_v2" } }) + expect(await getGoalInternal("ses_v2")).toMatchObject({ id: before?.id, status: before?.status, stopReason: before?.stopReason, closedAt: null }) + expect(mock.promptCalls).toHaveLength(0) + if (status === "paused" || status === "plan") { + expect((await setGoalStatus("ses_v2", "active"))?.status).toBe("active") + } + }) +} + test("V2 user cancellation persists even with auto-continue disabled", async () => { const mock = makeMockContext({ auto_continue: false }) await setupPlugin(mock as never) diff --git a/test/server.test.ts b/test/server.test.ts index d687651..4e12e69 100644 --- a/test/server.test.ts +++ b/test/server.test.ts @@ -6,6 +6,9 @@ import { z } from "zod" import plugin from "../src/server" import { accountUsage, + createGoal, + pauseGoalForPlanMode, + setGoalStatus, getGoal, getGoalInternal, recordContinuationResult, @@ -109,6 +112,34 @@ afterEach(async () => { await rm(dir, { recursive: true, force: true }) }) +for (const signal of ["session.error", "message.updated"]) { + for (const status of ["paused", "plan", "budgetLimited", "usageLimited"] as const) { + test(`V1 ${signal} preserves ${status} goals when a manual turn is aborted`, async () => { + const calls: unknown[] = [] + const hooks = await setupServer({ client: { session: { promptAsync: async (input: unknown) => { calls.push(input) } } } } as never) + const sessionID = "ses_manual_abort" + await createGoal(sessionID, "remain available after aborting a manual turn", { tokenBudget: status === "budgetLimited" ? 1 : null }) + if (status === "paused") await setGoalStatus(sessionID, "paused") + if (status === "plan") await pauseGoalForPlanMode(sessionID) + if (status === "budgetLimited") await accountUsage(sessionID, 2) + if (status === "usageLimited") { + await reserveContinuation(sessionID, 1, 0) + await reserveContinuation(sessionID, 1, 0) + } + const before = await getGoalInternal(sessionID) + expect(before?.status).toBe(status === "plan" ? "paused" : status) + const error = { name: "MessageAbortedError" } + const properties = signal === "session.error" ? { sessionID, error } : { info: { sessionID, role: "assistant", error } } + await hooks.event!({ event: { type: signal, properties } } as never) + expect(await getGoalInternal(sessionID)).toEqual(before) + expect(calls).toHaveLength(0) + if (status === "paused" || status === "plan") { + expect((await setGoalStatus(sessionID, "active"))?.status).toBe("active") + } + }) + } +} + test("V1 cancellation invalidates a continuation still reading the transcript", async () => { let releaseTranscript: (() => void) | undefined const calls: unknown[] = [] diff --git a/test/state.test.ts b/test/state.test.ts index b746364..1e8ea7c 100644 --- a/test/state.test.ts +++ b/test/state.test.ts @@ -5,6 +5,7 @@ import { tmpdir } from "node:os" import { accountUsage, cancelGoal, + cancelActiveGoal, clearGoal, completeGoal, createGoal, @@ -86,6 +87,17 @@ test("cancels, clears, and replaces goals while preserving per-session history", expect((await getGoalHistory("ses_1")).previous).toHaveLength(2) }) +test("host cancellation observes a queued pause atomically while explicit stop can still close it", async () => { + await createGoal("ses_1", "preserve the pause contract", null) + const [, result] = await Promise.all([setGoalStatus("ses_1", "paused"), cancelActiveGoal("ses_1")]) + expect(result?.status).toBe("paused") + expect((await setGoalStatus("ses_1", "active")).status).toBe("active") + expect((await cancelActiveGoal("ses_1"))?.status).toBe("cancelled") + await createGoal("ses_1", "explicit stop may close a paused goal", null) + await setGoalStatus("ses_1", "paused") + expect((await cancelGoal("ses_1"))?.status).toBe("cancelled") +}) + test("closed and cancelled goals cannot be edited or closed again", async () => { await createGoal("ses_1", "do not reopen", null) await cancelGoal("ses_1")