diff --git a/apps/mobile/src/features/threads/thread-work-log.tsx b/apps/mobile/src/features/threads/thread-work-log.tsx index 32e6d272b678..75a28348d5ff 100644 --- a/apps/mobile/src/features/threads/thread-work-log.tsx +++ b/apps/mobile/src/features/threads/thread-work-log.tsx @@ -53,6 +53,7 @@ import { scopeThreadRef } from "@t3tools/client-runtime/environment"; import { environmentThreadDetails, threadEnvironment } from "../../state/threads"; import { useAtomCommand } from "../../state/use-atom-command"; import { toolActivityFaviconUrl } from "@t3tools/shared/favicon"; +import { runInterruptSenderThreadId } from "@t3tools/shared/orchestrationV2Timeline"; import { AppText as Text } from "../../components/AppText"; import { T3Wordmark } from "../../components/T3Wordmark"; @@ -933,7 +934,9 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow( row.projectedItem.item.type === "notification" ? notificationChildThreadId(row.projectedItem.item.source) : undefined; - const canExpand = row.canExpand && notifiedSubagentThreadId === undefined; + const interruptSenderThreadId = runInterruptSenderThreadId(row.projectedItem.item); + const linkedThreadId = notifiedSubagentThreadId ?? interruptSenderThreadId; + const canExpand = row.canExpand && linkedThreadId === undefined; const reasoning = row.projectedItem.item.type === "reasoning" ? row.projectedItem.item : null; const fetchedItem = fetchedDetail.data?.item ?? null; // Reads keep their path list; the fetched file contents show as output. @@ -1000,26 +1003,26 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow( {...(isFreshRow(row.createdAt) ? { entering: FadeIn.duration(200) } : {})} > { - if (notifiedSubagentThreadId !== undefined) { + if (linkedThreadId !== undefined) { // Push, not navigate: navigate reuses this Thread route, so back // would skip the parent thread and land on Home (matches #15068). navigation.dispatch( StackActions.push("Thread", { environmentId: String(props.environmentId), - threadId: String(notifiedSubagentThreadId), + threadId: String(linkedThreadId), }), ); return; diff --git a/apps/server/src/mcp/OrchestratorMcpService.test.ts b/apps/server/src/mcp/OrchestratorMcpService.test.ts index 5b426477c240..196812fc76a7 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.test.ts @@ -499,10 +499,12 @@ describe("OrchestratorMcpService", () => { clientRequestId: "cancel-dispose-failed-task", }); assert.equal(result.status, "cancel_requested"); + const commands = yield* Ref.get(dispatched); assert.deepEqual( - (yield* Ref.get(dispatched)).map((command) => (command as { type: string }).type), + commands.map((command) => (command as { type: string }).type), ["thread.stop", "delegated_task.completion-delivery.dispose"], ); + assert.deepInclude(commands[0], { createdBy: "agent", senderThreadId: parentThreadId }); }).pipe(Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies)))); }), ); diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 872163b660d1..ac90349c5528 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -2054,12 +2054,17 @@ const make = Effect.gen(function* () { const stopChild = Effect.gen(function* () { const commandId = stableCommandId({ scope, requestKey: key, operation: "cancel-task" }); const reason = input.reason; + const attribution = { + createdBy: "agent", + senderThreadId: scope.thread.threadId, + } as const; yield* threadManagement .dispatch({ type: "thread.stop", commandId, threadId: current.childThreadId, ...(reason === undefined ? {} : { reason }), + ...attribution, }) .pipe( Effect.mapError((error) => @@ -2071,7 +2076,12 @@ const make = Effect.gen(function* () { ); // A retry with the same clientRequestId repeats only the stops that failed. yield* threadManagement - .stopDelegatedTasks({ threadId: current.childThreadId, commandId, reason }) + .stopDelegatedTasks({ + threadId: current.childThreadId, + commandId, + reason, + ...attribution, + }) .pipe( Effect.mapError((error) => failure( @@ -2518,6 +2528,8 @@ const make = Effect.gen(function* () { threadId: input.threadId, ...(input.runId === undefined ? {} : { runId: input.runId }), ...(input.reason === undefined ? {} : { reason: input.reason }), + createdBy: "agent", + ...(parent === undefined ? {} : { senderThreadId: parent.thread.id }), }) .pipe( Effect.mapError((error) => diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index e7e82eaddffb..21a1c9540104 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -2055,6 +2055,15 @@ describe("orchestrator MCP toolkit", () => { cancelledStatusCall.structuredContent, ).pipe(Effect.orDie); expect(cancelledStatus.status).toBe("interrupted"); + expect( + (yield* orchestrator.getThreadProjection(cancellable.childThreadId)).turnItems.find( + (item) => item.type === "run_interrupt_result", + ), + ).toMatchObject({ + createdBy: "agent", + senderThreadId: parentThreadId, + message: "Run interrupted by an agent", + }); // Explicit task_cancel disposes automatic delivery after the child // interrupt succeeds. The interrupted result remains readable, @@ -2571,6 +2580,15 @@ describe("orchestrator MCP toolkit", () => { interruptedWaitCall.structuredContent, ).pipe(Effect.orDie); expect(interruptedWait.status).toBe("interrupted"); + expect( + (yield* orchestrator.getThreadProjection(activeThread.threadId)).turnItems.find( + (item) => item.type === "run_interrupt_result" && item.runId === activeRun.id, + ), + ).toMatchObject({ + createdBy: "agent", + senderThreadId: parentThreadId, + message: "Run interrupted by an agent before provider start", + }); const repeatedInterruptCall = yield* invoke("t3_thread_interrupt", { threadId: activeThread.threadId, runId: activeRun.id, diff --git a/apps/server/src/orchestration-v2/BackgroundWorkStop.integration.test.ts b/apps/server/src/orchestration-v2/BackgroundWorkStop.integration.test.ts index 6d5e0b92d87f..1f51de81864c 100644 --- a/apps/server/src/orchestration-v2/BackgroundWorkStop.integration.test.ts +++ b/apps/server/src/orchestration-v2/BackgroundWorkStop.integration.test.ts @@ -351,6 +351,7 @@ const stopEarlierBackgroundWork = ({ type: "thread.stop", commandId: CommandId.make("stop-stalled-run"), threadId, + createdBy: "agent", }, ); yield* worker.drain(); @@ -381,10 +382,16 @@ const stopEarlierBackgroundWork = ({ after.turnItems.find((candidate) => candidate.type === "assistant_message")?.status, interrupted ? "interrupted" : "running", ); - assert.equal( - after.turnItems.filter((candidate) => candidate.type === "run_interrupt_result").length, - interrupted ? 1 : 0, + const results = after.turnItems.filter( + (candidate) => candidate.type === "run_interrupt_result", ); + assert.equal(results.length, interrupted ? 1 : 0); + if (interrupted) { + assert.deepInclude(results[0], { + createdBy: "agent", + message: "Run interrupted by an agent", + }); + } assert.isEmpty(after.runs.filter((candidate) => candidate.status === "waiting")); return; } diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 3f17a6faa55d..81f7bdf24842 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -536,6 +536,19 @@ function hasLiveRun(projection: Pick): ); } +/** A stop with no `createdBy` came from a client's Stop, so the user. */ +function runInterruptAttribution( + command: Extract< + OrchestrationV2ServerCommand, + { readonly type: "run.interrupt" | "thread.stop" } + >, +) { + return { + createdBy: command.createdBy ?? "user", + ...(command.senderThreadId === undefined ? {} : { senderThreadId: command.senderThreadId }), + } as const; +} + /** The link with its watch replaced, or removed when `watch` is undefined. */ function withPullRequestWatch( link: ThreadPullRequestLink, @@ -8530,6 +8543,20 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); } if (providerThread !== undefined) { + const requestId = idAllocator.derive.runSignalTurnItem({ + runId: run.id, + signal: "interrupt-request", + }); + // A dead session's Stop writes its request in this same command, before the store has it. + const pendingRequest = (yield* Ref.get(input.events)).findLast( + (event) => event.type === "turn-item.updated" && event.payload.id === requestId, + ); + const request = + pendingRequest?.type === "turn-item.updated" + ? pendingRequest.payload + : yield* projectionStore + .getTurnItem({ threadId: run.threadId, itemId: requestId }) + .pipe(Effect.catchCause(() => Effect.succeed(null))); yield* emitEvent({ ...base, type: "turn-item.updated", @@ -8538,6 +8565,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio run, rootNode, providerThread, + request: request?.type === "run_interrupt_request" ? request : null, completedAt: input.now, }), }); @@ -8638,6 +8666,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio completedAt: now, updatedAt: now, type: "run_interrupt_request", + ...runInterruptAttribution(command), message: command.reason ?? "Interrupt requested", }, }); @@ -8909,6 +8938,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); const emitEvent = emit(events, command); + const interruptAttribution = runInterruptAttribution(command); const interruptRequestItem: OrchestrationV2TurnItem = { id: idAllocator.derive.runSignalTurnItem({ runId: run.id, @@ -8928,6 +8958,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio completedAt: now, updatedAt: now, type: "run_interrupt_request", + ...interruptAttribution, message: command.reason ?? "Interrupt requested", }; @@ -9022,7 +9053,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio completedAt: now, updatedAt: now, type: "run_interrupt_result", - message: "Run interrupted before provider start", + ...interruptAttribution, + message: + interruptAttribution.createdBy === "agent" + ? "Run interrupted by an agent before provider start" + : interruptAttribution.createdBy === "user" + ? "Run interrupted by user before provider start" + : "Run interrupted before provider start", }; yield* emitEvent({ type: "turn-item.updated", @@ -9252,6 +9289,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio runId: target.id, holdQueue: true, ...(command.reason === undefined ? {} : { reason: command.reason }), + ...runInterruptAttribution(command), }, interruptEvents, interruptEffects, diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index e2c905b1e17a..7fee380f2eea 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -175,6 +175,19 @@ export const layer: Layer.Layer< }), ) .pipe(Effect.catchCause(() => Effect.succeed(false))), + loadRunInterruptRequest: () => + projectionStore + .getTurnItem({ + threadId: input.threadId, + itemId: idAllocator.derive.runSignalTurnItem({ + runId: input.runId, + signal: "interrupt-request", + }), + }) + .pipe( + Effect.map((item) => (item?.type === "run_interrupt_request" ? item : null)), + Effect.catchCause(() => Effect.succeed(null)), + ), }; }; @@ -1257,6 +1270,7 @@ export const layer: Layer.Layer< shouldStartProviderTurn: runControls.shouldStartProviderTurn, shouldFinalizeRun: runControls.shouldFinalizeRun, hasUnpairedRunInterruptRequest: runControls.hasUnpairedRunInterruptRequest, + loadRunInterruptRequest: runControls.loadRunInterruptRequest, message: { messageId: message.id, text: userText, diff --git a/apps/server/src/orchestration-v2/RunExecutionService.test.ts b/apps/server/src/orchestration-v2/RunExecutionService.test.ts index 708205663787..a077ce6aba2d 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.test.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.test.ts @@ -3352,22 +3352,51 @@ it.effect("does not overwrite Stop when ownership changes after the finalization }), ); -it.effect("emits run_interrupt_result when hard-stop finalizes the active attempt", () => - Effect.gen(function* () { - const { written, observed, committedEffects } = yield* captureRootRunTermination({ - key: "hard-stop", - shouldFinalizeRun: () => Effect.succeed(true), - }); - assert.deepEqual( - written.map((item) => item.type), - ["run_interrupt_result"], - ); - assert.deepEqual(observed, ["run:interrupted", "pull-requests-refreshed"]); - assert.deepEqual( - committedEffects.map((effect) => effect.request.type), - ["checkpoint.capture"], - ); - }), +const interruptingThreadId = ThreadId.make("thread:hard-stop:interrupting"); + +it.effect.each([ + { + name: "an agent's request", + request: { createdBy: "agent", senderThreadId: interruptingThreadId }, + message: "Run interrupted by an agent", + }, + { name: "the user's Stop", request: { createdBy: "user" }, message: "Run interrupted by user" }, + { name: "no request", request: null, message: "" }, +] as const)( + "emits run_interrupt_result attributed to $name when hard-stop finalizes the active attempt", + ({ name, request, message }) => + Effect.gen(function* () { + const { written, observed, committedEffects } = yield* captureRootRunTermination({ + key: `hard-stop:${name}`, + shouldFinalizeRun: () => Effect.succeed(true), + interruptRequest: + request === null + ? null + : ({ + type: "run_interrupt_request", + ...request, + } as RunExecutionService.RunInterruptRequestTurnItem), + }); + assert.deepEqual( + written.map((item) => item.type), + ["run_interrupt_result"], + ); + const result = written[0]; + if (result?.type !== "run_interrupt_result") { + return assert.fail("expected run_interrupt_result"); + } + assert.equal(result.message, message); + assert.equal(result.createdBy, request?.createdBy); + assert.equal( + result.senderThreadId, + request !== null && "senderThreadId" in request ? request.senderThreadId : undefined, + ); + assert.deepEqual(observed, ["run:interrupted", "pull-requests-refreshed"]); + assert.deepEqual( + committedEffects.map((effect) => effect.request.type), + ["checkpoint.capture"], + ); + }), ); it.effect.each(["completed", "interrupted", "cancelled", "failed"] as const)( @@ -3497,6 +3526,7 @@ function captureRootRunTermination(input: { readonly shouldFinalizeRun: () => Effect.Effect; readonly rejectTerminalWrite?: boolean; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; + readonly interruptRequest?: RunExecutionService.RunInterruptRequestTurnItem | null; readonly seedOpenSubagent?: boolean; readonly events?: ( ids: BackgroundScenarioIds, @@ -3653,6 +3683,9 @@ function captureRootRunTermination(input: { : { hasUnpairedRunInterruptRequest: input.hasUnpairedRunInterruptRequest, }), + ...(input.interruptRequest === undefined + ? {} + : { loadRunInterruptRequest: () => Effect.succeed(input.interruptRequest ?? null) }), message: { messageId: MessageId.make(`message:${input.key}`), text: "interrupt projection", diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index a316ce0ccd62..5fb847b7bc51 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -78,6 +78,11 @@ type ProviderTerminalEvent = Extract< { readonly type: "turn.terminal" } >; +export type RunInterruptRequestTurnItem = Extract< + OrchestrationV2TurnItem, + { readonly type: "run_interrupt_request" } +>; + function isTerminalProviderTurnStatus(status: OrchestrationV2ProviderTurn["status"]): boolean { return ( status === "completed" || @@ -518,6 +523,7 @@ export interface RunExecutionServiceV2StartRootRunInput { readonly shouldStartProviderTurn?: () => Effect.Effect; readonly shouldFinalizeRun?: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; + readonly loadRunInterruptRequest?: () => Effect.Effect; readonly message: ProviderAdapter.ProviderAdapterV2TurnMessage; readonly modelSelection: ModelSelection; readonly runtimePolicy: ProviderAdapter.ProviderAdapterV2RuntimePolicy; @@ -565,6 +571,10 @@ export const layer: Layer.Layer< readonly attempt: OrchestrationV2RunAttempt; readonly shouldFinalizeRun?: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; + readonly loadRunInterruptRequest?: () => Effect.Effect< + RunInterruptRequestTurnItem | null, + never + >; readonly openRunOwnedSubagents?: OpenRunOwnedSubagentProjection; readonly terminal: ProviderTerminalEvent; readonly failureItemPersisted: boolean; @@ -576,6 +586,7 @@ export const layer: Layer.Layer< }) => Effect.gen(function* () { const completedAt = yield* DateTime.now; + const loadInterruptRequest = input.loadRunInterruptRequest ?? (() => Effect.succeed(null)); const finalizedAttempt: OrchestrationV2RunAttempt | null = { ...input.attempt, status: input.terminal.status, @@ -609,6 +620,7 @@ export const layer: Layer.Layer< run: input.run, rootNode: input.rootNode, providerThread: input.providerThread, + request: yield* loadInterruptRequest(), completedAt, }), }, @@ -719,6 +731,7 @@ export const layer: Layer.Layer< run: input.run, rootNode: input.rootNode, providerThread: input.providerThread, + request: yield* loadInterruptRequest(), completedAt, }), }, @@ -988,6 +1001,9 @@ export const layer: Layer.Layer< : { hasUnpairedRunInterruptRequest: input.hasUnpairedRunInterruptRequest, }), + ...(input.loadRunInterruptRequest === undefined + ? {} + : { loadRunInterruptRequest: input.loadRunInterruptRequest }), openRunOwnedSubagents: openSubagents, terminal, failureItemPersisted: terminal.status === "failed", @@ -1469,11 +1485,18 @@ export const layer: Layer.Layer< }), ); +// Empty unless a user or agent asked for the stop, so clients show no attribution. +function interruptResultMessage(request: RunInterruptRequestTurnItem | null): string { + if (request === null || request.createdBy === "system") return ""; + return request.createdBy === "agent" ? "Run interrupted by an agent" : "Run interrupted by user"; +} + export function makeInterruptResultTurnItem(input: { readonly idAllocator: IdAllocator.IdAllocatorV2["Service"]; readonly run: OrchestrationV2Run; readonly rootNode: OrchestrationV2ExecutionNode; readonly providerThread: OrchestrationV2ProviderThread; + readonly request: RunInterruptRequestTurnItem | null; readonly completedAt: DateTime.Utc; }): OrchestrationV2TurnItem { return { @@ -1498,6 +1521,10 @@ export function makeInterruptResultTurnItem(input: { completedAt: input.completedAt, updatedAt: input.completedAt, type: "run_interrupt_result", - message: "Run interrupted by user", + ...(input.request?.createdBy === undefined ? {} : { createdBy: input.request.createdBy }), + ...(input.request?.senderThreadId === undefined + ? {} + : { senderThreadId: input.request.senderThreadId }), + message: interruptResultMessage(input.request), }; } diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.test.ts b/apps/server/src/orchestration-v2/ThreadManagementService.test.ts index 94e17bdc567e..954b970fde72 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.test.ts @@ -25,7 +25,7 @@ import * as LegacyV1ThreadImporter from "./legacy/LegacyV1ThreadImporter.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ThreadManagementService from "./ThreadManagementService.ts"; -it("stamps authoritative provenance on commands that create threads or messages", () => { +it("stamps authoritative provenance on commands that record who authored them", () => { const command: OrchestrationV2Command = { type: "thread.create", createdBy: "agent", @@ -53,15 +53,31 @@ it("stamps authoritative provenance on commands that create threads or messages" createdBy: "user", creationSource: "web", }); -}); - -it("leaves commands that do not create durable authored content unchanged", () => { - const command: OrchestrationV2Command = { + const interrupt: OrchestrationV2Command = { type: "run.interrupt", commandId: CommandId.make("command:thread-management:interrupt"), threadId: ThreadId.make("thread:thread-management:interrupt"), runId: RunId.make("run:thread-management:interrupt"), }; + expect( + ThreadManagementService.withCreationProvenance( + { + ...interrupt, + createdBy: "agent", + senderThreadId: ThreadId.make("thread:thread-management:spoofed-sender"), + }, + { createdBy: "user", creationSource: "web" }, + ), + ).toEqual({ ...interrupt, createdBy: "user" }); +}); + +it("leaves commands that do not create durable authored content unchanged", () => { + const command: OrchestrationV2Command = { + type: "prepared-run.retry", + commandId: CommandId.make("command:thread-management:retry"), + threadId: ThreadId.make("thread:thread-management:retry"), + runId: RunId.make("run:thread-management:retry"), + }; expect( ThreadManagementService.withCreationProvenance(command, { diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.ts b/apps/server/src/orchestration-v2/ThreadManagementService.ts index a12e6680ca7e..cafada2db6f1 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.ts @@ -60,6 +60,10 @@ export function withCreationProvenance( case "thread.merge_back": case "delegated_task.request": return { ...command, ...provenance }; + case "run.interrupt": { + const { senderThreadId: _senderThreadId, ...interrupt } = command; + return { ...interrupt, createdBy: provenance.createdBy }; + } default: return command; } @@ -150,6 +154,8 @@ export interface ThreadManagementInterruptInput { readonly threadId: ThreadId; readonly runId?: RunId; readonly reason?: string; + readonly createdBy?: OrchestrationV2Actor; + readonly senderThreadId?: ThreadId; } export type ThreadManagementInterruptResult = @@ -373,6 +379,8 @@ export interface ThreadManagementServiceShape { readonly threadId: ThreadId; readonly commandId: CommandId; readonly reason?: string | undefined; + readonly createdBy?: OrchestrationV2Actor | undefined; + readonly senderThreadId?: ThreadId | undefined; }) => Effect.Effect; readonly getThreadEventSequence: Orchestrator.OrchestratorV2["Service"]["getThreadEventSequence"]; readonly recoverDelegatedTask: Orchestrator.OrchestratorV2["Service"]["recoverDelegatedTask"]; @@ -838,6 +846,8 @@ const make = Effect.gen(function* () { threadId: input.threadId, runId: interruptibleRun.id, ...(input.reason === undefined ? {} : { reason: input.reason }), + ...(input.createdBy === undefined ? {} : { createdBy: input.createdBy }), + ...(input.senderThreadId === undefined ? {} : { senderThreadId: input.senderThreadId }), }); return { type: "interrupt_requested", run: interruptibleRun, dispatch } as const; }); @@ -854,6 +864,8 @@ const make = Effect.gen(function* () { commandId: CommandId.make(`${input.commandId}:stop:${threadId}`), threadId, ...(input.reason === undefined ? {} : { reason: input.reason }), + ...(input.createdBy === undefined ? {} : { createdBy: input.createdBy }), + ...(input.senderThreadId === undefined ? {} : { senderThreadId: input.senderThreadId }), }).pipe( Effect.andThen(stopDelegatedTasks({ ...input, threadId })), Effect.catch((error) => diff --git a/apps/server/src/orchestration-v2/ThreadStop.test.ts b/apps/server/src/orchestration-v2/ThreadStop.test.ts index 6c5874cb75f2..a6545e68628f 100644 --- a/apps/server/src/orchestration-v2/ThreadStop.test.ts +++ b/apps/server/src/orchestration-v2/ThreadStop.test.ts @@ -573,12 +573,21 @@ it.effect("thread.stop marks a turn it cannot interrupt so a late agent watch is }, }); + const senderThreadId = ThreadId.make("thread:stop-lost-session-sender"); yield* orchestrator.dispatch({ type: "thread.stop", commandId: CommandId.make("stop-lost-session"), threadId, + createdBy: "agent", + senderThreadId, }); assert.isTrue(Exit.isFailure(yield* Effect.exit(watch(threadId, 12)))); assert.deepEqual(yield* threadState(threadId), { runs: ["running"], watched: [] }); + assert.deepInclude( + (yield* orchestrator.getThreadProjection(threadId)).turnItems.find( + (item) => item.type === "run_interrupt_request", + ), + { createdBy: "agent", senderThreadId }, + ); }).pipe(Effect.provide(layerTest)), ); diff --git a/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt/codex_output.ts b/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt/codex_output.ts index a446c5abbba3..f0232cf5ceac 100644 --- a/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt/codex_output.ts +++ b/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt/codex_output.ts @@ -36,6 +36,10 @@ export function assertTurnInterruptOutput( assert.equal(interruptRequest.status, "completed"); assert.equal(interruptResult.status, "interrupted"); assert.equal(interruptResult.parentItemId, interruptRequest.id); + assert.equal( + interruptResult.type === "run_interrupt_result" ? interruptResult.message : undefined, + "Run interrupted by user", + ); assert.deepEqual( projection.attempts.map((attempt) => attempt.status), ["interrupted"], diff --git a/apps/web/src/components/chat/V2LifecycleRow.tsx b/apps/web/src/components/chat/V2LifecycleRow.tsx index 4d476e0db6dd..ceb3631b64d8 100644 --- a/apps/web/src/components/chat/V2LifecycleRow.tsx +++ b/apps/web/src/components/chat/V2LifecycleRow.tsx @@ -26,6 +26,7 @@ import { type ScopedThreadRef, } from "@t3tools/contracts"; import type { TimestampFormat } from "@t3tools/contracts/settings"; +import { runInterruptSenderThreadId } from "@t3tools/shared/orchestrationV2Timeline"; import { BotIcon, ChevronRightIcon, @@ -94,12 +95,19 @@ export function V2LifecycleRow(props: { ); } if (item.type === "run_interrupt_result") { + const senderThreadId = runInterruptSenderThreadId(item); return ( props.onOpenThread(senderThreadId), + })} /> ); } diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 7aa8c74fd80a..acf1140f4db2 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -93,6 +93,15 @@ const OrchestrationV2CreationFields = { creationSource: OrchestrationV2CreationSource, } as const; +/** + * Who asked to stop a run. A stop command or request without it is a client's Stop, so the user; + * a result without it had no request. + */ +const OrchestrationV2RunInterruptAttributionFields = { + createdBy: Schema.optionalKey(OrchestrationV2Actor), + senderThreadId: Schema.optionalKey(ThreadId), +} as const; + export const OrchestrationV2NativeRefStrength = Schema.Literals(["strong", "weak", "none"]); export type OrchestrationV2NativeRefStrength = typeof OrchestrationV2NativeRefStrength.Type; @@ -1482,11 +1491,13 @@ export const OrchestrationV2TurnItem = Schema.Union([ }), Schema.Struct({ ...OrchestrationV2TurnItemBaseFields, + ...OrchestrationV2RunInterruptAttributionFields, type: Schema.Literal("run_interrupt_request"), message: Schema.String, }), Schema.Struct({ ...OrchestrationV2TurnItemBaseFields, + ...OrchestrationV2RunInterruptAttributionFields, type: Schema.Literal("run_interrupt_result"), message: Schema.String, }), @@ -2261,11 +2272,13 @@ export const OrchestrationV2TurnItemJson = Schema.Union([ }), Schema.Struct({ ...OrchestrationV2TurnItemJsonBaseFields, + ...OrchestrationV2RunInterruptAttributionFields, type: Schema.Literal("run_interrupt_request"), message: Schema.String, }), Schema.Struct({ ...OrchestrationV2TurnItemJsonBaseFields, + ...OrchestrationV2RunInterruptAttributionFields, type: Schema.Literal("run_interrupt_result"), message: Schema.String, }), @@ -2921,6 +2934,7 @@ export const OrchestrationV2Command = Schema.Union([ threadId: ThreadId, runId: RunId, reason: Schema.optional(Schema.String), + ...OrchestrationV2RunInterruptAttributionFields, /** * Set by the Stop button. Stop also holds the queue, ends the thread's pull request * watches, and stops every delegated task under the thread. @@ -3119,6 +3133,7 @@ const OrchestrationV2InternalCommand = Schema.Union([ commandId: CommandId, threadId: ThreadId, reason: Schema.optional(Schema.String), + ...OrchestrationV2RunInterruptAttributionFields, }), /** * Records or updates a secret an agent asked the user for. Internal so no diff --git a/packages/shared/src/orchestrationV2Timeline.test.ts b/packages/shared/src/orchestrationV2Timeline.test.ts index 2c7f0d0c6049..053382a7677d 100644 --- a/packages/shared/src/orchestrationV2Timeline.test.ts +++ b/packages/shared/src/orchestrationV2Timeline.test.ts @@ -1,9 +1,10 @@ -import { NodeId, RunId } from "@t3tools/contracts"; +import { NodeId, RunId, ThreadId } from "@t3tools/contracts"; import { describe, expect, it } from "vite-plus/test"; import { createOrchestrationV2TurnItemVisibility, isOrchestrationV2TurnItemVisible, + runInterruptSenderThreadId, } from "./orchestrationV2Timeline.ts"; const runId = RunId.make("run:timeline-visibility"); @@ -118,3 +119,29 @@ describe.each([ ).toBe(true); }); }); + +describe("runInterruptSenderThreadId", () => { + const threadId = ThreadId.make("thread:timeline-interrupted"); + const senderThreadId = ThreadId.make("thread:timeline-interrupting"); + + it.each([ + { + name: "an agent in another thread", + createdBy: "agent", + sender: senderThreadId, + opens: senderThreadId, + }, + { name: "an agent in the same thread", createdBy: "agent", sender: threadId, opens: undefined }, + { name: "the user's Stop", createdBy: "user", sender: undefined, opens: undefined }, + { name: "no request", createdBy: undefined, sender: undefined, opens: undefined }, + ] as const)("opens the sender only for $name", ({ createdBy, sender, opens }) => { + expect( + runInterruptSenderThreadId({ + type: "run_interrupt_result", + threadId, + ...(createdBy === undefined ? {} : { createdBy }), + ...(sender === undefined ? {} : { senderThreadId: sender }), + }), + ).toBe(opens); + }); +}); diff --git a/packages/shared/src/orchestrationV2Timeline.ts b/packages/shared/src/orchestrationV2Timeline.ts index 215757d8f62b..e3bcca95bd5d 100644 --- a/packages/shared/src/orchestrationV2Timeline.ts +++ b/packages/shared/src/orchestrationV2Timeline.ts @@ -1,8 +1,10 @@ import type { + OrchestrationV2Actor, OrchestrationV2Run, OrchestrationV2RunAttempt, OrchestrationV2TurnItem, OrchestrationV2UserMessageInputIntent, + ThreadId, } from "@t3tools/contracts"; type TimelineRun = Pick; @@ -11,6 +13,19 @@ type TimelineTurnItem = Pick & { + readonly createdBy?: OrchestrationV2Actor | undefined; + readonly senderThreadId?: ThreadId | undefined; + }, +): ThreadId | undefined { + return item.type === "run_interrupt_result" && + item.createdBy === "agent" && + item.senderThreadId !== item.threadId + ? item.senderThreadId + : undefined; +} + export function isOrchestrationV2SupersededInterrupt(input: { readonly item: TimelineTurnItem; readonly attempts: ReadonlyArray;