From 5c80eb24a78b859919e9c8ff964bd58cd8c9a543 Mon Sep 17 00:00:00 2001 From: Mux Date: Tue, 1 Sep 2026 12:19:20 -0500 Subject: [PATCH 1/4] =?UTF-8?q?=F0=9F=A4=96=20fix:=20preserve=20workspace?= =?UTF-8?q?=20turns=20across=20synthetic=20wake=20ends?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- _Generated with `xum` • Model: `openai:gpt-5.6-sol` • Thinking: `xhigh` • Cost: `$0.00`_ --- src/node/services/taskHandleStore.ts | 6 + .../services/workspaceTurnManager.test.ts | 140 +++++++++++++ src/node/services/workspaceTurnManager.ts | 198 ++++++++++++++---- 3 files changed, 304 insertions(+), 40 deletions(-) diff --git a/src/node/services/taskHandleStore.ts b/src/node/services/taskHandleStore.ts index 1a3c115ae4..55c951b1e4 100644 --- a/src/node/services/taskHandleStore.ts +++ b/src/node/services/taskHandleStore.ts @@ -72,6 +72,11 @@ export interface WorkspaceTurnTaskHandleRecord { metadata: StreamEndEvent["metadata"]; }; deferredMessageIds?: string[]; + /** + * True only for a terminal result from an uncorrelated synthetic wake stream-end. + * A later correlated stream-end can replace this provisional result. + */ + provisionalOutcome?: boolean; error?: string; /** * How the owner workspace's stream-end treats this workspace turn while active. @@ -118,6 +123,7 @@ const WorkspaceTurnTaskHandleRecordSchema = z .passthrough() .optional(), deferredMessageIds: z.array(z.string().min(1)).optional(), + provisionalOutcome: z.boolean().optional(), error: z.string().optional(), attentionPolicy: BackgroundWorkAttentionPolicySchema.optional(), directParentResultDeliveryRequiredAt: z.string().optional(), diff --git a/src/node/services/workspaceTurnManager.test.ts b/src/node/services/workspaceTurnManager.test.ts index 20588945ad..205b716c0d 100644 --- a/src/node/services/workspaceTurnManager.test.ts +++ b/src/node/services/workspaceTurnManager.test.ts @@ -24,6 +24,7 @@ import type { AIService } from "@/node/services/aiService"; import type { WorkspaceHost, BackgroundableForegroundWaiter, + QueueCutAttributionSnapshot, WorkspaceTurnManagerHost, } from "@/node/services/taskWorkspaceSeam"; import type { InitStateManager } from "@/node/services/initStateManager"; @@ -362,6 +363,63 @@ describe("WorkspaceTurnManager", () => { }; } + async function finalizeWorkspaceTurnStreamEndForTest( + taskService: WorkspaceTurnManager, + event: StreamEndEvent + ): Promise { + const internal = taskService as unknown as { + captureQueueCutAttributionSnapshot: (workspaceId: string) => QueueCutAttributionSnapshot; + finalizeWorkspaceTurnFromStreamEnd: ( + event: StreamEndEvent, + queueCutSnapshot: QueueCutAttributionSnapshot + ) => Promise; + }; + return await internal.finalizeWorkspaceTurnFromStreamEnd( + event, + internal.captureQueueCutAttributionSnapshot(event.workspaceId) + ); + } + + async function appendUncorrelatedWakeHistory(params: { + historyService: HistoryService; + parentId: string; + taskId: string; + workspaceId: string; + inputSynthetic: boolean; + }): Promise { + const prompt = createMuxMessage("turn-prompt", "user", "Summarize the repo", { + muxMetadata: workspaceTurnMuxMetadata(params.parentId, params.taskId), + }); + expect((await params.historyService.appendToHistory(params.workspaceId, prompt)).success).toBe( + true + ); + const input = createMuxMessage("wake-input", "user", "Continue after a wake", { + synthetic: params.inputSynthetic, + }); + expect((await params.historyService.appendToHistory(params.workspaceId, input)).success).toBe( + true + ); + const output = createMuxMessage("wake-output", "assistant", "Wake result", { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + finishReason: "stop", + }); + expect((await params.historyService.appendToHistory(params.workspaceId, output)).success).toBe( + true + ); + return { + type: "stream-end", + workspaceId: params.workspaceId, + messageId: output.id, + metadata: { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + finishReason: "stop", + }, + parts: [{ type: "text", text: "Wake result" }], + }; + } + async function createWorkspaceLifecycleHarness( options: { archived?: boolean; @@ -3837,6 +3895,88 @@ describe("WorkspaceTurnManager", () => { expect(snapshot).toMatchObject({ status: "running", workspaceId: created.workspaceId }); }); + test("uncorrelated synthetic wake end stays active while a continuation is live", async () => { + let monitorWakePending = true; + const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest({ + hasPendingBashMonitorWakeContinuation: mock(() => monitorWakePending), + }); + const event = await appendUncorrelatedWakeHistory({ + historyService, + parentId, + taskId: created.taskId, + workspaceId: created.workspaceId, + inputSynthetic: true, + }); + + expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, event)).toBe(true); + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ + status: "running", + deferredMessageIds: ["wake-output"], + }); + + monitorWakePending = false; + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ + status: "completed", + messageId: "wake-output", + reportMarkdown: "Wake result", + }); + }); + + test("idle uncorrelated synthetic wake end settles provisionally from the wake output", async () => { + const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); + const event = await appendUncorrelatedWakeHistory({ + historyService, + parentId, + taskId: created.taskId, + workspaceId: created.workspaceId, + inputSynthetic: true, + }); + + expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, event)).toBe(true); + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ + status: "completed", + messageId: "wake-output", + reportMarkdown: "Wake result", + provisionalOutcome: true, + }); + + const correlatedFinal: StreamEndEvent = { + ...event, + messageId: "real-final", + metadata: { + ...event.metadata, + muxMetadata: workspaceTurnMuxMetadata(parentId, created.taskId), + }, + parts: [{ type: "text", text: "Real final result" }], + }; + expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, correlatedFinal)).toBe(true); + const corrected = await workspaceTurnSnapshot(taskService, parentId, created.taskId); + expect(corrected).toMatchObject({ + status: "completed", + messageId: "real-final", + reportMarkdown: "Real final result", + }); + expect(corrected?.provisionalOutcome).toBeUndefined(); + }); + + test("manual user input still supersedes an active workspace turn on uncorrelated end", async () => { + const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); + const event = await appendUncorrelatedWakeHistory({ + historyService, + parentId, + taskId: created.taskId, + workspaceId: created.workspaceId, + inputSynthetic: false, + }); + + expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, event)).toBe(true); + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ + status: "interrupted", + messageId: "wake-output", + error: "Workspace turn superseded by an uncorrelated workspace stream-end", + }); + }); + test("getWorkspaceTurnSnapshot recovers stale completed handles from matching history", async () => { const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); const appendResult = await historyService.appendToHistory( diff --git a/src/node/services/workspaceTurnManager.ts b/src/node/services/workspaceTurnManager.ts index b74f3b9d6e..0d5eef8f12 100644 --- a/src/node/services/workspaceTurnManager.ts +++ b/src/node/services/workspaceTurnManager.ts @@ -281,6 +281,11 @@ const WORKSPACE_TURN_STALE_RESTART_ERROR = "Workspace turn interrupted after res const WORKSPACE_TURN_SUPERSEDED_BY_NEW_INPUT_ERROR = "Workspace turn superseded by new input in the target workspace; the workspace continues under that input and this delegated turn will not report"; +/** A human-authored child input that redirects the delegated turn. */ +function isManualChildWorkspaceInput(message: MuxMessage): boolean { + return message.role === "user" && message.metadata?.synthetic !== true; +} + /** * Reason prefix persisted when the owner's OWN follow-up turn (task * kind="workspace", mode="existing", tool-end dispatch) cut its active @@ -2147,11 +2152,12 @@ export class WorkspaceTurnManager { // error / stale restart interrupt — never an explicit user interrupt) may be // corrected once by an explicitly allowed resettle, but only when the new settlement // actually changes the outcome (duplicate stream-end replays must stay idempotent). + const currentIsProvisional = current.provisionalOutcome === true; const resettleStaleTerminal = params.allowTerminalResettle === true && this.isTerminalWorkspaceTurnStatus(current.status) && - current.status !== "completed" && - isSelfHealEligibleSettledWorkspaceTurn(current) && + (current.status !== "completed" || currentIsProvisional) && + (currentIsProvisional || isSelfHealEligibleSettledWorkspaceTurn(current)) && (params.next.status !== current.status || params.next.messageId !== current.messageId); if (this.isTerminalWorkspaceTurnStatus(current.status) && !resettleStaleTerminal) { const active = this.activeWorkspaceTurnHandleByWorkspaceId.get(params.record.workspaceId); @@ -2231,6 +2237,7 @@ export class WorkspaceTurnManager { nextStatus: nextRecord.status, }); delete nextRecord.terminalAttentionNotifiedAt; + delete nextRecord.provisionalOutcome; } if (suppressedQuietSettlement) { // Downgrade-compatible suppression marker (upgrade↔downgrade rule): @@ -3804,7 +3811,9 @@ export class WorkspaceTurnManager { const active = this.activeWorkspaceTurnHandleByWorkspaceId.get(record.workspaceId); const hasRuntimeActivity = this.aiService.isStreaming(record.workspaceId) || - this.workspaceService.hasPendingQueuedOrPreparingTurn(record.workspaceId); + this.workspaceService.hasPendingQueuedOrPreparingTurn(record.workspaceId) || + this.workspaceService.hasPendingAutoRetry(record.workspaceId) || + this.workspaceService.hasPendingBashMonitorWakeContinuation(record.workspaceId); if (hasRuntimeActivity) { return true; } @@ -4065,6 +4074,26 @@ export class WorkspaceTurnManager { }; } + /** Rebuild a deferred uncorrelated wake end after its continuation evidence clears. */ + private buildDeferredWakeEndEventFromHistory( + record: WorkspaceTurnTaskHandleRecord, + message: MuxMessage + ): StreamEndEvent | null { + if (message.role !== "assistant" || message.metadata?.partial === true) { + return null; + } + return { + type: "stream-end", + workspaceId: record.workspaceId, + messageId: message.id, + metadata: { + ...message.metadata, + model: coerceNonEmptyString(message.metadata?.model) ?? record.modelString ?? defaultModel, + }, + parts: message.parts as StreamEndEvent["parts"], + }; + } + private buildTerminalWorkspaceTurnRecordFromEvent( record: WorkspaceTurnTaskHandleRecord, event: StreamEndEvent, @@ -4073,6 +4102,7 @@ export class WorkspaceTurnManager { const baseRecord = { ...record }; delete baseRecord.error; delete baseRecord.deferredMessageIds; + delete baseRecord.provisionalOutcome; // A "tool-calls" finish on a delegated turn backed by queue-dispatch // evidence is a queue cut: some other queued input (a manual user message, // /compact, the owner's own follow-up turn) dispatched at the tool boundary @@ -4159,10 +4189,13 @@ export class WorkspaceTurnManager { const allowDeferredMessages = !(await this.hasActiveWorkspaceTurnDeferredBlockers(record)); for (const message of historyResult.data.toReversed()) { - if (this.isDeferredWorkspaceTurnMessage(record, message.id) && !allowDeferredMessages) { + const isDeferred = this.isDeferredWorkspaceTurnMessage(record, message.id); + if (isDeferred && !allowDeferredMessages) { continue; } - const event = this.buildWorkspaceTurnStreamEndEventFromHistory(record, message); + const event = + this.buildWorkspaceTurnStreamEndEventFromHistory(record, message) ?? + (isDeferred ? this.buildDeferredWakeEndEventFromHistory(record, message) : null); if (event != null) { // History order alone cannot prove a queue cut (a later unrelated user // message is not causal evidence), so stale recovery conservatively @@ -4217,39 +4250,6 @@ export class WorkspaceTurnManager { }; } - private async isStreamEndBeforeWorkspaceTurnPrompt( - record: WorkspaceTurnTaskHandleRecord, - event: StreamEndEvent - ): Promise { - const historyResult = await this.historyService.getHistoryFromLatestBoundary(event.workspaceId); - if (!historyResult.success) { - log.warn("Could not compare uncorrelated stream-end history for workspace turn", { - workspaceId: event.workspaceId, - handleId: record.handleId, - error: historyResult.error, - }); - return false; - } - - let streamEndIndex = -1; - let promptIndex = -1; - for (const [index, message] of historyResult.data.entries()) { - if (message.id === event.messageId) { - streamEndIndex = index; - } - const metadata = this.getWorkspaceTurnMetadataFromValue(message.metadata?.muxMetadata); - if ( - metadata?.taskHandleId === record.handleId && - metadata.ownerWorkspaceId === record.ownerWorkspaceId && - metadata.turnId === record.turnId - ) { - promptIndex = index; - } - } - - return streamEndIndex !== -1 && promptIndex !== -1 && streamEndIndex < promptIndex; - } - private async interruptWorkspaceTurnFromUncorrelatedStreamEnd( event: StreamEndEvent ): Promise { @@ -4281,7 +4281,34 @@ export class WorkspaceTurnManager { return true; } - if (await this.isStreamEndBeforeWorkspaceTurnPrompt(record, event)) { + const historyResult = await this.historyService.getHistoryFromLatestBoundary(event.workspaceId); + if (!historyResult.success) { + log.warn("Could not compare uncorrelated stream-end history for workspace turn", { + workspaceId: event.workspaceId, + handleId: record.handleId, + error: historyResult.error, + }); + await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); + return true; + } + + let streamEndIndex = -1; + let promptIndex = -1; + for (const [index, message] of historyResult.data.entries()) { + if (message.id === event.messageId) { + streamEndIndex = index; + } + const metadata = this.getWorkspaceTurnMetadataFromValue(message.metadata?.muxMetadata); + if ( + metadata?.taskHandleId === record.handleId && + metadata.ownerWorkspaceId === record.ownerWorkspaceId && + metadata.turnId === record.turnId + ) { + promptIndex = index; + } + } + + if (streamEndIndex !== -1 && promptIndex !== -1 && streamEndIndex < promptIndex) { log.debug("Ignoring stale uncorrelated stream-end before queued workspace turn prompt", { workspaceId: event.workspaceId, taskHandleId: record.handleId, @@ -4290,6 +4317,96 @@ export class WorkspaceTurnManager { return true; } + if (promptIndex !== -1) { + const scanEnd = streamEndIndex === -1 ? historyResult.data.length : streamEndIndex; + const hasManualSupersessionInput = historyResult.data + .slice(promptIndex + 1, scanEnd) + .some(isManualChildWorkspaceInput); + if (hasManualSupersessionInput) { + await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); + return true; + } + + if (await this.hasLiveUncorrelatedWakeContinuationEvidence(event, record)) { + await this.persistDeferredUncorrelatedWakeEnd(record, event); + return true; + } + + const next = this.buildTerminalWorkspaceTurnRecordFromEvent(record, event, { + supersedeEvidence: null, + }); + next.provisionalOutcome = true; + await this.settleWorkspaceTurn({ + record, + next, + waiterSettlement: + next.status === "completed" + ? { status: "completed", result: this.buildWorkspaceTurnWaitResult(next) } + : { status: "error", error: new Error(next.error ?? "Workspace turn failed") }, + allowTerminalResettle: true, + }); + return true; + } + + await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); + return true; + } + + private async hasLiveUncorrelatedWakeContinuationEvidence( + event: StreamEndEvent, + record: WorkspaceTurnTaskHandleRecord + ): Promise { + if (await this.hasActiveWorkspaceTurnDeferredBlockers(record)) { + return true; + } + if ( + this.workspaceService.hasPendingQueuedOrPreparingTurn(event.workspaceId) || + this.workspaceService.hasPendingAutoRetry(event.workspaceId) || + this.workspaceService.hasPendingBashMonitorWakeContinuation(event.workspaceId) + ) { + return true; + } + if ( + this.workspaceService.hasPendingWorkspaceTurnContinuation(event.workspaceId, { + type: "workspace-turn-task", + taskHandleId: record.handleId, + ownerWorkspaceId: record.ownerWorkspaceId, + turnId: record.turnId, + }) + ) { + return true; + } + const activeStream = this.streamManager?.getStreamInfo(event.workspaceId); + return activeStream != null && activeStream.messageId !== event.messageId; + } + + private async persistDeferredUncorrelatedWakeEnd( + record: WorkspaceTurnTaskHandleRecord, + event: StreamEndEvent + ): Promise { + await this.workspaceTurnSettlementLocks.withLock(record.handleId, async () => { + const current = await this.taskHandleStore.getWorkspaceTurn( + record.ownerWorkspaceId, + record.handleId + ); + if (current == null || !isActiveWorkspaceTurnTaskStatus(current.status)) { + return; + } + if (this.isDeferredWorkspaceTurnMessage(current, event.messageId)) { + return; + } + await this.taskHandleStore.upsertWorkspaceTurn({ + ...current, + updatedAt: getIsoNow(), + deferredMessageIds: [...(current.deferredMessageIds ?? []), event.messageId], + }); + }); + } + + private async settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd( + record: WorkspaceTurnTaskHandleRecord, + event: StreamEndEvent + ): Promise { const error = "Workspace turn superseded by an uncorrelated workspace stream-end"; const next: WorkspaceTurnTaskHandleRecord = { ...record, @@ -4298,12 +4415,13 @@ export class WorkspaceTurnManager { messageId: event.messageId, error, }; + delete next.provisionalOutcome; + delete next.deferredMessageIds; await this.settleWorkspaceTurn({ record, next, waiterSettlement: { status: "error", error: new Error(error) }, }); - return true; } /** From 71efd6c2e06ccd51caa8c43de75eae7443ce39d4 Mon Sep 17 00:00:00 2001 From: Mux Date: Tue, 1 Sep 2026 15:19:26 -0500 Subject: [PATCH 2/4] =?UTF-8?q?=F0=9F=A4=96=20fix:=20keep=20synthetic=20wo?= =?UTF-8?q?rkspace=20wakes=20nonterminal?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- _Generated with `xum` • Model: `openai:gpt-5.6-sol` • Thinking: `xhigh` • Cost: `$0.00`_ --- src/node/services/taskHandleStore.ts | 6 - .../services/workspaceTurnManager.test.ts | 87 +++++---- src/node/services/workspaceTurnManager.ts | 166 +++++------------- 3 files changed, 93 insertions(+), 166 deletions(-) diff --git a/src/node/services/taskHandleStore.ts b/src/node/services/taskHandleStore.ts index 55c951b1e4..1a3c115ae4 100644 --- a/src/node/services/taskHandleStore.ts +++ b/src/node/services/taskHandleStore.ts @@ -72,11 +72,6 @@ export interface WorkspaceTurnTaskHandleRecord { metadata: StreamEndEvent["metadata"]; }; deferredMessageIds?: string[]; - /** - * True only for a terminal result from an uncorrelated synthetic wake stream-end. - * A later correlated stream-end can replace this provisional result. - */ - provisionalOutcome?: boolean; error?: string; /** * How the owner workspace's stream-end treats this workspace turn while active. @@ -123,7 +118,6 @@ const WorkspaceTurnTaskHandleRecordSchema = z .passthrough() .optional(), deferredMessageIds: z.array(z.string().min(1)).optional(), - provisionalOutcome: z.boolean().optional(), error: z.string().optional(), attentionPolicy: BackgroundWorkAttentionPolicySchema.optional(), directParentResultDeliveryRequiredAt: z.string().optional(), diff --git a/src/node/services/workspaceTurnManager.test.ts b/src/node/services/workspaceTurnManager.test.ts index 205b716c0d..9c254d2607 100644 --- a/src/node/services/workspaceTurnManager.test.ts +++ b/src/node/services/workspaceTurnManager.test.ts @@ -3895,34 +3895,7 @@ describe("WorkspaceTurnManager", () => { expect(snapshot).toMatchObject({ status: "running", workspaceId: created.workspaceId }); }); - test("uncorrelated synthetic wake end stays active while a continuation is live", async () => { - let monitorWakePending = true; - const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest({ - hasPendingBashMonitorWakeContinuation: mock(() => monitorWakePending), - }); - const event = await appendUncorrelatedWakeHistory({ - historyService, - parentId, - taskId: created.taskId, - workspaceId: created.workspaceId, - inputSynthetic: true, - }); - - expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, event)).toBe(true); - expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ - status: "running", - deferredMessageIds: ["wake-output"], - }); - - monitorWakePending = false; - expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ - status: "completed", - messageId: "wake-output", - reportMarkdown: "Wake result", - }); - }); - - test("idle uncorrelated synthetic wake end settles provisionally from the wake output", async () => { + test("uncorrelated synthetic wake end leaves the active handle running", async () => { const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); const event = await appendUncorrelatedWakeHistory({ historyService, @@ -3934,10 +3907,7 @@ describe("WorkspaceTurnManager", () => { expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, event)).toBe(true); expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ - status: "completed", - messageId: "wake-output", - reportMarkdown: "Wake result", - provisionalOutcome: true, + status: "running", }); const correlatedFinal: StreamEndEvent = { @@ -3950,13 +3920,60 @@ describe("WorkspaceTurnManager", () => { parts: [{ type: "text", text: "Real final result" }], }; expect(await finalizeWorkspaceTurnStreamEndForTest(taskService, correlatedFinal)).toBe(true); - const corrected = await workspaceTurnSnapshot(taskService, parentId, created.taskId); - expect(corrected).toMatchObject({ + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ status: "completed", messageId: "real-final", reportMarkdown: "Real final result", }); - expect(corrected?.provisionalOutcome).toBeUndefined(); + }); + + test("compaction-preserved turn anchor ignores an uncorrelated synthetic wake end", async () => { + const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); + const compactionSummary = createMuxMessage("compaction-summary", "user", "Compacted context", { + muxMetadata: { + type: "compaction-summary", + pendingFollowUp: { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + text: "Continue the delegated work", + workspaceTurnMetadata: workspaceTurnMuxMetadata(parentId, created.taskId), + }, + }, + }); + expect( + (await historyService.appendToHistory(created.workspaceId, compactionSummary)).success + ).toBe(true); + const wakeInput = createMuxMessage("wake-input", "user", "Continue after a wake", { + synthetic: true, + }); + expect((await historyService.appendToHistory(created.workspaceId, wakeInput)).success).toBe( + true + ); + const wakeOutput = createMuxMessage("wake-output", "assistant", "Wake result", { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + finishReason: "stop", + }); + expect((await historyService.appendToHistory(created.workspaceId, wakeOutput)).success).toBe( + true + ); + + expect( + await finalizeWorkspaceTurnStreamEndForTest(taskService, { + type: "stream-end", + workspaceId: created.workspaceId, + messageId: wakeOutput.id, + metadata: { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + finishReason: "stop", + }, + parts: [{ type: "text", text: "Wake result" }], + }) + ).toBe(true); + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ + status: "running", + }); }); test("manual user input still supersedes an active workspace turn on uncorrelated end", async () => { diff --git a/src/node/services/workspaceTurnManager.ts b/src/node/services/workspaceTurnManager.ts index 0d5eef8f12..1911718b41 100644 --- a/src/node/services/workspaceTurnManager.ts +++ b/src/node/services/workspaceTurnManager.ts @@ -2152,12 +2152,11 @@ export class WorkspaceTurnManager { // error / stale restart interrupt — never an explicit user interrupt) may be // corrected once by an explicitly allowed resettle, but only when the new settlement // actually changes the outcome (duplicate stream-end replays must stay idempotent). - const currentIsProvisional = current.provisionalOutcome === true; const resettleStaleTerminal = params.allowTerminalResettle === true && this.isTerminalWorkspaceTurnStatus(current.status) && - (current.status !== "completed" || currentIsProvisional) && - (currentIsProvisional || isSelfHealEligibleSettledWorkspaceTurn(current)) && + current.status !== "completed" && + isSelfHealEligibleSettledWorkspaceTurn(current) && (params.next.status !== current.status || params.next.messageId !== current.messageId); if (this.isTerminalWorkspaceTurnStatus(current.status) && !resettleStaleTerminal) { const active = this.activeWorkspaceTurnHandleByWorkspaceId.get(params.record.workspaceId); @@ -2237,7 +2236,6 @@ export class WorkspaceTurnManager { nextStatus: nextRecord.status, }); delete nextRecord.terminalAttentionNotifiedAt; - delete nextRecord.provisionalOutcome; } if (suppressedQuietSettlement) { // Downgrade-compatible suppression marker (upgrade↔downgrade rule): @@ -3811,9 +3809,7 @@ export class WorkspaceTurnManager { const active = this.activeWorkspaceTurnHandleByWorkspaceId.get(record.workspaceId); const hasRuntimeActivity = this.aiService.isStreaming(record.workspaceId) || - this.workspaceService.hasPendingQueuedOrPreparingTurn(record.workspaceId) || - this.workspaceService.hasPendingAutoRetry(record.workspaceId) || - this.workspaceService.hasPendingBashMonitorWakeContinuation(record.workspaceId); + this.workspaceService.hasPendingQueuedOrPreparingTurn(record.workspaceId); if (hasRuntimeActivity) { return true; } @@ -4074,26 +4070,6 @@ export class WorkspaceTurnManager { }; } - /** Rebuild a deferred uncorrelated wake end after its continuation evidence clears. */ - private buildDeferredWakeEndEventFromHistory( - record: WorkspaceTurnTaskHandleRecord, - message: MuxMessage - ): StreamEndEvent | null { - if (message.role !== "assistant" || message.metadata?.partial === true) { - return null; - } - return { - type: "stream-end", - workspaceId: record.workspaceId, - messageId: message.id, - metadata: { - ...message.metadata, - model: coerceNonEmptyString(message.metadata?.model) ?? record.modelString ?? defaultModel, - }, - parts: message.parts as StreamEndEvent["parts"], - }; - } - private buildTerminalWorkspaceTurnRecordFromEvent( record: WorkspaceTurnTaskHandleRecord, event: StreamEndEvent, @@ -4102,7 +4078,6 @@ export class WorkspaceTurnManager { const baseRecord = { ...record }; delete baseRecord.error; delete baseRecord.deferredMessageIds; - delete baseRecord.provisionalOutcome; // A "tool-calls" finish on a delegated turn backed by queue-dispatch // evidence is a queue cut: some other queued input (a manual user message, // /compact, the owner's own follow-up turn) dispatched at the tool boundary @@ -4189,13 +4164,10 @@ export class WorkspaceTurnManager { const allowDeferredMessages = !(await this.hasActiveWorkspaceTurnDeferredBlockers(record)); for (const message of historyResult.data.toReversed()) { - const isDeferred = this.isDeferredWorkspaceTurnMessage(record, message.id); - if (isDeferred && !allowDeferredMessages) { + if (this.isDeferredWorkspaceTurnMessage(record, message.id) && !allowDeferredMessages) { continue; } - const event = - this.buildWorkspaceTurnStreamEndEventFromHistory(record, message) ?? - (isDeferred ? this.buildDeferredWakeEndEventFromHistory(record, message) : null); + const event = this.buildWorkspaceTurnStreamEndEventFromHistory(record, message); if (event != null) { // History order alone cannot prove a queue cut (a later unrelated user // message is not causal evidence), so stale recovery conservatively @@ -4250,6 +4222,29 @@ export class WorkspaceTurnManager { }; } + private isWorkspaceTurnAnchorForRecord( + record: WorkspaceTurnTaskHandleRecord, + message: MuxMessage + ): boolean { + const muxMetadata = message.metadata?.muxMetadata; + if (muxMetadata?.type === "workspace-turn-task") { + return ( + muxMetadata.taskHandleId === record.handleId && + muxMetadata.ownerWorkspaceId === record.ownerWorkspaceId && + muxMetadata.turnId === record.turnId + ); + } + if (muxMetadata?.type === "compaction-summary") { + const preserved = muxMetadata.pendingFollowUp?.workspaceTurnMetadata; + return ( + preserved?.taskHandleId === record.handleId && + preserved.ownerWorkspaceId === record.ownerWorkspaceId && + preserved.turnId === record.turnId + ); + } + return false; + } + private async interruptWorkspaceTurnFromUncorrelatedStreamEnd( event: StreamEndEvent ): Promise { @@ -4293,22 +4288,21 @@ export class WorkspaceTurnManager { } let streamEndIndex = -1; - let promptIndex = -1; + let turnAnchorIndex = -1; for (const [index, message] of historyResult.data.entries()) { if (message.id === event.messageId) { streamEndIndex = index; } - const metadata = this.getWorkspaceTurnMetadataFromValue(message.metadata?.muxMetadata); - if ( - metadata?.taskHandleId === record.handleId && - metadata.ownerWorkspaceId === record.ownerWorkspaceId && - metadata.turnId === record.turnId - ) { - promptIndex = index; + if (this.isWorkspaceTurnAnchorForRecord(record, message)) { + turnAnchorIndex = index; } } - if (streamEndIndex !== -1 && promptIndex !== -1 && streamEndIndex < promptIndex) { + if (streamEndIndex === -1 || turnAnchorIndex === -1) { + await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); + return true; + } + if (streamEndIndex < turnAnchorIndex) { log.debug("Ignoring stale uncorrelated stream-end before queued workspace turn prompt", { workspaceId: event.workspaceId, taskHandleId: record.handleId, @@ -4317,92 +4311,15 @@ export class WorkspaceTurnManager { return true; } - if (promptIndex !== -1) { - const scanEnd = streamEndIndex === -1 ? historyResult.data.length : streamEndIndex; - const hasManualSupersessionInput = historyResult.data - .slice(promptIndex + 1, scanEnd) - .some(isManualChildWorkspaceInput); - if (hasManualSupersessionInput) { - await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); - return true; - } - - if (await this.hasLiveUncorrelatedWakeContinuationEvidence(event, record)) { - await this.persistDeferredUncorrelatedWakeEnd(record, event); - return true; - } - - const next = this.buildTerminalWorkspaceTurnRecordFromEvent(record, event, { - supersedeEvidence: null, - }); - next.provisionalOutcome = true; - await this.settleWorkspaceTurn({ - record, - next, - waiterSettlement: - next.status === "completed" - ? { status: "completed", result: this.buildWorkspaceTurnWaitResult(next) } - : { status: "error", error: new Error(next.error ?? "Workspace turn failed") }, - allowTerminalResettle: true, - }); - return true; + const hasManualSupersessionInput = historyResult.data + .slice(turnAnchorIndex + 1, streamEndIndex) + .some(isManualChildWorkspaceInput); + if (hasManualSupersessionInput) { + await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); } - - await this.settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd(record, event); return true; } - private async hasLiveUncorrelatedWakeContinuationEvidence( - event: StreamEndEvent, - record: WorkspaceTurnTaskHandleRecord - ): Promise { - if (await this.hasActiveWorkspaceTurnDeferredBlockers(record)) { - return true; - } - if ( - this.workspaceService.hasPendingQueuedOrPreparingTurn(event.workspaceId) || - this.workspaceService.hasPendingAutoRetry(event.workspaceId) || - this.workspaceService.hasPendingBashMonitorWakeContinuation(event.workspaceId) - ) { - return true; - } - if ( - this.workspaceService.hasPendingWorkspaceTurnContinuation(event.workspaceId, { - type: "workspace-turn-task", - taskHandleId: record.handleId, - ownerWorkspaceId: record.ownerWorkspaceId, - turnId: record.turnId, - }) - ) { - return true; - } - const activeStream = this.streamManager?.getStreamInfo(event.workspaceId); - return activeStream != null && activeStream.messageId !== event.messageId; - } - - private async persistDeferredUncorrelatedWakeEnd( - record: WorkspaceTurnTaskHandleRecord, - event: StreamEndEvent - ): Promise { - await this.workspaceTurnSettlementLocks.withLock(record.handleId, async () => { - const current = await this.taskHandleStore.getWorkspaceTurn( - record.ownerWorkspaceId, - record.handleId - ); - if (current == null || !isActiveWorkspaceTurnTaskStatus(current.status)) { - return; - } - if (this.isDeferredWorkspaceTurnMessage(current, event.messageId)) { - return; - } - await this.taskHandleStore.upsertWorkspaceTurn({ - ...current, - updatedAt: getIsoNow(), - deferredMessageIds: [...(current.deferredMessageIds ?? []), event.messageId], - }); - }); - } - private async settleWorkspaceTurnSupersededFromUncorrelatedStreamEnd( record: WorkspaceTurnTaskHandleRecord, event: StreamEndEvent @@ -4415,7 +4332,6 @@ export class WorkspaceTurnManager { messageId: event.messageId, error, }; - delete next.provisionalOutcome; delete next.deferredMessageIds; await this.settleWorkspaceTurn({ record, From ae252668a1c0b9b2716825a68dd973894670dac3 Mon Sep 17 00:00:00 2001 From: Mux Date: Tue, 1 Sep 2026 15:45:28 -0500 Subject: [PATCH 3/4] =?UTF-8?q?=F0=9F=A4=96=20fix:=20treat=20user=20compac?= =?UTF-8?q?tion=20as=20supersession?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../services/workspaceTurnManager.test.ts | 52 +++++++++++++++++++ src/node/services/workspaceTurnManager.ts | 13 ++++- 2 files changed, 64 insertions(+), 1 deletion(-) diff --git a/src/node/services/workspaceTurnManager.test.ts b/src/node/services/workspaceTurnManager.test.ts index 9c254d2607..31781d8cb6 100644 --- a/src/node/services/workspaceTurnManager.test.ts +++ b/src/node/services/workspaceTurnManager.test.ts @@ -3994,6 +3994,58 @@ describe("WorkspaceTurnManager", () => { }); }); + test("user-triggered auto-compaction still supersedes an active workspace turn", async () => { + const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); + const prompt = createMuxMessage("turn-prompt", "user", "Summarize the repo", { + muxMetadata: workspaceTurnMuxMetadata(parentId, created.taskId), + }); + expect((await historyService.appendToHistory(created.workspaceId, prompt)).success).toBe(true); + const compactionRequest = createMuxMessage( + "auto-compaction", + "user", + "Compacting before a new user prompt", + { + synthetic: true, + muxMetadata: { + type: "compaction-request", + rawCommand: "/compact", + parsed: {}, + source: "auto-compaction", + }, + } + ); + expect( + (await historyService.appendToHistory(created.workspaceId, compactionRequest)).success + ).toBe(true); + const wakeOutput = createMuxMessage("wake-output", "assistant", "Wake result", { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + finishReason: "stop", + }); + expect((await historyService.appendToHistory(created.workspaceId, wakeOutput)).success).toBe( + true + ); + + expect( + await finalizeWorkspaceTurnStreamEndForTest(taskService, { + type: "stream-end", + workspaceId: created.workspaceId, + messageId: wakeOutput.id, + metadata: { + model: "anthropic:claude-opus-4-6", + agentId: "exec", + finishReason: "stop", + }, + parts: [{ type: "text", text: "Wake result" }], + }) + ).toBe(true); + expect(await workspaceTurnSnapshot(taskService, parentId, created.taskId)).toMatchObject({ + status: "interrupted", + messageId: wakeOutput.id, + error: "Workspace turn superseded by an uncorrelated workspace stream-end", + }); + }); + test("getWorkspaceTurnSnapshot recovers stale completed handles from matching history", async () => { const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); const appendResult = await historyService.appendToHistory( diff --git a/src/node/services/workspaceTurnManager.ts b/src/node/services/workspaceTurnManager.ts index 1911718b41..00971f09ff 100644 --- a/src/node/services/workspaceTurnManager.ts +++ b/src/node/services/workspaceTurnManager.ts @@ -283,7 +283,18 @@ const WORKSPACE_TURN_SUPERSEDED_BY_NEW_INPUT_ERROR = /** A human-authored child input that redirects the delegated turn. */ function isManualChildWorkspaceInput(message: MuxMessage): boolean { - return message.role === "user" && message.metadata?.synthetic !== true; + if (message.role !== "user") { + return false; + } + if (message.metadata?.synthetic !== true) { + return true; + } + const muxMetadata = message.metadata.muxMetadata; + return ( + muxMetadata?.type === "compaction-request" && + muxMetadata.source === "auto-compaction" && + muxMetadata.parsed.followUpContent?.dispatchOptions?.source !== "internal-resume" + ); } /** From f9baa2fc9420fd2ce62164ea74380ef075696184 Mon Sep 17 00:00:00 2001 From: Mux Date: Tue, 1 Sep 2026 16:31:45 -0500 Subject: [PATCH 4/4] =?UTF-8?q?=F0=9F=A4=96=20fix:=20guard=20malformed=20c?= =?UTF-8?q?ompaction=20metadata?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Use the defensive compaction follow-up accessor when classifying manual child input. This keeps malformed persisted compaction rows from leaving workspace turns live. --- _Generated with `xum` • Model: `openai:gpt-5.6-sol` • Thinking: `xhigh` • Cost: `$0.06`_ --- src/node/services/workspaceTurnManager.test.ts | 8 ++++---- src/node/services/workspaceTurnManager.ts | 3 ++- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/src/node/services/workspaceTurnManager.test.ts b/src/node/services/workspaceTurnManager.test.ts index 31781d8cb6..0ec48e2e06 100644 --- a/src/node/services/workspaceTurnManager.test.ts +++ b/src/node/services/workspaceTurnManager.test.ts @@ -18,7 +18,7 @@ import { Ok, Err, type Result } from "@/common/types/result"; import { DEFAULT_TASK_SETTINGS } from "@/common/types/tasks"; import type { SendMessageError } from "@/common/types/errors"; import type { ErrorEvent, StreamEndEvent } from "@/common/types/stream"; -import { createMuxMessage } from "@/common/types/message"; +import { createMuxMessage, type MuxMessageMetadata } from "@/common/types/message"; import type { WorkspaceMetadata } from "@/common/types/workspace"; import type { AIService } from "@/node/services/aiService"; import type { @@ -3994,7 +3994,7 @@ describe("WorkspaceTurnManager", () => { }); }); - test("user-triggered auto-compaction still supersedes an active workspace turn", async () => { + test("malformed user-triggered auto-compaction still supersedes an active workspace turn", async () => { const { parentId, taskService, historyService, created } = await startWorkspaceTurnForTest(); const prompt = createMuxMessage("turn-prompt", "user", "Summarize the repo", { muxMetadata: workspaceTurnMuxMetadata(parentId, created.taskId), @@ -4009,9 +4009,9 @@ describe("WorkspaceTurnManager", () => { muxMetadata: { type: "compaction-request", rawCommand: "/compact", - parsed: {}, + parsed: null, source: "auto-compaction", - }, + } as unknown as MuxMessageMetadata, } ); expect( diff --git a/src/node/services/workspaceTurnManager.ts b/src/node/services/workspaceTurnManager.ts index 00971f09ff..0b4f7918c0 100644 --- a/src/node/services/workspaceTurnManager.ts +++ b/src/node/services/workspaceTurnManager.ts @@ -52,6 +52,7 @@ import { } from "@/common/types/backgroundWorkAttention"; import { createMuxMessage, + getCompactionFollowUpContent, parseWorkspaceTurnTaskCorrelation, type MuxMessage, type MuxMessageMetadata, @@ -293,7 +294,7 @@ function isManualChildWorkspaceInput(message: MuxMessage): boolean { return ( muxMetadata?.type === "compaction-request" && muxMetadata.source === "auto-compaction" && - muxMetadata.parsed.followUpContent?.dispatchOptions?.source !== "internal-resume" + getCompactionFollowUpContent(muxMetadata)?.dispatchOptions?.source !== "internal-resume" ); }