diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index 5761689fcedb..3ae75b24178a 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -29,6 +29,10 @@ import * as ProjectService from "../project/ProjectService.ts"; import { ProviderAuthService } from "../provider/Services/ProviderAuthService.ts"; import * as ContextHandoffService from "./ContextHandoffService.ts"; import * as EventSink from "./EventSink.ts"; +import { + ProviderAdapterEventStreamError, + type ProviderAdapterV2SessionRuntime, +} from "./ProviderAdapter.ts"; import * as IdAllocator from "./IdAllocator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import * as ProviderSessionManager from "./ProviderSessionManager.ts"; @@ -158,6 +162,8 @@ function makeLocalCommandHarness(input: { readonly interruptOpen?: boolean; readonly interruptRunBeforeOpenFailure?: boolean; readonly writeFailure?: unknown; + readonly postOpenFailure?: "stream" | "interrupt"; + readonly interruptRunBeforePostOpenFailure?: boolean; }) { const now = DateTime.makeUnsafe("2026-09-04T12:00:00Z"); const threadId = ThreadId.make("thread-native-account-command"); @@ -195,7 +201,14 @@ function makeLocalCommandHarness(input: { providerSessionId, appThreadId: threadId, ownerNodeId: null, - nativeThreadRef: null, + nativeThreadRef: + input.postOpenFailure === undefined + ? null + : { + driver: ProviderDriverKind.make("pi"), + nativeId: "pi-session", + strength: "strong" as const, + }, nativeConversationHeadRef: null, status: "not_loaded", firstRunOrdinal: 2, @@ -330,33 +343,65 @@ function makeLocalCommandHarness(input: { updatedAt: now, }; const events: Array = []; + const postOpenStreamError = new ProviderAdapterEventStreamError({ + driver: ProviderDriverKind.make("pi"), + providerSessionId, + cause: new Error("Pi RPC read failed: pi process exited with code 1."), + }); const open = vi.fn(() => input.interruptOpen === true ? Effect.interrupt - : "openFailure" in input - ? Effect.sync(() => { - if (input.interruptRunBeforeOpenFailure === true) { - projection = { - ...projection, - runs: projection.runs.map((candidate) => - candidate.id === runId - ? { ...candidate, status: "interrupted", completedAt: now } - : candidate, + : input.postOpenFailure !== undefined + ? Effect.succeed({ + driver: ProviderDriverKind.make("pi"), + resumeThread: () => + Effect.sync(() => { + if (input.interruptRunBeforePostOpenFailure === true) { + projection = { + ...projection, + runs: projection.runs.map((candidate) => + candidate.id === runId + ? { ...candidate, status: "interrupted" as const, completedAt: now } + : candidate, + ), + }; + } + }).pipe( + Effect.andThen( + input.postOpenFailure === "interrupt" + ? Effect.interrupt + : Effect.fail(postOpenStreamError), ), - }; - } - }).pipe( - Effect.andThen( - Effect.fail( - new ProviderSessionManager.ProviderSessionOpenError({ - instanceId: newInstanceId, - providerSessionId, - cause: input.openFailure, - }), ), - ), - ) - : Effect.die("A local command must not open a native session."), + ensureThread: () => + input.postOpenFailure === "interrupt" + ? Effect.interrupt + : Effect.fail(postOpenStreamError), + } as unknown as ProviderAdapterV2SessionRuntime) + : "openFailure" in input + ? Effect.sync(() => { + if (input.interruptRunBeforeOpenFailure === true) { + projection = { + ...projection, + runs: projection.runs.map((candidate) => + candidate.id === runId + ? { ...candidate, status: "interrupted", completedAt: now } + : candidate, + ), + }; + } + }).pipe( + Effect.andThen( + Effect.fail( + new ProviderSessionManager.ProviderSessionOpenError({ + instanceId: newInstanceId, + providerSessionId, + cause: input.openFailure, + }), + ), + ), + ) + : Effect.die("A local command must not open a native session."), ); const startRootRun = vi.fn(() => Effect.die("A local command must not start a native turn.")); const tryHandlePromptCommand = vi.fn(() => @@ -560,6 +605,114 @@ effectIt.effect("does not overwrite a run interrupted while its provider session }), ); +effectIt.effect( + "terminalizes a starting run when resume and ensure fail after the session opens", + () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ + text: "Continue", + postOpenFailure: "stream", + }); + + yield* harness.start; + + expect(harness.open).toHaveBeenCalledOnce(); + expect(harness.startRootRun).not.toHaveBeenCalled(); + expect(harness.writeIfRunCurrent).toHaveBeenCalledWith( + expect.objectContaining({ + activeAttemptId: harness.attemptId, + expectedStatus: "starting", + }), + ); + const projection = harness.projection(); + expect(projection.runs.at(-1)).toMatchObject({ status: "failed", startedAt: null }); + expect(projection.attempts[0]).toMatchObject({ status: "failed", startedAt: null }); + expect(projection.nodes[0]).toMatchObject({ status: "failed", startedAt: null }); + expect(projection.turnItems).toMatchObject([ + { + type: "error", + title: "Provider turn failed to start", + status: "failed", + failure: { + class: "provider_error", + message: "Pi RPC read failed: pi process exited with code 1.", + }, + }, + ]); + }), +); + +effectIt.effect("leaves the run starting when a post-open failure will be retried", () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ + text: "Continue", + postOpenFailure: "stream", + }); + + const error = yield* harness.startWithRetry.pipe(Effect.flip); + + expect(error._tag).toBe("ProviderTurnStartError"); + expect(harness.writeIfRunCurrent).not.toHaveBeenCalled(); + expect(harness.startRootRun).not.toHaveBeenCalled(); + expect(harness.projection().runs.at(-1)?.status).toBe("starting"); + }), +); + +effectIt.effect("keeps a post-open failure retryable when terminal persistence fails", () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ + text: "Continue", + postOpenFailure: "stream", + writeFailure: new Error("database unavailable"), + }); + + const error = yield* harness.start.pipe(Effect.flip); + + expect(error._tag).toBe("ProviderTurnStartError"); + expect(harness.writeIfRunCurrent).toHaveBeenCalledOnce(); + expect(harness.startRootRun).not.toHaveBeenCalled(); + expect(harness.projection().runs.at(-1)?.status).toBe("starting"); + expect(harness.events).toEqual([]); + }), +); + +effectIt.effect("does not terminalize a post-open interruption", () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ text: "Continue", postOpenFailure: "interrupt" }); + + const exit = yield* Effect.exit(harness.start); + + expect(exit._tag).toBe("Failure"); + if (exit._tag === "Failure") { + expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true); + } + expect(harness.writeIfRunCurrent).not.toHaveBeenCalled(); + expect(harness.startRootRun).not.toHaveBeenCalled(); + expect(harness.projection().runs.at(-1)?.status).toBe("starting"); + expect(harness.events).toEqual([]); + }), +); + +effectIt.effect("does not overwrite a run interrupted while provider setup fails", () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ + text: "Continue", + postOpenFailure: "stream", + interruptRunBeforePostOpenFailure: true, + }); + + const error = yield* harness.start.pipe(Effect.flip); + + expect(error._tag).toBe("ProviderTurnStartError"); + expect(harness.writeIfRunCurrent).not.toHaveBeenCalled(); + expect(harness.startRootRun).not.toHaveBeenCalled(); + const projection = harness.projection(); + expect(projection.runs.at(-1)?.status).toBe("interrupted"); + expect(projection.turnItems).toEqual([]); + expect(harness.events).toEqual([]); + }), +); + effectIt.effect( "signs out the existing native provider before opening the newly selected provider", () => diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 2932b53fbae8..ca4df3f44565 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -61,11 +61,28 @@ export class ProviderTurnStartError extends Schema.TaggedError candidate.id === runId); @@ -564,6 +584,7 @@ export const layer: Layer.Layer< return; } const session = sessionResult.success; + gate.sessionOpened = true; let effectiveHandoffs = handoffs; const loadedProviderThread = yield* Effect.gen(function* () { if (nativeForkTransfer !== undefined) { @@ -1168,15 +1189,124 @@ export const layer: Layer.Layer< }); }); + // Resume/ensure/fork run after open, while the run is still `starting`. + // Retryable attempts leave it starting. The last attempt settles it failed + // so a dead provider child cannot pin the thread with no Stop and no settle. + const settlePostOpenStartFailure = (input: { + readonly threadId: ThreadId; + readonly runId: RunId; + readonly cause: E; + }) => + Effect.gen(function* () { + const projection = yield* projectionStore.getTurnStartContext(input.threadId, input.runId); + const run = projection.runs.find((candidate) => candidate.id === input.runId); + const rootNode = projection.nodes.find((candidate) => candidate.id === run?.rootNodeId); + const attempt = projection.attempts.find( + (candidate) => candidate.id === run?.activeAttemptId, + ); + const providerThread = projection.providerThreads.find( + (candidate) => candidate.id === run?.providerThreadId, + ); + if ( + run === undefined || + run.status !== "starting" || + rootNode === undefined || + attempt === undefined || + providerThread === undefined + ) { + return yield* Effect.fail(input.cause); + } + const now = yield* DateTime.now; + const failure = makeProviderFailure({ + cause: input.cause, + message: providerFailureDetail(input.cause), + class: "provider_error", + }); + const item: OrchestrationV2TurnItem = { + id: idAllocator.derive.runSignalTurnItem({ + runId: input.runId, + signal: "provider-turn-start-failure", + }), + threadId: projection.thread.id, + runId: input.runId, + nodeId: rootNode.id, + providerThreadId: providerThread.id, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: + Math.max( + 0, + ...projection.turnItems + .filter((candidate) => candidate.runId === input.runId) + .map((candidate) => candidate.ordinal), + ) + 1, + status: "failed", + startedAt: now, + completedAt: now, + updatedAt: now, + type: "error", + title: "Provider turn failed to start", + failure, + }; + const eventPayloads = [ + { type: "turn-item.updated" as const, payload: item }, + { + type: "run.updated" as const, + payload: { ...run, status: "failed" as const, completedAt: now }, + }, + { + type: "run-attempt.updated" as const, + payload: { ...attempt, status: "failed" as const, completedAt: now }, + }, + { + type: "node.updated" as const, + payload: { ...rootNode, status: "failed" as const, completedAt: now }, + }, + ]; + const events = yield* Effect.forEach(eventPayloads, (event) => + Effect.gen(function* () { + return { + ...event, + id: yield* idAllocator.allocate.event({ threadId: projection.thread.id }), + threadId: projection.thread.id, + runId: input.runId, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + } satisfies OrchestrationV2DomainEvent; + }), + ); + yield* eventSink.writeIfRunCurrent({ + threadId: projection.thread.id, + runId: input.runId, + activeAttemptId: attempt.id, + expectedStatus: "starting", + events, + }); + }); + return ProviderTurnStartServiceV2.of({ - start: (input) => - start(input).pipe( + start: (input) => { + const gate: ProviderTurnSessionGate = { sessionOpened: false }; + return start(input, gate).pipe( + Effect.catch((cause) => { + if (input.willRetry === true || !gate.sessionOpened) { + return Effect.fail(cause); + } + return settlePostOpenStartFailure({ + threadId: input.threadId, + runId: input.runId, + cause, + }); + }), Effect.mapError((cause) => isProviderTurnStartError(cause) ? cause : new ProviderTurnStartError({ runId: input.runId, cause }), ), - ), + ); + }, }); }), ); diff --git a/apps/web/src/components/ChatView.logic.test.ts b/apps/web/src/components/ChatView.logic.test.ts index 0f48fa1b9f06..eec0a9daf0bb 100644 --- a/apps/web/src/components/ChatView.logic.test.ts +++ b/apps/web/src/components/ChatView.logic.test.ts @@ -65,6 +65,8 @@ import { MAX_HIDDEN_MOUNTED_TERMINAL_THREADS, branchMismatchKey, buildExpiredTerminalContextToastCopy, + canInterruptThreadSession, + composerPrimaryActionIsStop, createLocalDispatchSnapshot, deriveCommittedServerUserMessageIds, deriveComposerSendState, @@ -2136,3 +2138,93 @@ describe("waitForRevertedMessage", () => { vi.useRealTimers(); }); }); + +describe("canInterruptThreadSession", () => { + const runtime = ( + status: NonNullable["status"], + activeRunId: RunId | null, + ): NonNullable => ({ + status, + activeRunId, + providerInstanceId: ProviderInstanceId.make("provider-pi"), + providerName: "pi", + lastError: null, + updatedAt: "2026-09-24T06:32:35.405Z", + }); + + it("offers Stop for a stuck starting or preparing run", () => { + expect(canInterruptThreadSession("connecting", runtime("starting", RunId.make("run-7")))).toBe( + true, + ); + expect( + canInterruptThreadSession("connecting", runtime("preparing", RunId.make("run-prep"))), + ).toBe(true); + }); + + it("does not offer Stop for a queued run in the connecting phase", () => { + expect(canInterruptThreadSession("connecting", runtime("queued", null))).toBe(false); + }); + + it("keeps Stop for a running turn and a waiting checkpoint", () => { + expect(canInterruptThreadSession("running", runtime("running", RunId.make("run-live")))).toBe( + true, + ); + expect(canInterruptThreadSession("running", runtime("waiting", null))).toBe(true); + }); + + it("hides Stop once the thread is idle", () => { + expect(canInterruptThreadSession("ready", runtime("idle", null))).toBe(false); + expect(canInterruptThreadSession("disconnected", null)).toBe(false); + }); +}); + +describe("composerPrimaryActionIsStop", () => { + it("replaces the connecting spinner with Stop before the provider turn starts", () => { + expect( + composerPrimaryActionIsStop({ + isRunning: false, + canInterrupt: true, + hasSendableContent: false, + isEditingQueuedMessage: false, + }), + ).toBe(true); + expect( + composerPrimaryActionIsStop({ + isRunning: false, + canInterrupt: true, + hasSendableContent: true, + isEditingQueuedMessage: false, + }), + ).toBe(true); + }); + + it("keeps Steer available once a running turn has a draft", () => { + expect( + composerPrimaryActionIsStop({ + isRunning: true, + canInterrupt: true, + hasSendableContent: true, + isEditingQueuedMessage: false, + }), + ).toBe(false); + expect( + composerPrimaryActionIsStop({ + isRunning: true, + canInterrupt: true, + hasSendableContent: false, + isEditingQueuedMessage: false, + }), + ).toBe(true); + }); + + it("does not replace an in-progress queued-message edit", () => { + expect( + composerPrimaryActionIsStop({ + isRunning: false, + canInterrupt: true, + hasSendableContent: false, + isEditingQueuedMessage: true, + }), + ).toBe(false); + }); +}); diff --git a/apps/web/src/components/ChatView.logic.ts b/apps/web/src/components/ChatView.logic.ts index 255a4f7e8cd3..3ef008be69d0 100644 --- a/apps/web/src/components/ChatView.logic.ts +++ b/apps/web/src/components/ChatView.logic.ts @@ -39,6 +39,7 @@ import { codexArtifactTemplateUsePrompt, type CodexArtifactTemplate, } from "@t3tools/client-runtime/codex-artifact-templates"; +import { threadRuntimeHasInterruptibleRun } from "@t3tools/client-runtime/state/thread-execution"; import { presentThreadShell } from "@t3tools/client-runtime/state/shell"; import { type ChatMessage, @@ -1232,6 +1233,34 @@ export function deriveCommittedServerUserMessageIds( ); } +/** + * Stop is offered for a running turn, and for a preparing or starting run. + * Those setup states stay in the connecting phase, but `run.interrupt` already + * accepts them. Queued runs are not interruptible. + */ +export function canInterruptThreadSession( + phase: SessionPhase, + runtime: Thread["runtime"] | null | undefined, +): boolean { + return phase === "running" || threadRuntimeHasInterruptibleRun(runtime); +} + +/** + * The primary composer control is Stop while a run can be interrupted and has + * not started provider work. A running turn keeps Steer once the draft has + * something to send. + */ +export function composerPrimaryActionIsStop(input: { + readonly isRunning: boolean; + readonly canInterrupt: boolean; + readonly hasSendableContent: boolean; + readonly isEditingQueuedMessage: boolean; +}): boolean { + if (input.isEditingQueuedMessage) return false; + if (input.isRunning) return !input.hasSendableContent; + return input.canInterrupt; +} + export function hasServerAcknowledgedLocalDispatch(input: { localDispatch: LocalDispatchSnapshot | null; phase: SessionPhase; diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 7ba2c3843e98..e76a9771c507 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -480,6 +480,7 @@ import { dismissBranchMismatchForSession, hasEnvironmentReconnectWarningGraceElapsed, scheduleEnvironmentReconnectWarning, + canInterruptThreadSession, hasServerAcknowledgedLocalDispatch, isBranchMismatchDismissedForSession, shouldShowBranchMismatchBanner, @@ -4314,7 +4315,8 @@ export default function ChatView(props: ChatViewProps) { const focusComposer = useCallback(() => { composerRef.current?.focusAtEnd(); }, [composerRef]); - const canInterruptRunningThread = activeThread !== undefined && phase === "running"; + const canInterruptRunningThread = + activeThread !== undefined && canInterruptThreadSession(phase, activeRuntime); const onInterrupt = useCallback(async () => { if (!activeThread) return; const result = await interruptThreadTurn({ diff --git a/apps/web/src/components/chat/ChatComposer.tsx b/apps/web/src/components/chat/ChatComposer.tsx index 893b0bb50b59..9248b6fbb238 100644 --- a/apps/web/src/components/chat/ChatComposer.tsx +++ b/apps/web/src/components/chat/ChatComposer.tsx @@ -87,6 +87,7 @@ import { deriveComposerSendState, getAntigravitySendBlockReason, readFileAsDataUrl, + canInterruptThreadSession, resolveComposerInteractionMode, resolveComposerProviderSelection, threadShellHasStarted, @@ -1354,6 +1355,7 @@ const ComposerFooterPrimaryActions = memo(function ComposerFooterPrimaryActions( isComplete: boolean; } | null; isRunning: boolean; + canInterrupt: boolean; followUpBehavior: "queue" | "steer"; alternateShortcutLabel: string | null; showPlanFollowUpPrompt: boolean; @@ -1392,6 +1394,7 @@ const ComposerFooterPrimaryActions = memo(function ComposerFooterPrimaryActions( compact={props.compact} pendingAction={props.pendingAction} isRunning={props.isRunning} + canInterrupt={props.canInterrupt} followUpBehavior={props.followUpBehavior} alternateShortcutLabel={props.alternateShortcutLabel} showPlanFollowUpPrompt={props.showPlanFollowUpPrompt} @@ -7433,6 +7436,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) activeThreadModelDisplayName={activeThreadModelDisplayName} pendingAction={pendingPrimaryAction} isRunning={phase === "running"} + canInterrupt={canInterruptThreadSession(phase, activeThread?.runtime)} followUpBehavior={settings.followUpBehavior} alternateShortcutLabel={shortcutLabelForCommand( keybindings, diff --git a/apps/web/src/components/chat/ComposerPrimaryActions.tsx b/apps/web/src/components/chat/ComposerPrimaryActions.tsx index 990520befd0b..d53021b02bda 100644 --- a/apps/web/src/components/chat/ComposerPrimaryActions.tsx +++ b/apps/web/src/components/chat/ComposerPrimaryActions.tsx @@ -20,6 +20,7 @@ import { alternateComposerDispatchAction, resolveComposerDispatchMode, } from "@t3tools/client-runtime/state/composer-dispatch"; +import { composerPrimaryActionIsStop } from "../ChatView.logic"; interface PendingActionState { questionIndex: number; @@ -33,6 +34,11 @@ interface ComposerPrimaryActionsProps { compact: boolean; pendingAction: PendingActionState | null; isRunning: boolean; + /** + * Preparing and starting runs stay in the connecting phase, but Stop must + * still be offered. Defaults to `isRunning` for callers that only know that. + */ + canInterrupt?: boolean; followUpBehavior?: "queue" | "steer"; alternateShortcutLabel?: string | null; showPlanFollowUpPrompt: boolean; @@ -84,6 +90,7 @@ export const ComposerPrimaryActions = memo(function ComposerPrimaryActions({ compact, pendingAction, isRunning, + canInterrupt, followUpBehavior = "steer", alternateShortcutLabel = null, showPlanFollowUpPrompt, @@ -108,6 +115,7 @@ export const ComposerPrimaryActions = memo(function ComposerPrimaryActions({ : undefined; const environmentIdentificationMode = useEnvironmentIdentificationMode(); const shortcutModifiers = useShortcutModifierState(); + const stopAvailable = canInterrupt ?? isRunning; const isQueuing = !isEditingQueuedMessage && resolveComposerDispatchMode({ @@ -148,7 +156,7 @@ export const ComposerPrimaryActions = memo(function ComposerPrimaryActions({ if (pendingAction) { return (
- {isRunning ? renderStopGenerationButton(true) : null} + {stopAvailable ? renderStopGenerationButton(true) : null} {pendingAction.questionIndex > 0 ? ( compact ? (