From e25a405f8c0c0fe090e7cbc7acf6a071cf931644 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Sat, 3 Oct 2026 17:32:08 +0300 Subject: [PATCH 01/11] feat(web): stop active subagents from Lineage --- .../Orchestrator.control-reads.test.ts | 152 ++++++++++++++++++ .../src/orchestration-v2/Orchestrator.ts | 69 ++++++++ .../src/orchestration-v2/ProjectionStore.ts | 4 + .../testkit/OrchestratorScenario.ts | 1 + apps/server/src/relay/AgentAwarenessRelay.ts | 1 + ...ThreadRelationshipsControl.agents.test.tsx | 80 ++++++++- .../chat/ThreadRelationshipsControl.tsx | 78 ++++++++- .../client-runtime/src/operations/commands.ts | 16 ++ .../state/orchestrationV2Projection.test.ts | 13 ++ .../src/state/orchestrationV2Projection.ts | 2 + .../src/state/threadCommands.ts | 8 + packages/contracts/src/orchestrationV2.ts | 16 ++ 12 files changed, 436 insertions(+), 4 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index ad41ae6ffde6..e4dfbf5b49ac 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -23,6 +23,7 @@ import * as Layer from "effect/Layer"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; +import * as EffectOutbox from "./EffectOutbox.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; @@ -42,6 +43,7 @@ const database = SqlitePersistenceMemory; const testLayer = Layer.mergeAll( database, ProjectionStore.layer.pipe(Layer.provide(database)), + EffectOutbox.layer.pipe(Layer.provide(database)), makeOrchestratorV2ReplayLayerWithRegistry( { name: "control-reads" }, ProviderAdapterRegistry.makeLayer([adapter]), @@ -49,6 +51,156 @@ const testLayer = Layer.mergeAll( ), ); +it.effect("interrupts only the selected running native Codex subagent", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const outbox = yield* EffectOutbox.EffectOutboxV2; + const parentThreadId = ThreadId.make("parent:stop-subagent"); + const childThreadId = ThreadId.make("child:stop-subagent"); + const subagentId = NodeId.make("subagent:stop-subagent"); + const providerThreadId = ProviderThreadId.make("provider-thread:stop-subagent"); + const providerTurnId = ProviderTurnId.make("provider-turn:stop-subagent"); + const providerSessionId = ProviderSessionId.make("session:stop-subagent"); + const now = yield* DateTime.now; + for (const threadId of [parentThreadId, childThreadId]) { + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make(`create:${threadId}`), + threadId, + projectId: ProjectId.make("project:stop-subagent"), + title: "Agent", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + } + const subagent = { + id: subagentId, + threadId: parentThreadId, + runId: null, + parentNodeId: NodeId.make("root:stop-subagent"), + origin: "provider_native" as const, + createdBy: "agent" as const, + driver: adapter.driver, + providerInstanceId: instanceId, + providerThreadId: null, + childThreadId, + nativeTaskRef: null, + prompt: "Check this", + title: "Worker", + model: "gpt-5.1-codex", + status: "running" as const, + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }; + yield* projections.apply({ + id: EventId.make("subagent:stop-subagent"), + type: "subagent.updated", + threadId: parentThreadId, + occurredAt: now, + payload: subagent, + }); + yield* projections.apply({ + id: EventId.make("provider-thread:stop-subagent"), + type: "provider-thread.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + id: providerThreadId, + driver: adapter.driver, + providerInstanceId: instanceId, + providerSessionId, + appThreadId: childThreadId, + ownerNodeId: subagentId, + nativeThreadRef: null, + nativeConversationHeadRef: null, + status: "active", + firstRunOrdinal: null, + lastRunOrdinal: null, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }); + yield* projections.apply({ + id: EventId.make("provider-turn:stop-subagent"), + type: "provider-turn.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + id: providerTurnId, + providerThreadId, + nodeId: NodeId.make("root:child-stop-subagent"), + runAttemptId: null, + nativeTurnRef: null, + ordinal: 1, + status: "running", + startedAt: now, + completedAt: null, + }, + }); + const commandId = CommandId.make("interrupt:stop-subagent"); + const accepted = yield* orchestrator.dispatch({ + type: "subagent.interrupt", + commandId, + threadId: parentThreadId, + subagentId, + }); + assert.deepEqual( + accepted.storedEvents.map((event) => event.event.type), + ["subagent.interrupt-requested"], + ); + assert.equal( + (yield* projections.getThreadProjection(parentThreadId)).subagents[0]?.status, + "running", + ); + assert.deepEqual( + (yield* outbox.listByCommandId(commandId)).map(({ threadId, request }) => ({ + threadId, + request, + })), + [ + { + threadId: childThreadId, + request: { + type: "provider-turn.interrupt", + providerSessionId, + providerThreadId, + providerTurnId, + }, + }, + ], + ); + yield* projections.apply({ + id: EventId.make("subagent:stop-subagent:completed"), + type: "subagent.updated", + threadId: parentThreadId, + occurredAt: now, + payload: { ...subagent, status: "completed", completedAt: now }, + }); + const settledCommandId = CommandId.make("interrupt:stop-subagent:settled"); + const rejected = yield* orchestrator + .dispatch({ + type: "subagent.interrupt", + commandId: settledCommandId, + threadId: parentThreadId, + subagentId, + }) + .pipe(Effect.flip); + assert.instanceOf(rejected, Orchestrator.OrchestratorDispatchError); + assert.equal(rejected.cause, undefined); + assert.deepEqual(yield* outbox.listByCommandId(settledCommandId), []); + }).pipe(Effect.provide(testLayer)), +); + it.effect( "dispatches metadata, queue resume and request controls without hydrating unrelated history", () => diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 2707fe5b8344..16ed35030e8c 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -387,6 +387,7 @@ function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { case "prepared-run.progress": case "prepared-run.fail": case "run.interrupt": + case "subagent.interrupt": case "queued-message.promote-to-steer": case "queue.resume": case "queued-run.reorder": @@ -8204,6 +8205,71 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return undefined; }); + const dispatchSubagentInterrupt = ( + command: Extract, + events: Ref.Ref>, + effects: Ref.Ref>, + ) => + Effect.gen(function* () { + const parent = yield* loadProjectionForCommand(command, ["subagents"]); + const subagent = parent.subagents.find((entry) => entry.id === command.subagentId); + if ( + subagent?.origin !== "provider_native" || + subagent.driver !== "codex" || + subagent.status !== "running" || + subagent.childThreadId === null + ) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + }); + } + const childThreadId = subagent.childThreadId; + const child = yield* projectionStore + .getThreadRecords(childThreadId, ["providerThreads", "providerTurns"]) + .pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: childThreadId, cause }), + ), + ); + const turn = child.providerTurns.findLast((entry) => entry.status === "running"); + const providerThread = child.providerThreads.find( + (entry) => entry.id === turn?.providerThreadId, + ); + if (turn === undefined || providerThread?.providerSessionId == null) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + }); + } + const providerSessionId = providerThread.providerSessionId; + const emitEvent = emit(events, command); + yield* emitEvent({ + type: "subagent.interrupt-requested", + threadId: command.threadId, + ...(subagent.runId === null ? {} : { runId: subagent.runId }), + nodeId: subagent.id, + driver: subagent.driver, + providerInstanceId: subagent.providerInstanceId, + occurredAt: yield* DateTime.now, + payload: subagent.id, + }); + yield* Ref.update(effects, (existing) => [ + ...existing, + { + id: `effect:${command.commandId}:provider-turn.interrupt:${turn.id}`, + commandId: command.commandId, + threadId: childThreadId, + request: { + type: "provider-turn.interrupt", + providerSessionId, + providerThreadId: providerThread.id, + providerTurnId: turn.id, + }, + } satisfies PendingOrchestrationEffectV2, + ]); + }); + const dispatchCheckpointRollback = ( command: Extract, events: Ref.Ref>, @@ -9332,6 +9398,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "run.interrupt": cancelUnsettledEffects = yield* dispatchRunInterrupt(command, events, effects); break; + case "subagent.interrupt": + yield* dispatchSubagentInterrupt(command, events, effects); + break; case "queued-message.promote-to-steer": yield* dispatchQueuedMessagePromoteToSteer(command, events, effects); break; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index ac0624b05806..8f42cc499ce1 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -638,6 +638,8 @@ export function applyToProjection( }; switch (event.type) { + case "subagent.interrupt-requested": + return projection; case "thread.created": case "thread.archived": case "thread.unarchived": @@ -1663,6 +1665,8 @@ export const layer: Layer.Layer = const apply: ProjectionStoreV2Shape["apply"] = (event) => Effect.gen(function* () { switch (event.type) { + case "subagent.interrupt-requested": + break; case "thread.created": case "thread.archived": case "thread.unarchived": diff --git a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts index b9cad107409a..552e02653c35 100644 --- a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts +++ b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts @@ -166,6 +166,7 @@ function commandThreadIds(command: OrchestrationV2Command): ReadonlyArray ({ projects: [] as unknown[], configs: new Map(), showTooltips: false, + command: vi.fn().mockResolvedValue({ _tag: "Success" }), })); vi.mock("@tanstack/react-router", () => ({ useNavigate: () => state.navigate })); @@ -23,7 +24,7 @@ vi.mock("../../state/entities", () => ({ vi.mock("../../lib/archivedThreadsState", () => ({ useArchivedThreadSnapshots: () => ({ snapshots: [] }), })); -vi.mock("../../state/use-atom-command", () => ({ useAtomCommand: () => vi.fn() })); +vi.mock("../../state/use-atom-command", () => ({ useAtomCommand: () => state.command })); vi.mock("../ui/tooltip", () => ({ Tooltip: ({ children }: { children: ReactNode }) => children, TooltipTrigger: ({ render, children }: { render: ReactElement; children: ReactNode }) => @@ -42,6 +43,83 @@ afterEach(async () => { state.projects = []; state.configs.clear(); state.showTooltips = false; + state.command.mockClear(); + state.projection = null; +}); + +it("stops only active subagents from lineage without opening their thread", async () => { + vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); + const parent = { + id: "parent", + lineage: { relationshipToParent: null }, + activeProviderThreadId: null, + }; + const child = { + id: "child", + title: "Worker", + lineage: { parentThreadId: "parent", relationshipToParent: "subagent" }, + }; + const agent = { + id: "agent", + childThreadId: "child", + origin: "app_owned", + driver: "codex", + providerInstanceId: "codex", + title: "Worker", + prompt: "Check the change", + model: "gpt-5.4", + status: "running", + progress: null, + result: null, + startedAt: DateTime.makeUnsafe("2026-09-16T12:00:00Z"), + completedAt: null, + updatedAt: DateTime.makeUnsafe("2026-09-16T12:00:00Z"), + }; + state.shells = [{ environmentId: "test", source: child }]; + const projection = { + thread: parent, + runs: [], + providerThreads: [], + providerSessions: [], + contextTransfers: [], + subagents: [agent], + }; + state.projection = projection; + const panel = ( + + ); + await act(async () => { + renderer = create(panel); + }); + const stopButton = () => renderer.root.findByProps({ "aria-label": "Stop subagent Worker" }); + await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenCalledWith({ + environmentId: "test", + input: { threadId: "child" }, + }); + expect(state.navigate).not.toHaveBeenCalled(); + + state.command.mockClear(); + state.projection = { ...projection, subagents: [{ ...agent, origin: "provider_native" }] }; + await act(async () => renderer.update(cloneElement(panel))); + await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenCalledWith({ + environmentId: "test", + input: { threadId: "parent", subagentId: "agent" }, + }); + + state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; + await act(async () => renderer.update(cloneElement(panel))); + expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); + state.projection = { + ...projection, + subagents: [{ ...agent, origin: "provider_native", driver: "claudeAgent" }], + }; + await act(async () => renderer.update(cloneElement(panel))); + expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); }); it("shows the matching child agent details and refreshes them when the agent settles", async () => { diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx index ac0235f85d1c..d941abd9dd4f 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx @@ -21,7 +21,12 @@ import { canDetachThreadProviderSession, resolveLatestMergeBackRun, } from "@t3tools/client-runtime/state/thread-workflows"; -import type { EnvironmentId, OrchestrationV2ThreadShell, ThreadId } from "@t3tools/contracts"; +import type { + EnvironmentId, + NodeId, + OrchestrationV2ThreadShell, + ThreadId, +} from "@t3tools/contracts"; import { groupBy } from "effect/Array"; import { useNavigate } from "@tanstack/react-router"; import { @@ -32,6 +37,7 @@ import { LoaderCircleIcon, MoreHorizontalIcon, PlusIcon, + SquareIcon, UnplugIcon, } from "lucide-react"; import { useMemo, useState, type ReactNode } from "react"; @@ -51,6 +57,7 @@ import { ThreadRelationshipIcon, threadRelationshipStatusLabel } from "./ThreadR import { Menu, MenuItem, MenuPopup, MenuTrigger } from "../ui/menu"; import { Tooltip, TooltipPopup, TooltipTrigger } from "../ui/tooltip"; +import { toastManager } from "../ui/toast"; import { THREAD_DETAILS_PANEL_LINK_SPLIT_GROUP_CLASS, THREAD_DETAILS_PANEL_ROW_CONTENT_CLASS, @@ -178,6 +185,8 @@ export function ThreadRelationshipsPanel(props: { subagent.childThreadId, { ...projectedSubagentsToRuntime([subagent])[0]!, + subagentId: subagent.id, + origin: subagent.origin, driver: subagent.driver, providerInstanceId: subagent.providerInstanceId, }, @@ -205,7 +214,10 @@ export function ThreadRelationshipsPanel(props: { const navigate = useNavigate(); const mergeBack = useAtomCommand(threadEnvironment.mergeBack); const stopSession = useAtomCommand(threadEnvironment.stopSession); + const interruptTurn = useAtomCommand(threadEnvironment.interruptTurn); + const interruptSubagent = useAtomCommand(threadEnvironment.interruptSubagent); const [busyAction, setBusyAction] = useState<"merge" | "detach" | null>(null); + const [stoppingSubagentId, setStoppingSubagentId] = useState(null); const latestMergeBackRun = projection === null ? null : resolveLatestMergeBackRun(projection); const mergeTargetThreadId = resolveMergeBackTargetThreadId(projection); const relationshipRows = useMemo( @@ -279,6 +291,29 @@ export function ThreadRelationshipsPanel(props: { setBusyAction(null); }; + const stopSubagent = async ( + subagentId: NodeId, + origin: "app_owned" | "provider_native", + childThreadId: ThreadId, + ) => { + if (stoppingSubagentId !== null) return; + setStoppingSubagentId(subagentId); + const result = + origin === "app_owned" + ? await interruptTurn({ + environmentId: props.environmentId, + input: { threadId: childThreadId }, + }) + : await interruptSubagent({ + environmentId: props.environmentId, + input: { threadId: props.threadId, subagentId }, + }); + setStoppingSubagentId(null); + if (result._tag === "Failure") { + toastManager.add({ type: "error", title: "Could not stop subagent" }); + } + }; + const parentTitle = mergeTargetThreadId === null ? null @@ -331,6 +366,13 @@ export function ThreadRelationshipsPanel(props: { : GitForkIcon; const relationship = relationshipLabel(edge, props.threadId); const agent = isSubagent && !isParent ? subagentsByThreadId.get(threadId) : undefined; + const canStop = + agent?.startedAt && + ((agent.origin === "app_owned" && + ["pending", "running", "waiting"].includes(agent.status)) || + (agent.origin === "provider_native" && + agent.driver === "codex" && + agent.status === "running")); const threadTitle = relationshipThreadTitle({ title: node?.thread?.title ?? agent?.title ?? threadId, isSubagent, @@ -379,7 +421,9 @@ export function ThreadRelationshipsPanel(props: { {agent ? ( agent.startedAt ? ( - + ) : null @@ -394,7 +438,7 @@ export function ThreadRelationshipsPanel(props: { ); return ( -
  • +
  • {isMergeTarget ? (
    @@ -473,6 +517,34 @@ export function ThreadRelationshipsPanel(props: { {relationshipTooltip} )} + {canStop && agent ? ( +
    + + + void stopSubagent(agent.subagentId, agent.origin, threadId) + } + /> + } + > + {stoppingSubagentId === agent.subagentId ? ( + + ) : ( + + )} + + Stop subagent + +
    + ) : null}
  • ); }) diff --git a/packages/client-runtime/src/operations/commands.ts b/packages/client-runtime/src/operations/commands.ts index c0cb8e67504c..33171320f18b 100644 --- a/packages/client-runtime/src/operations/commands.ts +++ b/packages/client-runtime/src/operations/commands.ts @@ -9,6 +9,7 @@ import { WS_METHODS, type ChatAttachment, type MessageId, + type NodeId, type ModelSelection, type OrchestrationV2Command, type OrchestrationV2CreationSource, @@ -186,6 +187,10 @@ export interface InterruptThreadTurnInput extends ThreadCommandInput { readonly turnId?: string; } +export interface InterruptSubagentInput extends ThreadCommandInput { + readonly subagentId: NodeId; +} + export interface RespondToThreadApprovalInput extends ThreadCommandInput { readonly requestId: RuntimeRequestId; readonly decision: ProviderApprovalDecision; @@ -805,6 +810,17 @@ export const interruptThreadTurn = Effect.fn("EnvironmentCommands.interruptThrea }); }); +export const interruptSubagent = Effect.fn("EnvironmentCommands.interruptSubagent")(function* ( + input: InterruptSubagentInput, +) { + return yield* dispatch({ + type: "subagent.interrupt", + commandId: yield* allocateCommandId(input), + threadId: input.threadId, + subagentId: input.subagentId, + }); +}); + export const respondToThreadApproval = Effect.fn("EnvironmentCommands.respondToThreadApproval")( function* (input: RespondToThreadApprovalInput) { return yield* dispatch({ diff --git a/packages/client-runtime/src/state/orchestrationV2Projection.test.ts b/packages/client-runtime/src/state/orchestrationV2Projection.test.ts index 1a53afb8d7e8..b6690c57c235 100644 --- a/packages/client-runtime/src/state/orchestrationV2Projection.test.ts +++ b/packages/client-runtime/src/state/orchestrationV2Projection.test.ts @@ -105,6 +105,19 @@ const emptyProjection = { } as OrchestrationV2ThreadProjection; describe("applyOrchestrationV2ProjectionEvent", () => { + it("leaves subagent state unchanged when a stop is requested", () => { + const event = { + id: "event-subagent-stop", + type: "subagent.interrupt-requested", + threadId, + nodeId: NodeId.make("subagent-stop"), + occurredAt: DateTime.makeUnsafe("2026-06-20T01:00:00.000Z"), + payload: NodeId.make("subagent-stop"), + } as OrchestrationV2DomainEvent; + + expect(applyOrchestrationV2ProjectionEvent(emptyProjection, event)).toBe(emptyProjection); + }); + it("keeps live token usage when the terminal provider turn omits it", () => { const providerTurnId = ProviderTurnId.make("provider-turn-reducer"); const running = { diff --git a/packages/client-runtime/src/state/orchestrationV2Projection.ts b/packages/client-runtime/src/state/orchestrationV2Projection.ts index 7e96a3039353..8e406901b207 100644 --- a/packages/client-runtime/src/state/orchestrationV2Projection.ts +++ b/packages/client-runtime/src/state/orchestrationV2Projection.ts @@ -161,6 +161,8 @@ export function applyOrchestrationV2ProjectionEvent( const latestLocalTurnOrdinal = options?.latestLocalTurnOrdinal; const base = { ...projection, updatedAt: event.occurredAt }; switch (event.type) { + case "subagent.interrupt-requested": + return projection; case "thread.created": case "thread.archived": case "thread.unarchived": diff --git a/packages/client-runtime/src/state/threadCommands.ts b/packages/client-runtime/src/state/threadCommands.ts index e3ad5c07f583..689b78ba43f0 100644 --- a/packages/client-runtime/src/state/threadCommands.ts +++ b/packages/client-runtime/src/state/threadCommands.ts @@ -25,6 +25,7 @@ import { type DeleteThreadInput, type EditQueuedRunInput, type InterruptThreadTurnInput, + type InterruptSubagentInput, type MarkThreadUnreadInput, type ForkThreadFromRunInput, type MergeThreadBackInput, @@ -59,6 +60,7 @@ import { deleteThread, editQueuedRun, interruptThreadTurn, + interruptSubagent, forkThreadFromRun, markThreadUnread, mergeThreadBack, @@ -284,6 +286,12 @@ export function createThreadEnvironmentAtoms( scheduler, concurrency, }), + interruptSubagent: createEnvironmentCommand(runtime, { + label: "environment-data:commands:thread:interrupt-subagent", + execute: (input: InterruptSubagentInput) => interruptSubagent(input), + scheduler, + concurrency, + }), respondToApproval: createEnvironmentCommand(runtime, { label: "environment-data:commands:thread:respond-to-approval", execute: (input: RespondToThreadApprovalInput) => respondToThreadApproval(input), diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 6004ab119917..5e668b73b8b6 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -1564,6 +1564,11 @@ export const OrchestrationV2DomainEvent = Schema.Union([ type: Schema.Literal("subagent.updated"), payload: OrchestrationV2Subagent, }), + Schema.Struct({ + ...OrchestrationV2EventBase.fields, + type: Schema.Literal("subagent.interrupt-requested"), + payload: NodeId, + }), Schema.Struct({ ...OrchestrationV2EventBase.fields, type: Schema.Literals(["provider-session.attached", "provider-session.updated"]), @@ -2355,6 +2360,11 @@ export const OrchestrationV2DomainEventJson = Schema.Union([ type: Schema.Literal("subagent.updated"), payload: OrchestrationV2SubagentJson, }), + Schema.Struct({ + ...OrchestrationV2JsonEventBaseFields, + type: Schema.Literal("subagent.interrupt-requested"), + payload: NodeId, + }), Schema.Struct({ ...OrchestrationV2JsonEventBaseFields, type: Schema.Literals(["provider-session.attached", "provider-session.updated"]), @@ -2733,6 +2743,12 @@ export const OrchestrationV2Command = Schema.Union([ reason: Schema.optional(Schema.String), holdQueue: Schema.optional(Schema.Boolean), }), + Schema.Struct({ + type: Schema.Literal("subagent.interrupt"), + commandId: CommandId, + threadId: ThreadId, + subagentId: NodeId, + }), Schema.Struct({ type: Schema.Literal("queued-message.promote-to-steer"), commandId: CommandId, From ac1c00dbe0db50b1d65aed9d9b339a724c782f54 Mon Sep 17 00:00:00 2001 From: Bil0000 <62337003+Bil0000@users.noreply.github.com> Date: Sun, 4 Oct 2026 04:28:48 +0200 Subject: [PATCH 02/11] fix(server): preserve parent activity on subagent stop requests --- .../orchestration-v2/ProjectionStore.test.ts | 17 ++++++++++++++++- .../src/orchestration-v2/ProjectionStore.ts | 1 + 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index 3e979544c87b..2e22d81bfe2f 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -1596,7 +1596,7 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { }), ); - it.effect("does not treat visited or marked-unread state as thread activity", () => + it.effect("does not treat read state or subagent interrupt requests as thread activity", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStore.ProjectionStoreV2; const createdAt = yield* DateTime.now; @@ -1662,6 +1662,21 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { const markedUnread = yield* projectionStore.getThreadProjection(threadId); assert.isNull(markedUnread.thread.lastVisitedAt); assert.deepEqual(markedUnread.thread.updatedAt, createdAt); + + yield* projectionStore.apply({ + id: EventId.make("event:projection-read-state:subagent-interrupt"), + type: "subagent.interrupt-requested", + threadId, + occurredAt: DateTime.add(createdAt, { seconds: 3 }), + payload: NodeId.make("subagent:projection-read-state"), + }); + + const interrupted = yield* projectionStore.getThreadProjection(threadId); + assert.deepEqual(interrupted.thread.updatedAt, createdAt); + const shell = (yield* projectionStore.getShellSnapshot()).threads.find( + (entry) => entry.id === threadId, + ); + assert.deepEqual(shell?.updatedAt, createdAt); }), ); diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 65e2daee0ba8..5efc4e7cb9cb 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -2507,6 +2507,7 @@ export const layer: Layer.Layer = } if ( + event.type !== "subagent.interrupt-requested" && event.type !== "thread.created" && event.type !== "thread.archived" && event.type !== "thread.unarchived" && From 015a3f6b9e242525a56bb314dc512f98c36f2075 Mon Sep 17 00:00:00 2001 From: Bil0000 <62337003+Bil0000@users.noreply.github.com> Date: Sun, 4 Oct 2026 04:28:58 +0200 Subject: [PATCH 03/11] fix(lineage): gate native subagent stops on server capability --- .../src/environment/ServerEnvironment.test.ts | 1 + .../src/environment/ServerEnvironment.ts | 1 + ...ThreadRelationshipsControl.agents.test.tsx | 15 +++++++- .../chat/ThreadRelationshipsControl.tsx | 7 ++-- .../src/operations/commands.test.ts | 35 +++++++++++++++++++ .../client-runtime/src/operations/commands.ts | 12 ++++++- packages/contracts/src/environment.test.ts | 12 +++++++ packages/contracts/src/environment.ts | 1 + 8 files changed, 80 insertions(+), 4 deletions(-) diff --git a/apps/server/src/environment/ServerEnvironment.test.ts b/apps/server/src/environment/ServerEnvironment.test.ts index f02e8e0ec7e3..dad7a0e05f60 100644 --- a/apps/server/src/environment/ServerEnvironment.test.ts +++ b/apps/server/src/environment/ServerEnvironment.test.ts @@ -169,6 +169,7 @@ it.layer(NodeServices.layer)("ServerEnvironmentLive", (it) => { expect(first.environmentId).toBe(second.environmentId); expect(first.orchestrationProtocolVersion).toBe(ORCHESTRATION_PROTOCOL_VERSION); expect(second.capabilities.repositoryIdentity).toBe(true); + expect(second.capabilities.subagentInterrupt).toBe(true); expect(second.capabilities.connectionProbe).toBe(true); expect(second.capabilities.attachmentUploads).toBe(true); expect(second.capabilities.fileAttachments).toEqual({ maxUploadBytes: 50 * 1024 * 1024 }); diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index 16a4a93b3f20..afa9494d8f11 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -216,6 +216,7 @@ export const make = Effect.gen(function* () { capabilities: { repositoryIdentity: true, connectionProbe: true, + subagentInterrupt: true, attachmentUploads: true, questionAttachments: true, fileAttachments: { maxUploadBytes: PROVIDER_SEND_TURN_MAX_FILE_BYTES }, diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx index 6ed6e19429ce..98d027dee8c1 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx @@ -104,7 +104,19 @@ it("stops only active subagents from lineage without opening their thread", asyn state.command.mockClear(); state.projection = { ...projection, subagents: [{ ...agent, origin: "provider_native" }] }; - await act(async () => renderer.update(cloneElement(panel))); + for (const supported of [undefined, false, true]) { + state.configs.set("test", { + environment: { + capabilities: supported === undefined ? {} : { subagentInterrupt: supported }, + }, + }); + await act(async () => renderer.update(cloneElement(panel))); + expect( + renderer.root.findAll( + (node) => node.type === "button" && node.props["aria-label"] === "Stop subagent Worker", + ), + ).toHaveLength(supported === true ? 1 : 0); + } await act(async () => stopButton().props.onClick()); expect(state.command).toHaveBeenCalledWith({ environmentId: "test", @@ -289,6 +301,7 @@ it("shows readable models and only differing workspace details in agent tooltips ]; state.shells = [{ environmentId: "test", source: child }]; state.configs.set("test", { + environment: { capabilities: {} }, providers: [ { instanceId: "codex", diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx index c58c8f3b00f3..608fa679ed55 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx @@ -203,7 +203,9 @@ export function ThreadRelationshipsPanel(props: { }) { const ref = scopeThreadRef(props.environmentId, props.threadId); const projection = useThreadProjection(ref)?.projection ?? null; - const providers = useServerConfigs().get(props.environmentId)?.providers; + const config = useServerConfigs().get(props.environmentId); + const providers = config?.providers; + const supportsSubagentInterrupt = config?.environment.capabilities.subagentInterrupt === true; const subagentsByThreadId = useMemo( () => new Map( @@ -403,7 +405,8 @@ export function ThreadRelationshipsPanel(props: { agent?.startedAt && ((agent.origin === "app_owned" && ["pending", "running", "waiting"].includes(agent.status)) || - (agent.origin === "provider_native" && + (supportsSubagentInterrupt && + agent.origin === "provider_native" && agent.driver === "codex" && agent.status === "running")); const threadTitle = relationshipThreadTitle({ diff --git a/packages/client-runtime/src/operations/commands.test.ts b/packages/client-runtime/src/operations/commands.test.ts index a9833f7499d7..3f8a675a41fd 100644 --- a/packages/client-runtime/src/operations/commands.test.ts +++ b/packages/client-runtime/src/operations/commands.test.ts @@ -25,6 +25,7 @@ import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Result from "effect/Result"; import * as SubscriptionRef from "effect/SubscriptionRef"; import { @@ -44,6 +45,7 @@ import { editQueuedRun, forkThreadFromRun, interruptThreadTurn, + interruptSubagent, mergeThreadBack, promoteQueuedRun, reorderActiveThread, @@ -78,6 +80,7 @@ const makeSupervisor = Effect.fn("TestEnvironmentCommands.makeSupervisor")(funct readonly projection?: OrchestrationV2ThreadProjection; readonly projectionRequests?: ThreadId[]; readonly advertiseServerResolvedCommandContext?: boolean; + readonly subagentInterrupt?: boolean; }) { const client = { [ORCHESTRATION_V2_WS_METHODS.dispatchCommand]: (command: OrchestrationV2Command) => @@ -125,6 +128,9 @@ const makeSupervisor = Effect.fn("TestEnvironmentCommands.makeSupervisor")(funct environment: { capabilities: { repositoryIdentity: true, + ...(input.subagentInterrupt === undefined + ? {} + : { subagentInterrupt: input.subagentInterrupt }), ...(input.advertiseServerResolvedCommandContext === false ? {} : { serverResolvedCommandContext: true }), @@ -148,6 +154,35 @@ const makeSupervisor = Effect.fn("TestEnvironmentCommands.makeSupervisor")(funct }); describe("V2 environment commands", () => { + it.effect.each([undefined, false, true])( + "dispatches a subagent interrupt only when support is advertised as %s", + (supported) => + Effect.gen(function* () { + const commands: OrchestrationV2Command[] = []; + const supervisor = yield* makeSupervisor({ + commands, + projects: [], + ...(supported === undefined ? {} : { subagentInterrupt: supported }), + }); + const commandId = CommandId.make("subagent-interrupt"); + const subagentId = NodeId.make("subagent-1"); + const result = yield* interruptSubagent({ + commandId, + threadId: v2ThreadId, + subagentId, + }).pipe( + Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), + Effect.result, + ); + expect(Result.isSuccess(result)).toBe(supported === true); + expect(commands).toEqual( + supported === true + ? [{ type: "subagent.interrupt", commandId, threadId: v2ThreadId, subagentId }] + : [], + ); + }).pipe(Effect.provide(TEST_CRYPTO_LAYER)), + ); + it.effect("routes projects through the event-sourced project transport", () => Effect.gen(function* () { const projects: ProjectMutation[] = []; diff --git a/packages/client-runtime/src/operations/commands.ts b/packages/client-runtime/src/operations/commands.ts index 33171320f18b..1207010c3ec4 100644 --- a/packages/client-runtime/src/operations/commands.ts +++ b/packages/client-runtime/src/operations/commands.ts @@ -6,6 +6,7 @@ import { CheckpointScopeId, ORCHESTRATION_V2_WS_METHODS, OrchestrationV2CheckpointUnavailableError, + OrchestrationV2DispatchCommandError, WS_METHODS, type ChatAttachment, type MessageId, @@ -813,9 +814,18 @@ export const interruptThreadTurn = Effect.fn("EnvironmentCommands.interruptThrea export const interruptSubagent = Effect.fn("EnvironmentCommands.interruptSubagent")(function* ( input: InterruptSubagentInput, ) { + const config = yield* getInitialServerConfig(); + const commandId = yield* allocateCommandId(input); + if (config.environment.capabilities.subagentInterrupt !== true) { + return yield* new OrchestrationV2DispatchCommandError({ + commandId, + commandType: "subagent.interrupt", + message: "This server does not support stopping native subagents.", + }); + } return yield* dispatch({ type: "subagent.interrupt", - commandId: yield* allocateCommandId(input), + commandId, threadId: input.threadId, subagentId: input.subagentId, }); diff --git a/packages/contracts/src/environment.test.ts b/packages/contracts/src/environment.test.ts index 3f624adf6fe2..242c1a348d64 100644 --- a/packages/contracts/src/environment.test.ts +++ b/packages/contracts/src/environment.test.ts @@ -14,6 +14,18 @@ const descriptor = { } as const; describe("ExecutionEnvironmentDescriptor", () => { + it.each([undefined, false, true])("preserves subagent interrupt support as %s", (supported) => { + expect( + decodeDescriptor({ + ...descriptor, + capabilities: { + ...descriptor.capabilities, + ...(supported === undefined ? {} : { subagentInterrupt: supported }), + }, + }).capabilities.subagentInterrupt, + ).toBe(supported); + }); + it("requires an advertised required-worktree bootstrap capability", () => { expect(decodeDescriptor(descriptor).capabilities.requiredWorktreeBootstrap).toBeUndefined(); expect( diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index 9a31f0888e67..dad3505bc0ff 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -91,6 +91,7 @@ export type ServerSelfUpdateCapability = typeof ServerSelfUpdateCapability.Type; export const ExecutionEnvironmentCapabilities = Schema.Struct({ repositoryIdentity: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), connectionProbe: Schema.optionalKey(Schema.Boolean), + subagentInterrupt: Schema.optionalKey(Schema.Boolean), /** Missing on older servers, which still accept inline image attachments. */ attachmentUploads: Schema.optionalKey(Schema.Boolean), /** Uploaded files may accompany question answers. */ From 9497d2f362632dd087809048781c02b44e4badc0 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Sun, 4 Oct 2026 17:01:59 +0300 Subject: [PATCH 04/11] fix(relay): publish native subagent stop requests --- apps/server/src/relay/AgentAwarenessRelay.test.ts | 9 +++++++++ apps/server/src/relay/AgentAwarenessRelay.ts | 2 +- 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index 4467fd9768ea..0e1a23fd343a 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -306,6 +306,15 @@ describe("AgentAwarenessRelay", () => { } }); + it("publishes subagent interrupts with a node ID payload", () => { + assert.isTrue( + AgentAwarenessRelay.shouldPublishAgentAwarenessEvent({ + type: "subagent.interrupt-requested", + payload: NodeId.make("native-subagent"), + }), + ); + }); + it("does not publish imported thread creation as new agent activity", () => { assert.isFalse( AgentAwarenessRelay.shouldPublishAgentAwarenessEvent({ diff --git a/apps/server/src/relay/AgentAwarenessRelay.ts b/apps/server/src/relay/AgentAwarenessRelay.ts index 957396a44ca2..fb0d868594cf 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.ts @@ -105,6 +105,7 @@ export function shouldPublishAgentAwarenessEvent( case "runtime-request.updated": case "subagent.updated": case "provider-thread.updated": + case "subagent.interrupt-requested": return true; case "thread.settled": case "thread.unsettled": @@ -123,7 +124,6 @@ export function shouldPublishAgentAwarenessEvent( case "run-attempt.created": case "run-attempt.updated": case "node.updated": - case "subagent.interrupt-requested": case "provider-session.attached": case "provider-session.updated": case "provider-session.detached": From 7d5217fbb9c60e5781d7c0b1c68c665b3505af1e Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Sun, 4 Oct 2026 17:02:39 +0300 Subject: [PATCH 05/11] fix(server): explain native subagent stop rejections --- .../Orchestrator.control-reads.test.ts | 101 +++++++++++++----- .../src/orchestration-v2/Orchestrator.ts | 2 + 2 files changed, 74 insertions(+), 29 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index e4dfbf5b49ac..53b80482d396 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -107,45 +107,47 @@ it.effect("interrupts only the selected running native Codex subagent", () => occurredAt: now, payload: subagent, }); + const providerThread = { + id: providerThreadId, + driver: adapter.driver, + providerInstanceId: instanceId, + providerSessionId, + appThreadId: childThreadId, + ownerNodeId: subagentId, + nativeThreadRef: null, + nativeConversationHeadRef: null, + status: "active", + firstRunOrdinal: null, + lastRunOrdinal: null, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + } as const; yield* projections.apply({ id: EventId.make("provider-thread:stop-subagent"), type: "provider-thread.updated", threadId: childThreadId, occurredAt: now, - payload: { - id: providerThreadId, - driver: adapter.driver, - providerInstanceId: instanceId, - providerSessionId, - appThreadId: childThreadId, - ownerNodeId: subagentId, - nativeThreadRef: null, - nativeConversationHeadRef: null, - status: "active", - firstRunOrdinal: null, - lastRunOrdinal: null, - handoffIds: [], - forkedFrom: null, - createdAt: now, - updatedAt: now, - }, + payload: providerThread, }); + const providerTurn = { + id: providerTurnId, + providerThreadId, + nodeId: NodeId.make("root:child-stop-subagent"), + runAttemptId: null, + nativeTurnRef: null, + ordinal: 1, + status: "running", + startedAt: now, + completedAt: null, + } as const; yield* projections.apply({ id: EventId.make("provider-turn:stop-subagent"), type: "provider-turn.updated", threadId: childThreadId, occurredAt: now, - payload: { - id: providerTurnId, - providerThreadId, - nodeId: NodeId.make("root:child-stop-subagent"), - runAttemptId: null, - nativeTurnRef: null, - ordinal: 1, - status: "running", - startedAt: now, - completedAt: null, - }, + payload: providerTurn, }); const commandId = CommandId.make("interrupt:stop-subagent"); const accepted = yield* orchestrator.dispatch({ @@ -179,6 +181,44 @@ it.effect("interrupts only the selected running native Codex subagent", () => }, ], ); + for (const missing of ["turn", "session"] as const) { + yield* projections.apply({ + id: EventId.make(`provider-turn:stop-subagent:missing-${missing}`), + type: "provider-turn.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + ...providerTurn, + status: missing === "turn" ? "completed" : "running", + completedAt: missing === "turn" ? now : null, + }, + }); + yield* projections.apply({ + id: EventId.make(`provider-thread:stop-subagent:missing-${missing}`), + type: "provider-thread.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + ...providerThread, + providerSessionId: missing === "session" ? null : providerSessionId, + }, + }); + const rejectedCommandId = CommandId.make(`interrupt:stop-subagent:missing-${missing}`); + const rejected = yield* orchestrator + .dispatch({ + type: "subagent.interrupt", + commandId: rejectedCommandId, + threadId: parentThreadId, + subagentId, + }) + .pipe(Effect.flip); + assert.instanceOf(rejected, Orchestrator.OrchestratorDispatchError); + assert.equal( + rejected.cause, + `Child thread ${childThreadId} for subagent ${subagentId} has no running provider turn or provider session.`, + ); + assert.deepEqual(yield* outbox.listByCommandId(rejectedCommandId), []); + } yield* projections.apply({ id: EventId.make("subagent:stop-subagent:completed"), type: "subagent.updated", @@ -196,7 +236,10 @@ it.effect("interrupts only the selected running native Codex subagent", () => }) .pipe(Effect.flip); assert.instanceOf(rejected, Orchestrator.OrchestratorDispatchError); - assert.equal(rejected.cause, undefined); + assert.equal( + rejected.cause, + `Subagent ${subagentId} is not a running native Codex subagent with a child thread.`, + ); assert.deepEqual(yield* outbox.listByCommandId(settledCommandId), []); }).pipe(Effect.provide(testLayer)), ); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index a4bbbbc06860..54d6da458b8c 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -8408,6 +8408,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return yield* new OrchestratorDispatchError({ commandId: command.commandId, commandType: command.type, + cause: `Subagent ${command.subagentId} is not a running native Codex subagent with a child thread.`, }); } const childThreadId = subagent.childThreadId; @@ -8426,6 +8427,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return yield* new OrchestratorDispatchError({ commandId: command.commandId, commandType: command.type, + cause: `Child thread ${childThreadId} for subagent ${command.subagentId} has no running provider turn or provider session.`, }); } const providerSessionId = providerThread.providerSessionId; From fd4981059a2fc72649415b777409cebf5ca7e3e1 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Tue, 6 Oct 2026 21:07:32 +0300 Subject: [PATCH 06/11] fix(server): settle native subagents after session loss --- .../Orchestrator.control-reads.test.ts | 250 ++++++++++++++++++ .../src/orchestration-v2/Orchestrator.ts | 125 +++++++++ 2 files changed, 375 insertions(+) diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index 6e88480362f4..ffbb48869982 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -24,10 +24,12 @@ import * as SqlClient from "effect/sql/SqlClient"; import * as SqlitePersistence from "../persistence/Sqlite.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as EffectOutbox from "./EffectOutbox.ts"; +import * as EffectWorker from "./EffectWorker.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; +import { makeSubagentChildThread } from "./SubagentProjection.ts"; import * as ProviderReplayHarness from "./testkit/ProviderReplayHarness.ts"; const instanceId = ProviderInstanceId.make("codex"); @@ -79,6 +81,24 @@ it.effect("interrupts only the selected running native Codex subagent", () => creationSource: "web", }); } + yield* projections.apply({ + id: EventId.make("child:stop-subagent:lineage"), + type: "thread.metadata-updated", + threadId: childThreadId, + occurredAt: now, + payload: makeSubagentChildThread({ + parentThread: yield* projections.getThread(parentThreadId), + childThreadId, + parentNodeId: subagentId, + activeProviderThreadId: providerThreadId, + providerInstanceId: instanceId, + modelSelection, + title: "Worker", + now, + createdBy: "agent", + creationSource: "provider", + }), + }); const subagent = { id: subagentId, threadId: parentThreadId, @@ -149,6 +169,126 @@ it.effect("interrupts only the selected running native Codex subagent", () => occurredAt: now, payload: providerTurn, }); + for (const [threadId, nodeId, kind] of [ + [parentThreadId, subagentId, "subagent"], + [childThreadId, providerTurn.nodeId, "root_turn"], + ] as const) { + yield* projections.apply({ + id: EventId.make(`node:${nodeId}`), + type: "node.updated", + threadId, + occurredAt: now, + payload: { + id: nodeId, + threadId, + runId: null, + parentNodeId: null, + rootNodeId: nodeId, + kind, + status: "running", + countsForRun: false, + providerThreadId, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }); + } + const item = { + id: TurnItemId.make("item:stop-subagent"), + threadId: parentThreadId, + runId: null, + nodeId: subagentId, + providerThreadId, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running" as const, + title: subagent.title, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "subagent" as const, + subagentId, + origin: subagent.origin, + driver: subagent.driver, + providerInstanceId: instanceId, + childThreadId, + prompt: subagent.prompt, + result: null, + }; + yield* projections.apply({ + id: EventId.make("item:stop-subagent"), + type: "turn-item.updated", + threadId: parentThreadId, + occurredAt: now, + payload: item, + }); + yield* projections.apply({ + id: EventId.make("command:stop-subagent"), + type: "turn-item.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + id: TurnItemId.make("command:stop-subagent"), + threadId: childThreadId, + runId: null, + nodeId: providerTurn.nodeId, + providerThreadId, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "command_execution", + input: "sleep 100", + }, + }); + yield* projections.apply({ + id: EventId.make("message:stop-subagent"), + type: "message.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + id: MessageId.make("message:stop-subagent"), + threadId: childThreadId, + runId: null, + nodeId: providerTurn.nodeId, + role: "assistant", + text: "Working", + attachments: [], + streaming: true, + createdBy: "agent", + creationSource: "provider", + createdAt: now, + updatedAt: now, + }, + }); + yield* projections.apply({ + id: EventId.make("request:stop-subagent"), + type: "runtime-request.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + id: RuntimeRequestId.make("request:stop-subagent"), + nodeId: providerTurn.nodeId, + providerTurnId, + nativeRequestRef: null, + kind: "user_input", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }, + }); const commandId = CommandId.make("interrupt:stop-subagent"); const accepted = yield* orchestrator.dispatch({ type: "subagent.interrupt", @@ -181,6 +321,116 @@ it.effect("interrupts only the selected running native Codex subagent", () => }, ], ); + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + yield* worker.drain(); + const child = yield* projections.getThreadProjection(childThreadId); + const parent = yield* projections.getThreadProjection(parentThreadId); + assert.equal(child.providerTurns[0]?.status, "interrupted"); + assert.equal(child.providerThreads[0]?.status, "idle"); + assert.equal(child.nodes[0]?.status, "interrupted"); + assert.equal(child.turnItems[0]?.status, "interrupted"); + assert.equal(child.messages[0]?.streaming, false); + assert.equal(child.runtimeRequests[0]?.status, "cancelled"); + assert.equal(parent.subagents[0]?.status, "interrupted"); + assert.equal(parent.nodes[0]?.status, "interrupted"); + assert.equal(parent.turnItems[0]?.status, "interrupted"); + for (const race of [ + "completed-turn", + "completed-task", + "waiting-task", + "newer-turn", + ] as const) { + yield* projections.apply({ + id: EventId.make(`turn:${race}:running`), + type: "provider-turn.updated", + threadId: childThreadId, + occurredAt: now, + payload: providerTurn, + }); + yield* projections.apply({ + id: EventId.make(`task:${race}:running`), + type: "subagent.updated", + threadId: parentThreadId, + occurredAt: now, + payload: subagent, + }); + yield* projections.apply({ + id: EventId.make(`provider-thread:${race}:active`), + type: "provider-thread.updated", + threadId: childThreadId, + occurredAt: now, + payload: providerThread, + }); + yield* orchestrator.dispatch({ + type: "subagent.interrupt", + commandId: CommandId.make(`interrupt:${race}`), + threadId: parentThreadId, + subagentId, + }); + const newerTurnId = ProviderTurnId.make("turn:stop-subagent:newer"); + if (race === "completed-turn" || race === "newer-turn") { + yield* projections.apply({ + id: EventId.make(`turn:${race}:raced`), + type: "provider-turn.updated", + threadId: childThreadId, + occurredAt: now, + payload: + race === "newer-turn" + ? { ...providerTurn, id: newerTurnId, ordinal: 2 } + : { ...providerTurn, status: "completed", completedAt: now }, + }); + } + if (race !== "newer-turn") { + yield* projections.apply({ + id: EventId.make(`task:${race}:completed`), + type: "subagent.updated", + threadId: parentThreadId, + occurredAt: now, + payload: { + ...subagent, + status: race === "waiting-task" ? "waiting" : "completed", + completedAt: race === "waiting-task" ? null : now, + }, + }); + } + yield* worker.drain(); + const racedChild = yield* projections.getThreadProjection(childThreadId); + assert.equal( + racedChild.providerTurns.find((turn) => turn.id === providerTurnId)?.status, + race === "completed-turn" ? "completed" : "interrupted", + ); + assert.equal( + (yield* projections.getThreadProjection(parentThreadId)).subagents[0]?.status, + race === "newer-turn" ? "running" : race === "waiting-task" ? "interrupted" : "completed", + ); + if (race === "newer-turn") { + assert.equal( + racedChild.providerTurns.find((turn) => turn.id === newerTurnId)?.status, + "running", + ); + assert.equal(racedChild.providerThreads[0]?.status, "active"); + yield* projections.apply({ + id: EventId.make("turn:newer:completed"), + type: "provider-turn.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + ...providerTurn, + id: newerTurnId, + ordinal: 2, + status: "completed", + completedAt: now, + }, + }); + } + } + yield* projections.apply({ + id: EventId.make("subagent:stop-subagent:running"), + type: "subagent.updated", + threadId: parentThreadId, + occurredAt: now, + payload: subagent, + }); for (const missing of ["turn", "session"] as const) { yield* projections.apply({ id: EventId.make(`provider-turn:stop-subagent:missing-${missing}`), diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index e402791adf14..0b25f7d7856f 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -16,6 +16,7 @@ import { type ChatAttachment, CommandId, isProviderNativeSubagentThread, + isOrchestrationV2WorkActive, MessageId, type ModelSelection, OrchestrationV2Command, @@ -8491,6 +8492,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio "runtimeRequests", "turnItems", "providerThreads", + "providerTurns", ], { turnItemTypes: [ @@ -8512,6 +8514,129 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), ), ); + const turn = stopped.providerTurn; + if ( + turn?.runAttemptId === null && + turn.status === "running" && + stopped.providerThread?.driver === "codex" && + isProviderNativeSubagentThread(projection.thread) + ) { + const now = yield* DateTime.now; + const emitEvent = emit(events, command); + yield* emitEvent({ + type: "provider-turn.updated", + threadId: command.threadId, + nodeId: turn.nodeId, + occurredAt: now, + payload: { ...turn, status: "interrupted", completedAt: now }, + }); + const newerTurn = projection.providerTurns.some( + (candidate) => + candidate.providerThreadId === turn.providerThreadId && + candidate.ordinal > turn.ordinal, + ); + const parentThreadId = projection.thread.lineage.parentThreadId; + const parent = + parentThreadId === null || newerTurn + ? undefined + : yield* projectionStore + .getThreadRecords(parentThreadId, ["subagents", "nodes", "turnItems"], { + turnItemTypes: ["subagent"], + turnItemStatuses: ["pending", "running", "waiting"], + }) + .pipe(mapDispatchError(command)); + const task = parent?.subagents.find( + (candidate) => + candidate.origin === "provider_native" && + candidate.driver === "codex" && + candidate.childThreadId === command.threadId && + isOrchestrationV2WorkActive(candidate.status), + ); + if (task !== undefined) { + yield* emitEvent({ + type: "subagent.updated", + threadId: task.threadId, + ...(task.runId === null ? {} : { runId: task.runId }), + nodeId: task.id, + occurredAt: now, + payload: { ...task, status: "interrupted", completedAt: now, updatedAt: now }, + }); + } + if (!newerTurn && stopped.providerThread.status === "active") { + yield* emitEvent({ + type: "provider-thread.updated", + threadId: command.threadId, + occurredAt: now, + payload: { ...stopped.providerThread, status: "idle", updatedAt: now }, + }); + } + const nodes = [ + ...projection.nodes.filter((node) => node.providerTurnId === turn.id), + ...(parent?.nodes.filter((node) => node.id === task?.id) ?? []), + ]; + for (const node of nodes) { + if (!isOrchestrationV2WorkActive(node.status)) continue; + yield* emitEvent({ + type: "node.updated", + threadId: node.threadId, + ...(node.runId === null ? {} : { runId: node.runId }), + nodeId: node.id, + occurredAt: now, + payload: { ...node, status: "interrupted", completedAt: now }, + }); + } + const items = [ + ...projection.turnItems.filter((item) => item.providerTurnId === turn.id), + ...(parent?.turnItems.filter( + (item) => item.type === "subagent" && item.subagentId === task?.id, + ) ?? []), + ]; + for (const item of items) { + yield* emitEvent({ + type: "turn-item.updated", + threadId: item.threadId, + ...(item.runId === null ? {} : { runId: item.runId }), + ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), + occurredAt: now, + payload: { + ...item, + ...("streaming" in item ? { streaming: false } : {}), + status: "interrupted", + completedAt: now, + updatedAt: now, + }, + }); + } + const output = yield* projectionStore + .getThreadRecords(command.threadId, ["messages"], { messageRoles: ["assistant"] }) + .pipe(mapDispatchError(command)); + for (const message of output.messages) { + if (!message.streaming || !nodes.some((node) => node.id === message.nodeId)) continue; + yield* emitEvent({ + type: "message.updated", + threadId: command.threadId, + occurredAt: now, + payload: { ...message, streaming: false, updatedAt: now }, + }); + } + for (const request of projection.runtimeRequests) { + if (request.status !== "pending" || !nodes.some((node) => node.id === request.nodeId)) + continue; + yield* emitEvent({ + type: "runtime-request.updated", + threadId: command.threadId, + nodeId: request.nodeId, + occurredAt: now, + payload: { + ...request, + status: "cancelled", + resolvedAt: now, + responseCapability: { type: "not_resumable", reason: "The turn was interrupted." }, + }, + }); + } + return; + } const stoppedRunId = projection.attempts.find( (attempt) => attempt.id === stopped.providerTurn?.runAttemptId, )?.runId; From 680a0258679eee32c979a527228b9d1fabd06b75 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Tue, 6 Oct 2026 21:27:09 +0300 Subject: [PATCH 07/11] fix(server): serialize native subagent parent settlement --- .../Orchestrator.control-reads.test.ts | 161 +++++++++++++++++- .../src/orchestration-v2/Orchestrator.ts | 126 +++++++++----- .../orchestration-v2/ProviderEventIngestor.ts | 14 ++ 3 files changed, 262 insertions(+), 39 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index ffbb48869982..a66037ae2c2e 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -18,6 +18,9 @@ import { TurnItemId, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; +import * as Deferred from "effect/Deferred"; +import * as Fiber from "effect/Fiber"; +import { vi } from "vite-plus/test"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as SqlClient from "effect/sql/SqlClient"; @@ -25,6 +28,10 @@ import * as SqlitePersistence from "../persistence/Sqlite.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as EffectOutbox from "./EffectOutbox.ts"; import * as EffectWorker from "./EffectWorker.ts"; +import * as EventSink from "./EventSink.ts"; +import * as IdAllocator from "./IdAllocator.ts"; +import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; +import * as ThreadCommandExecutor from "./ThreadCommandExecutor.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; @@ -42,7 +49,9 @@ const adapter = { openSession: () => Effect.die("No provider process needed for metadata controls"), } as ProviderAdapterV2Shape; const layerDatabase = SqlitePersistence.layerMemory; -const layerTest = Layer.mergeAll( +const layerControls = Layer.mergeAll( + ThreadCommandExecutor.layer, + IdAllocator.layer, layerDatabase, ProjectionStore.layer.pipe(Layer.provide(layerDatabase)), EffectOutbox.layer.pipe(Layer.provide(layerDatabase)), @@ -53,6 +62,11 @@ const layerTest = Layer.mergeAll( ), ); +const layerTest = Layer.merge( + layerControls, + ProviderEventIngestor.layer.pipe(Layer.provide(layerControls)), +); + it.effect("interrupts only the selected running native Codex subagent", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; @@ -334,6 +348,151 @@ it.effect("interrupts only the selected running native Codex subagent", () => assert.equal(parent.subagents[0]?.status, "interrupted"); assert.equal(parent.nodes[0]?.status, "interrupted"); assert.equal(parent.turnItems[0]?.status, "interrupted"); + const eventSink = yield* EventSink.EventSinkV2; + const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor; + const ingestor = yield* ProviderEventIngestor.ProviderEventIngestorV2; + for (const order of ["completion-first", "settlement-first"] as const) { + for (const event of [ + { type: "provider-turn.updated", threadId: childThreadId, payload: providerTurn }, + { type: "subagent.updated", threadId: parentThreadId, payload: subagent }, + { + type: "node.updated", + threadId: parentThreadId, + payload: { ...parent.nodes[0]!, status: "running", completedAt: null }, + }, + { type: "turn-item.updated", threadId: parentThreadId, payload: item }, + ] as const) { + yield* projections.apply({ + ...event, + id: EventId.make(`reset:${order}:${event.type}`), + occurredAt: now, + }); + } + const settleCommandId = CommandId.make(`settle:${order}`); + const reachedCommit = yield* Deferred.make(); + const releaseCommit = yield* Deferred.make(); + const providerQueued = yield* Deferred.make(); + const commitCommand = eventSink.commitCommand; + const writeWithEffects = eventSink.writeWithEffects; + const withLock = executor.withLock; + let parentWritesQueued = 0; + const commitSpy = vi + .spyOn(eventSink, "commitCommand") + .mockImplementation((input) => + order === "completion-first" && input.commandId === settleCommandId + ? Deferred.succeed(reachedCommit, undefined).pipe( + Effect.andThen(Deferred.await(releaseCommit)), + Effect.andThen(commitCommand(input)), + ) + : commitCommand(input), + ); + const writeSpy = vi + .spyOn(eventSink, "writeWithEffects") + .mockImplementation((input) => + order === "settlement-first" && + input.events.some((event) => event.type === "subagent.updated") + ? Deferred.succeed(reachedCommit, undefined).pipe( + Effect.andThen(Deferred.await(releaseCommit)), + Effect.andThen(writeWithEffects(input)), + ) + : writeWithEffects(input), + ); + const observeLock: ThreadCommandExecutor.ThreadCommandExecutor["Service"]["withLock"] = ( + key, + effect, + ) => + key === parentThreadId && ++parentWritesQueued === 4 + ? Deferred.succeed(providerQueued, undefined).pipe(Effect.andThen(withLock(key, effect))) + : withLock(key, effect); + const lockSpy = vi.spyOn(executor, "withLock").mockImplementation(observeLock); + yield* Effect.gen(function* () { + const settle = yield* orchestrator + .dispatch({ + type: "thread.background-work.settle", + commandId: settleCommandId, + threadId: childThreadId, + providerThreadId, + providerTurnId, + }) + .pipe(Effect.forkChild({ startImmediately: true })); + yield* Effect.raceFirst( + Deferred.await(reachedCommit), + Fiber.join(settle).pipe( + Effect.andThen(Effect.die("Settlement bypassed the commit barrier.")), + ), + ); + const completions = [ + { + type: "subagent.updated", + driver: adapter.driver, + subagent: { + ...subagent, + status: "completed", + completedAt: now, + completionWake: "always", + }, + }, + { + type: "node.updated", + driver: adapter.driver, + node: { ...parent.nodes[0]!, status: "completed", completedAt: now }, + }, + { + type: "turn_item.updated", + driver: adapter.driver, + turnItem: { ...item, status: "completed", completedAt: now, result: "Done" }, + }, + ] as const; + const complete = Effect.forEach( + completions, + (event) => + ingestor.ingestNormalized({ + providerSessionId, + providerInstanceId: instanceId, + threadId: parentThreadId, + event, + }), + { concurrency: "unbounded" }, + ); + if (order === "completion-first") { + yield* complete; + yield* Deferred.succeed(releaseCommit, undefined); + yield* Fiber.join(settle); + } else { + const completion = yield* complete.pipe(Effect.forkChild({ startImmediately: true })); + yield* Effect.raceFirst( + Deferred.await(providerQueued), + Fiber.join(completion).pipe( + Effect.andThen(Effect.die("Provider updates bypassed the parent lock.")), + ), + ); + yield* Deferred.succeed(releaseCommit, undefined); + yield* Fiber.join(settle); + yield* Fiber.join(completion); + } + const completedParent = yield* projections.getThreadProjection(parentThreadId); + assert.equal(completedParent.subagents[0]?.status, "completed"); + assert.equal(completedParent.subagents[0]?.completionWake, "always"); + assert.equal(completedParent.nodes[0]?.status, "completed"); + assert.equal(completedParent.turnItems[0]?.status, "completed"); + assert.equal( + completedParent.turnItems[0]?.type === "subagent" && completedParent.turnItems[0].result, + "Done", + ); + assert.equal( + (yield* projections.getThreadProjection(childThreadId)).providerTurns[0]?.status, + "interrupted", + ); + }).pipe( + Effect.ensuring( + Effect.sync(() => { + commitSpy.mockRestore(); + writeSpy.mockRestore(); + lockSpy.mockRestore(); + }), + ), + ); + } for (const race of [ "completed-turn", "completed-task", diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 0b25f7d7856f..c1c74bdce386 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -8535,33 +8535,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio candidate.providerThreadId === turn.providerThreadId && candidate.ordinal > turn.ordinal, ); - const parentThreadId = projection.thread.lineage.parentThreadId; - const parent = - parentThreadId === null || newerTurn - ? undefined - : yield* projectionStore - .getThreadRecords(parentThreadId, ["subagents", "nodes", "turnItems"], { - turnItemTypes: ["subagent"], - turnItemStatuses: ["pending", "running", "waiting"], - }) - .pipe(mapDispatchError(command)); - const task = parent?.subagents.find( - (candidate) => - candidate.origin === "provider_native" && - candidate.driver === "codex" && - candidate.childThreadId === command.threadId && - isOrchestrationV2WorkActive(candidate.status), - ); - if (task !== undefined) { - yield* emitEvent({ - type: "subagent.updated", - threadId: task.threadId, - ...(task.runId === null ? {} : { runId: task.runId }), - nodeId: task.id, - occurredAt: now, - payload: { ...task, status: "interrupted", completedAt: now, updatedAt: now }, - }); - } if (!newerTurn && stopped.providerThread.status === "active") { yield* emitEvent({ type: "provider-thread.updated", @@ -8570,10 +8543,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio payload: { ...stopped.providerThread, status: "idle", updatedAt: now }, }); } - const nodes = [ - ...projection.nodes.filter((node) => node.providerTurnId === turn.id), - ...(parent?.nodes.filter((node) => node.id === task?.id) ?? []), - ]; + const nodes = projection.nodes.filter((node) => node.providerTurnId === turn.id); for (const node of nodes) { if (!isOrchestrationV2WorkActive(node.status)) continue; yield* emitEvent({ @@ -8585,12 +8555,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio payload: { ...node, status: "interrupted", completedAt: now }, }); } - const items = [ - ...projection.turnItems.filter((item) => item.providerTurnId === turn.id), - ...(parent?.turnItems.filter( - (item) => item.type === "subagent" && item.subagentId === task?.id, - ) ?? []), - ]; + const items = projection.turnItems.filter((item) => item.providerTurnId === turn.id); for (const item of items) { yield* emitEvent({ type: "turn-item.updated", @@ -10728,8 +10693,93 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } satisfies OrchestratorV2DispatchResult; }); + const settleNativeSubagentParent = ( + command: Extract< + OrchestrationV2InternalCommand, + { readonly type: "thread.background-work.settle" } + >, + parentThreadId: ThreadId, + ) => + Effect.gen(function* () { + const child = yield* projectionStore.getThreadRecords(command.threadId, ["providerTurns"]); + const turn = child.providerTurns.find((candidate) => candidate.id === command.providerTurnId); + if ( + turn?.runAttemptId !== null || + turn.status !== "interrupted" || + !isProviderNativeSubagentThread(child.thread) || + child.providerTurns.some( + (candidate) => + candidate.providerThreadId === turn.providerThreadId && + candidate.ordinal > turn.ordinal, + ) + ) + return; + const parent = yield* projectionStore.getThreadRecords( + parentThreadId, + ["subagents", "nodes", "turnItems"], + { turnItemTypes: ["subagent"], turnItemStatuses: ["pending", "running", "waiting"] }, + ); + const task = parent.subagents.find( + (candidate) => + candidate.origin === "provider_native" && + candidate.driver === "codex" && + candidate.childThreadId === command.threadId && + isOrchestrationV2WorkActive(candidate.status), + ); + if (task === undefined) return; + const now = yield* DateTime.now; + yield* writeSystemEvents([ + { + type: "subagent.updated", + threadId: parentThreadId, + ...(task.runId === null ? {} : { runId: task.runId }), + nodeId: task.id, + occurredAt: now, + payload: { ...task, status: "interrupted", completedAt: now, updatedAt: now }, + }, + ...parent.nodes + .filter((node) => node.id === task.id && isOrchestrationV2WorkActive(node.status)) + .map((node) => ({ + type: "node.updated" as const, + threadId: parentThreadId, + ...(node.runId === null ? {} : { runId: node.runId }), + nodeId: node.id, + occurredAt: now, + payload: { ...node, status: "interrupted" as const, completedAt: now }, + })), + ...parent.turnItems + .filter((item) => item.type === "subagent" && item.subagentId === task.id) + .map((item) => ({ + type: "turn-item.updated" as const, + threadId: parentThreadId, + ...(item.runId === null ? {} : { runId: item.runId }), + ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), + occurredAt: now, + payload: { ...item, status: "interrupted" as const, completedAt: now, updatedAt: now }, + })), + ]); + }).pipe(mapDispatchError(command)); + const dispatchWithReceipt = (command: OrchestrationV2ServerCommand) => - threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command)); + Effect.gen(function* () { + const result = yield* threadDispatch.withLock( + commandThreadId(command), + dispatchWithReceiptEffect(command), + ); + if (command.type === "thread.background-work.settle") { + const child = yield* projectionStore + .getThread(command.threadId) + .pipe(mapDispatchError(command)); + const parentThreadId = child.lineage.parentThreadId; + if (isProviderNativeSubagentThread(child) && parentThreadId !== null) { + yield* threadDispatch.withLock( + parentThreadId, + settleNativeSubagentParent(command, parentThreadId), + ); + } + } + return result; + }); const handleTerminalRun = (stored: OrchestrationV2StoredEvent) => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts index ebf6243d357d..b94fb3099727 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts @@ -591,6 +591,20 @@ export const layer: Layer.Layer< .pipe(Effect.mapError(mapWriteError)); return result.storedEvents; }).pipe( + (write) => { + const event = input.event; + const parentThreadId = + event.type === "subagent.updated" + ? event.subagent.threadId + : event.type === "node.updated" && event.node.kind === "subagent" + ? event.node.threadId + : event.type === "turn_item.updated" && event.turnItem.type === "subagent" + ? event.turnItem.threadId + : undefined; + return parentThreadId === undefined + ? write + : threadCommands.withLock(parentThreadId, write); + }, Effect.flatMap((storedEvents) => storedEvents.length === 0 || input.event.type !== "subagent.updated" ? Effect.succeed(storedEvents) From 1cc05844b84c1d87bb1020149ec07b3d8561cdb0 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Tue, 6 Oct 2026 21:29:29 +0300 Subject: [PATCH 08/11] fix(web): allow stopping waiting Codex subagents --- .../Orchestrator.control-reads.test.ts | 31 ++++++++++++++++++- .../src/orchestration-v2/Orchestrator.ts | 4 +-- ...ThreadRelationshipsControl.agents.test.tsx | 13 ++++++++ .../chat/ThreadRelationshipsControl.tsx | 2 +- 4 files changed, 46 insertions(+), 4 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index a66037ae2c2e..9703977bc55b 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -348,6 +348,35 @@ it.effect("interrupts only the selected running native Codex subagent", () => assert.equal(parent.subagents[0]?.status, "interrupted"); assert.equal(parent.nodes[0]?.status, "interrupted"); assert.equal(parent.turnItems[0]?.status, "interrupted"); + yield* projections.apply({ + id: EventId.make("turn:waiting-stop"), + type: "provider-turn.updated", + threadId: childThreadId, + occurredAt: now, + payload: providerTurn, + }); + yield* projections.apply({ + id: EventId.make("task:waiting-stop"), + type: "subagent.updated", + threadId: parentThreadId, + occurredAt: now, + payload: { ...subagent, status: "waiting" }, + }); + yield* orchestrator.dispatch({ + type: "subagent.interrupt", + commandId: CommandId.make("interrupt:waiting-stop"), + threadId: parentThreadId, + subagentId, + }); + yield* worker.drain(); + assert.equal( + (yield* projections.getThreadProjection(childThreadId)).providerTurns[0]?.status, + "interrupted", + ); + assert.equal( + (yield* projections.getThreadProjection(parentThreadId)).subagents[0]?.status, + "interrupted", + ); const eventSink = yield* EventSink.EventSinkV2; const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor; const ingestor = yield* ProviderEventIngestor.ProviderEventIngestorV2; @@ -647,7 +676,7 @@ it.effect("interrupts only the selected running native Codex subagent", () => assert.instanceOf(rejected, Orchestrator.OrchestratorDispatchError); assert.equal( rejected.cause, - `Subagent ${subagentId} is not a running native Codex subagent with a child thread.`, + `Subagent ${subagentId} is not an active native Codex subagent with a child thread.`, ); assert.deepEqual(yield* outbox.listByCommandId(settledCommandId), []); }).pipe(Effect.provide(layerTest)), diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index c1c74bdce386..9d7f17c5a25a 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -9050,13 +9050,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio if ( subagent?.origin !== "provider_native" || subagent.driver !== "codex" || - subagent.status !== "running" || + (subagent.status !== "running" && subagent.status !== "waiting") || subagent.childThreadId === null ) { return yield* new OrchestratorDispatchError({ commandId: command.commandId, commandType: command.type, - cause: `Subagent ${command.subagentId} is not a running native Codex subagent with a child thread.`, + cause: `Subagent ${command.subagentId} is not an active native Codex subagent with a child thread.`, }); } const childThreadId = subagent.childThreadId; diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx index 98d027dee8c1..a6a802c3b348 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx @@ -123,6 +123,19 @@ it("stops only active subagents from lineage without opening their thread", asyn input: { threadId: "parent", subagentId: "agent" }, }); + state.command.mockClear(); + state.projection = { + ...projection, + subagents: [{ ...agent, origin: "provider_native", status: "waiting" }], + }; + await act(async () => renderer.update(cloneElement(panel))); + await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenCalledWith({ + environmentId: "test", + input: { threadId: "parent", subagentId: "agent" }, + }); + expect(state.navigate).not.toHaveBeenCalled(); + state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; await act(async () => renderer.update(cloneElement(panel))); expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx index f3a7e2acc1e3..41edf13a0211 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx @@ -408,7 +408,7 @@ export function ThreadRelationshipsPanel(props: { (supportsSubagentInterrupt && agent.origin === "provider_native" && agent.driver === "codex" && - agent.status === "running")); + (agent.status === "running" || agent.status === "waiting"))); const threadTitle = relationshipThreadTitle({ title: node?.thread?.title ?? agent?.title ?? threadId, isSubagent, From 86a5ab32d944251b8e85e50f3456634d3cfe8653 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Tue, 6 Oct 2026 21:34:42 +0300 Subject: [PATCH 09/11] fix(web): gate native agent stops on stored task status --- ...ThreadRelationshipsControl.agents.test.tsx | 26 +++++++++++++++++++ .../chat/ThreadRelationshipsControl.tsx | 8 +++--- 2 files changed, 30 insertions(+), 4 deletions(-) diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx index a6a802c3b348..a87376a00d36 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx @@ -136,6 +136,32 @@ it("stops only active subagents from lineage without opening their thread", asyn }); expect(state.navigate).not.toHaveBeenCalled(); + for (const status of ["running", "waiting"] as const) { + state.shells = [ + { + environmentId: "test", + source: { + ...child, + activityRunStatus: status, + activityRunStartedAt: DateTime.makeUnsafe("2026-09-16T12:05:00Z"), + }, + }, + ]; + state.projection = { + ...projection, + subagents: [{ ...agent, origin: "provider_native", status: "completed" }], + }; + await act(async () => renderer.update(cloneElement(panel))); + expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); + state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; + await act(async () => renderer.update(cloneElement(panel))); + await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenLastCalledWith({ + environmentId: "test", + input: { threadId: "child" }, + }); + } + state.shells = [{ environmentId: "test", source: child }]; state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; await act(async () => renderer.update(cloneElement(panel))); expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx index 41edf13a0211..1400ffa11a19 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx @@ -397,10 +397,9 @@ export function ThreadRelationshipsPanel(props: { ? BotIcon : GitForkIcon; const relationship = relationshipLabel(edge, props.threadId); - const agent = liveSubagent( - isSubagent && !isParent ? subagentsByThreadId.get(threadId) : undefined, - node?.thread, - ); + const storedAgent = + isSubagent && !isParent ? subagentsByThreadId.get(threadId) : undefined; + const agent = liveSubagent(storedAgent, node?.thread); const canStop = agent?.startedAt && ((agent.origin === "app_owned" && @@ -408,6 +407,7 @@ export function ThreadRelationshipsPanel(props: { (supportsSubagentInterrupt && agent.origin === "provider_native" && agent.driver === "codex" && + (storedAgent?.status === "running" || storedAgent?.status === "waiting") && (agent.status === "running" || agent.status === "waiting"))); const threadTitle = relationshipThreadTitle({ title: node?.thread?.title ?? agent?.title ?? threadId, From 62378883d5c27fbee043cb9815f0748a4f4c3c24 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Tue, 6 Oct 2026 22:43:58 +0300 Subject: [PATCH 10/11] fix(web): limit Lineage stop to T3-owned subagents --- .../src/environment/ServerEnvironment.test.ts | 1 - .../src/environment/ServerEnvironment.ts | 1 - .../Orchestrator.control-reads.test.ts | 635 +----------------- .../src/orchestration-v2/Orchestrator.ts | 248 +------ .../orchestration-v2/ProjectionStore.test.ts | 17 +- .../src/orchestration-v2/ProjectionStore.ts | 5 - .../orchestration-v2/ProviderEventIngestor.ts | 14 - .../testkit/OrchestratorScenario.ts | 1 - .../src/relay/AgentAwarenessRelay.test.ts | 9 - apps/server/src/relay/AgentAwarenessRelay.ts | 1 - ...ThreadRelationshipsControl.agents.test.tsx | 213 +++--- .../chat/ThreadRelationshipsControl.tsx | 71 +- .../src/operations/commands.test.ts | 35 - .../client-runtime/src/operations/commands.ts | 26 - .../state/orchestrationV2Projection.test.ts | 13 - .../src/state/orchestrationV2Projection.ts | 2 - .../src/state/threadCommands.ts | 8 - packages/contracts/src/environment.test.ts | 12 - packages/contracts/src/environment.ts | 1 - packages/contracts/src/orchestrationV2.ts | 16 - 20 files changed, 121 insertions(+), 1208 deletions(-) diff --git a/apps/server/src/environment/ServerEnvironment.test.ts b/apps/server/src/environment/ServerEnvironment.test.ts index bb281c975ddb..41c18b907437 100644 --- a/apps/server/src/environment/ServerEnvironment.test.ts +++ b/apps/server/src/environment/ServerEnvironment.test.ts @@ -216,7 +216,6 @@ it.layer(NodeServices.layer)("ServerEnvironmentLive", (it) => { expect(first.environmentId).toBe(second.environmentId); expect(first.orchestrationProtocolVersion).toBe(ORCHESTRATION_PROTOCOL_VERSION); expect(second.capabilities.repositoryIdentity).toBe(true); - expect(second.capabilities.subagentInterrupt).toBe(true); expect(second.capabilities.connectionProbe).toBe(true); expect(second.capabilities.attachmentUploads).toBe(true); expect(second.capabilities.fileAttachments).toEqual({ maxUploadBytes: 50 * 1024 * 1024 }); diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index e66fdcb460f9..6381547468ae 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -218,7 +218,6 @@ export const make = Effect.gen(function* () { capabilities: { repositoryIdentity: true, connectionProbe: true, - subagentInterrupt: true, attachmentUploads: true, questionAttachments: true, fileAttachments: { maxUploadBytes: PROVIDER_SEND_TURN_MAX_FILE_BYTES }, diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index 9703977bc55b..c68fc0ace5e9 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -18,25 +18,15 @@ import { TurnItemId, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; -import * as Deferred from "effect/Deferred"; -import * as Fiber from "effect/Fiber"; -import { vi } from "vite-plus/test"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as SqlClient from "effect/sql/SqlClient"; import * as SqlitePersistence from "../persistence/Sqlite.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; -import * as EffectOutbox from "./EffectOutbox.ts"; -import * as EffectWorker from "./EffectWorker.ts"; -import * as EventSink from "./EventSink.ts"; -import * as IdAllocator from "./IdAllocator.ts"; -import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; -import * as ThreadCommandExecutor from "./ThreadCommandExecutor.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; -import { makeSubagentChildThread } from "./SubagentProjection.ts"; import * as ProviderReplayHarness from "./testkit/ProviderReplayHarness.ts"; const instanceId = ProviderInstanceId.make("codex"); @@ -49,12 +39,9 @@ const adapter = { openSession: () => Effect.die("No provider process needed for metadata controls"), } as ProviderAdapterV2Shape; const layerDatabase = SqlitePersistence.layerMemory; -const layerControls = Layer.mergeAll( - ThreadCommandExecutor.layer, - IdAllocator.layer, +const layerTest = Layer.mergeAll( layerDatabase, ProjectionStore.layer.pipe(Layer.provide(layerDatabase)), - EffectOutbox.layer.pipe(Layer.provide(layerDatabase)), ProviderReplayHarness.layerWithRegistry( { name: "control-reads" }, ProviderAdapterRegistry.layerFromAdapters([adapter]), @@ -62,626 +49,6 @@ const layerControls = Layer.mergeAll( ), ); -const layerTest = Layer.merge( - layerControls, - ProviderEventIngestor.layer.pipe(Layer.provide(layerControls)), -); - -it.effect("interrupts only the selected running native Codex subagent", () => - Effect.gen(function* () { - const orchestrator = yield* Orchestrator.OrchestratorV2; - const projections = yield* ProjectionStore.ProjectionStoreV2; - const outbox = yield* EffectOutbox.EffectOutboxV2; - const parentThreadId = ThreadId.make("parent:stop-subagent"); - const childThreadId = ThreadId.make("child:stop-subagent"); - const subagentId = NodeId.make("subagent:stop-subagent"); - const providerThreadId = ProviderThreadId.make("provider-thread:stop-subagent"); - const providerTurnId = ProviderTurnId.make("provider-turn:stop-subagent"); - const providerSessionId = ProviderSessionId.make("session:stop-subagent"); - const now = yield* DateTime.now; - for (const threadId of [parentThreadId, childThreadId]) { - yield* orchestrator.dispatch({ - type: "thread.create", - commandId: CommandId.make(`create:${threadId}`), - threadId, - projectId: ProjectId.make("project:stop-subagent"), - title: "Agent", - modelSelection, - runtimeMode: "full-access", - interactionMode: "default", - branch: null, - worktreePath: null, - createdBy: "user", - creationSource: "web", - }); - } - yield* projections.apply({ - id: EventId.make("child:stop-subagent:lineage"), - type: "thread.metadata-updated", - threadId: childThreadId, - occurredAt: now, - payload: makeSubagentChildThread({ - parentThread: yield* projections.getThread(parentThreadId), - childThreadId, - parentNodeId: subagentId, - activeProviderThreadId: providerThreadId, - providerInstanceId: instanceId, - modelSelection, - title: "Worker", - now, - createdBy: "agent", - creationSource: "provider", - }), - }); - const subagent = { - id: subagentId, - threadId: parentThreadId, - runId: null, - parentNodeId: NodeId.make("root:stop-subagent"), - origin: "provider_native" as const, - createdBy: "agent" as const, - driver: adapter.driver, - providerInstanceId: instanceId, - providerThreadId: null, - childThreadId, - nativeTaskRef: null, - prompt: "Check this", - title: "Worker", - model: "gpt-5.1-codex", - status: "running" as const, - result: null, - startedAt: now, - completedAt: null, - updatedAt: now, - }; - yield* projections.apply({ - id: EventId.make("subagent:stop-subagent"), - type: "subagent.updated", - threadId: parentThreadId, - occurredAt: now, - payload: subagent, - }); - const providerThread = { - id: providerThreadId, - driver: adapter.driver, - providerInstanceId: instanceId, - providerSessionId, - appThreadId: childThreadId, - ownerNodeId: subagentId, - nativeThreadRef: null, - nativeConversationHeadRef: null, - status: "active", - firstRunOrdinal: null, - lastRunOrdinal: null, - handoffIds: [], - forkedFrom: null, - createdAt: now, - updatedAt: now, - } as const; - yield* projections.apply({ - id: EventId.make("provider-thread:stop-subagent"), - type: "provider-thread.updated", - threadId: childThreadId, - occurredAt: now, - payload: providerThread, - }); - const providerTurn = { - id: providerTurnId, - providerThreadId, - nodeId: NodeId.make("root:child-stop-subagent"), - runAttemptId: null, - nativeTurnRef: null, - ordinal: 1, - status: "running", - startedAt: now, - completedAt: null, - } as const; - yield* projections.apply({ - id: EventId.make("provider-turn:stop-subagent"), - type: "provider-turn.updated", - threadId: childThreadId, - occurredAt: now, - payload: providerTurn, - }); - for (const [threadId, nodeId, kind] of [ - [parentThreadId, subagentId, "subagent"], - [childThreadId, providerTurn.nodeId, "root_turn"], - ] as const) { - yield* projections.apply({ - id: EventId.make(`node:${nodeId}`), - type: "node.updated", - threadId, - occurredAt: now, - payload: { - id: nodeId, - threadId, - runId: null, - parentNodeId: null, - rootNodeId: nodeId, - kind, - status: "running", - countsForRun: false, - providerThreadId, - providerTurnId, - nativeItemRef: null, - runtimeRequestId: null, - checkpointScopeId: null, - startedAt: now, - completedAt: null, - }, - }); - } - const item = { - id: TurnItemId.make("item:stop-subagent"), - threadId: parentThreadId, - runId: null, - nodeId: subagentId, - providerThreadId, - providerTurnId: null, - nativeItemRef: null, - parentItemId: null, - ordinal: 1, - status: "running" as const, - title: subagent.title, - startedAt: now, - completedAt: null, - updatedAt: now, - type: "subagent" as const, - subagentId, - origin: subagent.origin, - driver: subagent.driver, - providerInstanceId: instanceId, - childThreadId, - prompt: subagent.prompt, - result: null, - }; - yield* projections.apply({ - id: EventId.make("item:stop-subagent"), - type: "turn-item.updated", - threadId: parentThreadId, - occurredAt: now, - payload: item, - }); - yield* projections.apply({ - id: EventId.make("command:stop-subagent"), - type: "turn-item.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - id: TurnItemId.make("command:stop-subagent"), - threadId: childThreadId, - runId: null, - nodeId: providerTurn.nodeId, - providerThreadId, - providerTurnId, - nativeItemRef: null, - parentItemId: null, - ordinal: 1, - status: "running", - title: null, - startedAt: now, - completedAt: null, - updatedAt: now, - type: "command_execution", - input: "sleep 100", - }, - }); - yield* projections.apply({ - id: EventId.make("message:stop-subagent"), - type: "message.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - id: MessageId.make("message:stop-subagent"), - threadId: childThreadId, - runId: null, - nodeId: providerTurn.nodeId, - role: "assistant", - text: "Working", - attachments: [], - streaming: true, - createdBy: "agent", - creationSource: "provider", - createdAt: now, - updatedAt: now, - }, - }); - yield* projections.apply({ - id: EventId.make("request:stop-subagent"), - type: "runtime-request.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - id: RuntimeRequestId.make("request:stop-subagent"), - nodeId: providerTurn.nodeId, - providerTurnId, - nativeRequestRef: null, - kind: "user_input", - status: "pending", - responseCapability: { type: "live", providerSessionId }, - createdAt: now, - resolvedAt: null, - }, - }); - const commandId = CommandId.make("interrupt:stop-subagent"); - const accepted = yield* orchestrator.dispatch({ - type: "subagent.interrupt", - commandId, - threadId: parentThreadId, - subagentId, - }); - assert.deepEqual( - accepted.storedEvents.map((event) => event.event.type), - ["subagent.interrupt-requested"], - ); - assert.equal( - (yield* projections.getThreadProjection(parentThreadId)).subagents[0]?.status, - "running", - ); - assert.deepEqual( - (yield* outbox.listByCommandId(commandId)).map(({ threadId, request }) => ({ - threadId, - request, - })), - [ - { - threadId: childThreadId, - request: { - type: "provider-turn.interrupt", - providerSessionId, - providerThreadId, - providerTurnId, - }, - }, - ], - ); - const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; - yield* worker.drain(); - const child = yield* projections.getThreadProjection(childThreadId); - const parent = yield* projections.getThreadProjection(parentThreadId); - assert.equal(child.providerTurns[0]?.status, "interrupted"); - assert.equal(child.providerThreads[0]?.status, "idle"); - assert.equal(child.nodes[0]?.status, "interrupted"); - assert.equal(child.turnItems[0]?.status, "interrupted"); - assert.equal(child.messages[0]?.streaming, false); - assert.equal(child.runtimeRequests[0]?.status, "cancelled"); - assert.equal(parent.subagents[0]?.status, "interrupted"); - assert.equal(parent.nodes[0]?.status, "interrupted"); - assert.equal(parent.turnItems[0]?.status, "interrupted"); - yield* projections.apply({ - id: EventId.make("turn:waiting-stop"), - type: "provider-turn.updated", - threadId: childThreadId, - occurredAt: now, - payload: providerTurn, - }); - yield* projections.apply({ - id: EventId.make("task:waiting-stop"), - type: "subagent.updated", - threadId: parentThreadId, - occurredAt: now, - payload: { ...subagent, status: "waiting" }, - }); - yield* orchestrator.dispatch({ - type: "subagent.interrupt", - commandId: CommandId.make("interrupt:waiting-stop"), - threadId: parentThreadId, - subagentId, - }); - yield* worker.drain(); - assert.equal( - (yield* projections.getThreadProjection(childThreadId)).providerTurns[0]?.status, - "interrupted", - ); - assert.equal( - (yield* projections.getThreadProjection(parentThreadId)).subagents[0]?.status, - "interrupted", - ); - const eventSink = yield* EventSink.EventSinkV2; - const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor; - const ingestor = yield* ProviderEventIngestor.ProviderEventIngestorV2; - for (const order of ["completion-first", "settlement-first"] as const) { - for (const event of [ - { type: "provider-turn.updated", threadId: childThreadId, payload: providerTurn }, - { type: "subagent.updated", threadId: parentThreadId, payload: subagent }, - { - type: "node.updated", - threadId: parentThreadId, - payload: { ...parent.nodes[0]!, status: "running", completedAt: null }, - }, - { type: "turn-item.updated", threadId: parentThreadId, payload: item }, - ] as const) { - yield* projections.apply({ - ...event, - id: EventId.make(`reset:${order}:${event.type}`), - occurredAt: now, - }); - } - const settleCommandId = CommandId.make(`settle:${order}`); - const reachedCommit = yield* Deferred.make(); - const releaseCommit = yield* Deferred.make(); - const providerQueued = yield* Deferred.make(); - const commitCommand = eventSink.commitCommand; - const writeWithEffects = eventSink.writeWithEffects; - const withLock = executor.withLock; - let parentWritesQueued = 0; - const commitSpy = vi - .spyOn(eventSink, "commitCommand") - .mockImplementation((input) => - order === "completion-first" && input.commandId === settleCommandId - ? Deferred.succeed(reachedCommit, undefined).pipe( - Effect.andThen(Deferred.await(releaseCommit)), - Effect.andThen(commitCommand(input)), - ) - : commitCommand(input), - ); - const writeSpy = vi - .spyOn(eventSink, "writeWithEffects") - .mockImplementation((input) => - order === "settlement-first" && - input.events.some((event) => event.type === "subagent.updated") - ? Deferred.succeed(reachedCommit, undefined).pipe( - Effect.andThen(Deferred.await(releaseCommit)), - Effect.andThen(writeWithEffects(input)), - ) - : writeWithEffects(input), - ); - const observeLock: ThreadCommandExecutor.ThreadCommandExecutor["Service"]["withLock"] = ( - key, - effect, - ) => - key === parentThreadId && ++parentWritesQueued === 4 - ? Deferred.succeed(providerQueued, undefined).pipe(Effect.andThen(withLock(key, effect))) - : withLock(key, effect); - const lockSpy = vi.spyOn(executor, "withLock").mockImplementation(observeLock); - yield* Effect.gen(function* () { - const settle = yield* orchestrator - .dispatch({ - type: "thread.background-work.settle", - commandId: settleCommandId, - threadId: childThreadId, - providerThreadId, - providerTurnId, - }) - .pipe(Effect.forkChild({ startImmediately: true })); - yield* Effect.raceFirst( - Deferred.await(reachedCommit), - Fiber.join(settle).pipe( - Effect.andThen(Effect.die("Settlement bypassed the commit barrier.")), - ), - ); - const completions = [ - { - type: "subagent.updated", - driver: adapter.driver, - subagent: { - ...subagent, - status: "completed", - completedAt: now, - completionWake: "always", - }, - }, - { - type: "node.updated", - driver: adapter.driver, - node: { ...parent.nodes[0]!, status: "completed", completedAt: now }, - }, - { - type: "turn_item.updated", - driver: adapter.driver, - turnItem: { ...item, status: "completed", completedAt: now, result: "Done" }, - }, - ] as const; - const complete = Effect.forEach( - completions, - (event) => - ingestor.ingestNormalized({ - providerSessionId, - providerInstanceId: instanceId, - threadId: parentThreadId, - event, - }), - { concurrency: "unbounded" }, - ); - if (order === "completion-first") { - yield* complete; - yield* Deferred.succeed(releaseCommit, undefined); - yield* Fiber.join(settle); - } else { - const completion = yield* complete.pipe(Effect.forkChild({ startImmediately: true })); - yield* Effect.raceFirst( - Deferred.await(providerQueued), - Fiber.join(completion).pipe( - Effect.andThen(Effect.die("Provider updates bypassed the parent lock.")), - ), - ); - yield* Deferred.succeed(releaseCommit, undefined); - yield* Fiber.join(settle); - yield* Fiber.join(completion); - } - const completedParent = yield* projections.getThreadProjection(parentThreadId); - assert.equal(completedParent.subagents[0]?.status, "completed"); - assert.equal(completedParent.subagents[0]?.completionWake, "always"); - assert.equal(completedParent.nodes[0]?.status, "completed"); - assert.equal(completedParent.turnItems[0]?.status, "completed"); - assert.equal( - completedParent.turnItems[0]?.type === "subagent" && completedParent.turnItems[0].result, - "Done", - ); - assert.equal( - (yield* projections.getThreadProjection(childThreadId)).providerTurns[0]?.status, - "interrupted", - ); - }).pipe( - Effect.ensuring( - Effect.sync(() => { - commitSpy.mockRestore(); - writeSpy.mockRestore(); - lockSpy.mockRestore(); - }), - ), - ); - } - for (const race of [ - "completed-turn", - "completed-task", - "waiting-task", - "newer-turn", - ] as const) { - yield* projections.apply({ - id: EventId.make(`turn:${race}:running`), - type: "provider-turn.updated", - threadId: childThreadId, - occurredAt: now, - payload: providerTurn, - }); - yield* projections.apply({ - id: EventId.make(`task:${race}:running`), - type: "subagent.updated", - threadId: parentThreadId, - occurredAt: now, - payload: subagent, - }); - yield* projections.apply({ - id: EventId.make(`provider-thread:${race}:active`), - type: "provider-thread.updated", - threadId: childThreadId, - occurredAt: now, - payload: providerThread, - }); - yield* orchestrator.dispatch({ - type: "subagent.interrupt", - commandId: CommandId.make(`interrupt:${race}`), - threadId: parentThreadId, - subagentId, - }); - const newerTurnId = ProviderTurnId.make("turn:stop-subagent:newer"); - if (race === "completed-turn" || race === "newer-turn") { - yield* projections.apply({ - id: EventId.make(`turn:${race}:raced`), - type: "provider-turn.updated", - threadId: childThreadId, - occurredAt: now, - payload: - race === "newer-turn" - ? { ...providerTurn, id: newerTurnId, ordinal: 2 } - : { ...providerTurn, status: "completed", completedAt: now }, - }); - } - if (race !== "newer-turn") { - yield* projections.apply({ - id: EventId.make(`task:${race}:completed`), - type: "subagent.updated", - threadId: parentThreadId, - occurredAt: now, - payload: { - ...subagent, - status: race === "waiting-task" ? "waiting" : "completed", - completedAt: race === "waiting-task" ? null : now, - }, - }); - } - yield* worker.drain(); - const racedChild = yield* projections.getThreadProjection(childThreadId); - assert.equal( - racedChild.providerTurns.find((turn) => turn.id === providerTurnId)?.status, - race === "completed-turn" ? "completed" : "interrupted", - ); - assert.equal( - (yield* projections.getThreadProjection(parentThreadId)).subagents[0]?.status, - race === "newer-turn" ? "running" : race === "waiting-task" ? "interrupted" : "completed", - ); - if (race === "newer-turn") { - assert.equal( - racedChild.providerTurns.find((turn) => turn.id === newerTurnId)?.status, - "running", - ); - assert.equal(racedChild.providerThreads[0]?.status, "active"); - yield* projections.apply({ - id: EventId.make("turn:newer:completed"), - type: "provider-turn.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - ...providerTurn, - id: newerTurnId, - ordinal: 2, - status: "completed", - completedAt: now, - }, - }); - } - } - yield* projections.apply({ - id: EventId.make("subagent:stop-subagent:running"), - type: "subagent.updated", - threadId: parentThreadId, - occurredAt: now, - payload: subagent, - }); - for (const missing of ["turn", "session"] as const) { - yield* projections.apply({ - id: EventId.make(`provider-turn:stop-subagent:missing-${missing}`), - type: "provider-turn.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - ...providerTurn, - status: missing === "turn" ? "completed" : "running", - completedAt: missing === "turn" ? now : null, - }, - }); - yield* projections.apply({ - id: EventId.make(`provider-thread:stop-subagent:missing-${missing}`), - type: "provider-thread.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - ...providerThread, - providerSessionId: missing === "session" ? null : providerSessionId, - }, - }); - const rejectedCommandId = CommandId.make(`interrupt:stop-subagent:missing-${missing}`); - const rejected = yield* orchestrator - .dispatch({ - type: "subagent.interrupt", - commandId: rejectedCommandId, - threadId: parentThreadId, - subagentId, - }) - .pipe(Effect.flip); - assert.instanceOf(rejected, Orchestrator.OrchestratorDispatchError); - assert.equal( - rejected.cause, - `Child thread ${childThreadId} for subagent ${subagentId} has no running provider turn or provider session.`, - ); - assert.deepEqual(yield* outbox.listByCommandId(rejectedCommandId), []); - } - yield* projections.apply({ - id: EventId.make("subagent:stop-subagent:completed"), - type: "subagent.updated", - threadId: parentThreadId, - occurredAt: now, - payload: { ...subagent, status: "completed", completedAt: now }, - }); - const settledCommandId = CommandId.make("interrupt:stop-subagent:settled"); - const rejected = yield* orchestrator - .dispatch({ - type: "subagent.interrupt", - commandId: settledCommandId, - threadId: parentThreadId, - subagentId, - }) - .pipe(Effect.flip); - assert.instanceOf(rejected, Orchestrator.OrchestratorDispatchError); - assert.equal( - rejected.cause, - `Subagent ${subagentId} is not an active native Codex subagent with a child thread.`, - ); - assert.deepEqual(yield* outbox.listByCommandId(settledCommandId), []); - }).pipe(Effect.provide(layerTest)), -); - it.effect( "dispatches metadata, queue resume and request controls without hydrating unrelated history", () => diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 9d7f17c5a25a..de100d291b8f 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -16,7 +16,6 @@ import { type ChatAttachment, CommandId, isProviderNativeSubagentThread, - isOrchestrationV2WorkActive, MessageId, type ModelSelection, OrchestrationV2Command, @@ -419,7 +418,6 @@ function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { case "prepared-run.fail": case "prepared-run.retry": case "run.interrupt": - case "subagent.interrupt": case "queued-message.promote-to-steer": case "queue.resume": case "queued-run.reorder": @@ -8492,7 +8490,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio "runtimeRequests", "turnItems", "providerThreads", - "providerTurns", ], { turnItemTypes: [ @@ -8514,94 +8511,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), ), ); - const turn = stopped.providerTurn; - if ( - turn?.runAttemptId === null && - turn.status === "running" && - stopped.providerThread?.driver === "codex" && - isProviderNativeSubagentThread(projection.thread) - ) { - const now = yield* DateTime.now; - const emitEvent = emit(events, command); - yield* emitEvent({ - type: "provider-turn.updated", - threadId: command.threadId, - nodeId: turn.nodeId, - occurredAt: now, - payload: { ...turn, status: "interrupted", completedAt: now }, - }); - const newerTurn = projection.providerTurns.some( - (candidate) => - candidate.providerThreadId === turn.providerThreadId && - candidate.ordinal > turn.ordinal, - ); - if (!newerTurn && stopped.providerThread.status === "active") { - yield* emitEvent({ - type: "provider-thread.updated", - threadId: command.threadId, - occurredAt: now, - payload: { ...stopped.providerThread, status: "idle", updatedAt: now }, - }); - } - const nodes = projection.nodes.filter((node) => node.providerTurnId === turn.id); - for (const node of nodes) { - if (!isOrchestrationV2WorkActive(node.status)) continue; - yield* emitEvent({ - type: "node.updated", - threadId: node.threadId, - ...(node.runId === null ? {} : { runId: node.runId }), - nodeId: node.id, - occurredAt: now, - payload: { ...node, status: "interrupted", completedAt: now }, - }); - } - const items = projection.turnItems.filter((item) => item.providerTurnId === turn.id); - for (const item of items) { - yield* emitEvent({ - type: "turn-item.updated", - threadId: item.threadId, - ...(item.runId === null ? {} : { runId: item.runId }), - ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), - occurredAt: now, - payload: { - ...item, - ...("streaming" in item ? { streaming: false } : {}), - status: "interrupted", - completedAt: now, - updatedAt: now, - }, - }); - } - const output = yield* projectionStore - .getThreadRecords(command.threadId, ["messages"], { messageRoles: ["assistant"] }) - .pipe(mapDispatchError(command)); - for (const message of output.messages) { - if (!message.streaming || !nodes.some((node) => node.id === message.nodeId)) continue; - yield* emitEvent({ - type: "message.updated", - threadId: command.threadId, - occurredAt: now, - payload: { ...message, streaming: false, updatedAt: now }, - }); - } - for (const request of projection.runtimeRequests) { - if (request.status !== "pending" || !nodes.some((node) => node.id === request.nodeId)) - continue; - yield* emitEvent({ - type: "runtime-request.updated", - threadId: command.threadId, - nodeId: request.nodeId, - occurredAt: now, - payload: { - ...request, - status: "cancelled", - resolvedAt: now, - responseCapability: { type: "not_resumable", reason: "The turn was interrupted." }, - }, - }); - } - return; - } const stoppedRunId = projection.attempts.find( (attempt) => attempt.id === stopped.providerTurn?.runAttemptId, )?.runId; @@ -9039,73 +8948,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return undefined; }); - const dispatchSubagentInterrupt = ( - command: Extract, - events: Ref.Ref>, - effects: Ref.Ref>, - ) => - Effect.gen(function* () { - const parent = yield* loadProjectionForCommand(command, ["subagents"]); - const subagent = parent.subagents.find((entry) => entry.id === command.subagentId); - if ( - subagent?.origin !== "provider_native" || - subagent.driver !== "codex" || - (subagent.status !== "running" && subagent.status !== "waiting") || - subagent.childThreadId === null - ) { - return yield* new OrchestratorDispatchError({ - commandId: command.commandId, - commandType: command.type, - cause: `Subagent ${command.subagentId} is not an active native Codex subagent with a child thread.`, - }); - } - const childThreadId = subagent.childThreadId; - const child = yield* projectionStore - .getThreadRecords(childThreadId, ["providerThreads", "providerTurns"]) - .pipe( - Effect.mapError( - (cause) => new OrchestratorProjectionError({ threadId: childThreadId, cause }), - ), - ); - const turn = child.providerTurns.findLast((entry) => entry.status === "running"); - const providerThread = child.providerThreads.find( - (entry) => entry.id === turn?.providerThreadId, - ); - if (turn === undefined || providerThread?.providerSessionId == null) { - return yield* new OrchestratorDispatchError({ - commandId: command.commandId, - commandType: command.type, - cause: `Child thread ${childThreadId} for subagent ${command.subagentId} has no running provider turn or provider session.`, - }); - } - const providerSessionId = providerThread.providerSessionId; - const emitEvent = emit(events, command); - yield* emitEvent({ - type: "subagent.interrupt-requested", - threadId: command.threadId, - ...(subagent.runId === null ? {} : { runId: subagent.runId }), - nodeId: subagent.id, - driver: subagent.driver, - providerInstanceId: subagent.providerInstanceId, - occurredAt: yield* DateTime.now, - payload: subagent.id, - }); - yield* Ref.update(effects, (existing) => [ - ...existing, - { - id: `effect:${command.commandId}:provider-turn.interrupt:${turn.id}`, - commandId: command.commandId, - threadId: childThreadId, - request: { - type: "provider-turn.interrupt", - providerSessionId, - providerThreadId: providerThread.id, - providerTurnId: turn.id, - }, - } satisfies PendingOrchestrationEffectV2, - ]); - }); - /** * Stop for a thread whatever it is doing, the way the Stop button stops the run it * shows: the running turn (or one whose background work is still pending) is interrupted, @@ -10390,9 +10232,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "thread.stop": cancelUnsettledEffects = yield* dispatchThreadStop(command, events, effects); break; - case "subagent.interrupt": - yield* dispatchSubagentInterrupt(command, events, effects); - break; case "queued-message.promote-to-steer": yield* dispatchQueuedMessagePromoteToSteer(command, events, effects); break; @@ -10693,93 +10532,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } satisfies OrchestratorV2DispatchResult; }); - const settleNativeSubagentParent = ( - command: Extract< - OrchestrationV2InternalCommand, - { readonly type: "thread.background-work.settle" } - >, - parentThreadId: ThreadId, - ) => - Effect.gen(function* () { - const child = yield* projectionStore.getThreadRecords(command.threadId, ["providerTurns"]); - const turn = child.providerTurns.find((candidate) => candidate.id === command.providerTurnId); - if ( - turn?.runAttemptId !== null || - turn.status !== "interrupted" || - !isProviderNativeSubagentThread(child.thread) || - child.providerTurns.some( - (candidate) => - candidate.providerThreadId === turn.providerThreadId && - candidate.ordinal > turn.ordinal, - ) - ) - return; - const parent = yield* projectionStore.getThreadRecords( - parentThreadId, - ["subagents", "nodes", "turnItems"], - { turnItemTypes: ["subagent"], turnItemStatuses: ["pending", "running", "waiting"] }, - ); - const task = parent.subagents.find( - (candidate) => - candidate.origin === "provider_native" && - candidate.driver === "codex" && - candidate.childThreadId === command.threadId && - isOrchestrationV2WorkActive(candidate.status), - ); - if (task === undefined) return; - const now = yield* DateTime.now; - yield* writeSystemEvents([ - { - type: "subagent.updated", - threadId: parentThreadId, - ...(task.runId === null ? {} : { runId: task.runId }), - nodeId: task.id, - occurredAt: now, - payload: { ...task, status: "interrupted", completedAt: now, updatedAt: now }, - }, - ...parent.nodes - .filter((node) => node.id === task.id && isOrchestrationV2WorkActive(node.status)) - .map((node) => ({ - type: "node.updated" as const, - threadId: parentThreadId, - ...(node.runId === null ? {} : { runId: node.runId }), - nodeId: node.id, - occurredAt: now, - payload: { ...node, status: "interrupted" as const, completedAt: now }, - })), - ...parent.turnItems - .filter((item) => item.type === "subagent" && item.subagentId === task.id) - .map((item) => ({ - type: "turn-item.updated" as const, - threadId: parentThreadId, - ...(item.runId === null ? {} : { runId: item.runId }), - ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), - occurredAt: now, - payload: { ...item, status: "interrupted" as const, completedAt: now, updatedAt: now }, - })), - ]); - }).pipe(mapDispatchError(command)); - const dispatchWithReceipt = (command: OrchestrationV2ServerCommand) => - Effect.gen(function* () { - const result = yield* threadDispatch.withLock( - commandThreadId(command), - dispatchWithReceiptEffect(command), - ); - if (command.type === "thread.background-work.settle") { - const child = yield* projectionStore - .getThread(command.threadId) - .pipe(mapDispatchError(command)); - const parentThreadId = child.lineage.parentThreadId; - if (isProviderNativeSubagentThread(child) && parentThreadId !== null) { - yield* threadDispatch.withLock( - parentThreadId, - settleNativeSubagentParent(command, parentThreadId), - ); - } - } - return result; - }); + threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command)); const handleTerminalRun = (stored: OrchestrationV2StoredEvent) => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index 7fd2095142c0..f47b4d6829e8 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -1684,7 +1684,7 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => { }), ); - it.effect("does not treat read state or subagent interrupt requests as thread activity", () => + it.effect("does not treat visited or marked-unread state as thread activity", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStore.ProjectionStoreV2; const createdAt = yield* DateTime.now; @@ -1750,21 +1750,6 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => { const markedUnread = yield* projectionStore.getThreadProjection(threadId); assert.isNull(markedUnread.thread.lastVisitedAt); assert.deepEqual(markedUnread.thread.updatedAt, createdAt); - - yield* projectionStore.apply({ - id: EventId.make("event:projection-read-state:subagent-interrupt"), - type: "subagent.interrupt-requested", - threadId, - occurredAt: DateTime.add(createdAt, { seconds: 3 }), - payload: NodeId.make("subagent:projection-read-state"), - }); - - const interrupted = yield* projectionStore.getThreadProjection(threadId); - assert.deepEqual(interrupted.thread.updatedAt, createdAt); - const shell = (yield* projectionStore.getShellSnapshot()).threads.find( - (entry) => entry.id === threadId, - ); - assert.deepEqual(shell?.updatedAt, createdAt); }), ); diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 1ac261753cdd..ec75937a9a62 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -649,8 +649,6 @@ export function applyToProjection( }; switch (event.type) { - case "subagent.interrupt-requested": - return projection; case "thread.created": case "thread.archived": case "thread.unarchived": @@ -1755,8 +1753,6 @@ export const layer: Layer.Layer = const apply: ProjectionStoreV2Shape["apply"] = (event) => Effect.gen(function* () { switch (event.type) { - case "subagent.interrupt-requested": - break; case "thread.created": case "thread.archived": case "thread.unarchived": @@ -2589,7 +2585,6 @@ export const layer: Layer.Layer = } if ( - event.type !== "subagent.interrupt-requested" && event.type !== "thread.created" && event.type !== "thread.archived" && event.type !== "thread.unarchived" && diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts index b94fb3099727..ebf6243d357d 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts @@ -591,20 +591,6 @@ export const layer: Layer.Layer< .pipe(Effect.mapError(mapWriteError)); return result.storedEvents; }).pipe( - (write) => { - const event = input.event; - const parentThreadId = - event.type === "subagent.updated" - ? event.subagent.threadId - : event.type === "node.updated" && event.node.kind === "subagent" - ? event.node.threadId - : event.type === "turn_item.updated" && event.turnItem.type === "subagent" - ? event.turnItem.threadId - : undefined; - return parentThreadId === undefined - ? write - : threadCommands.withLock(parentThreadId, write); - }, Effect.flatMap((storedEvents) => storedEvents.length === 0 || input.event.type !== "subagent.updated" ? Effect.succeed(storedEvents) diff --git a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts index 0b59a927a8b4..0bb2d5053b4b 100644 --- a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts +++ b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts @@ -167,7 +167,6 @@ function commandThreadIds(command: OrchestrationV2Command): ReadonlyArray { } }); - it("publishes subagent interrupts with a node ID payload", () => { - assert.isTrue( - AgentAwarenessRelay.shouldPublishAgentAwarenessEvent({ - type: "subagent.interrupt-requested", - payload: NodeId.make("native-subagent"), - }), - ); - }); - it("does not publish imported thread creation as new agent activity", () => { assert.isFalse( AgentAwarenessRelay.shouldPublishAgentAwarenessEvent({ diff --git a/apps/server/src/relay/AgentAwarenessRelay.ts b/apps/server/src/relay/AgentAwarenessRelay.ts index b94ccbd81456..f4d3e7b57723 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.ts @@ -105,7 +105,6 @@ export function shouldPublishAgentAwarenessEvent( case "runtime-request.updated": case "subagent.updated": case "provider-thread.updated": - case "subagent.interrupt-requested": return true; case "thread.settled": case "thread.unsettled": diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx index a87376a00d36..2a5afdebc816 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx @@ -47,131 +47,109 @@ afterEach(async () => { state.projection = null; }); -it("stops only active subagents from lineage without opening their thread", async () => { - vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); - const parent = { - id: "parent", - lineage: { relationshipToParent: null }, - activeProviderThreadId: null, - }; - const child = { - id: "child", - title: "Worker", - lineage: { parentThreadId: "parent", relationshipToParent: "subagent" }, - }; - const agent = { - id: "agent", - childThreadId: "child", - origin: "app_owned", - driver: "codex", - providerInstanceId: "codex", - title: "Worker", - prompt: "Check the change", - model: "gpt-5.4", - status: "running", - progress: null, - result: null, - startedAt: DateTime.makeUnsafe("2026-09-16T12:00:00Z"), - completedAt: null, - updatedAt: DateTime.makeUnsafe("2026-09-16T12:00:00Z"), - }; - state.shells = [{ environmentId: "test", source: child }]; - const projection = { - thread: parent, - runs: [], - providerThreads: [], - providerSessions: [], - contextTransfers: [], - subagents: [agent], - }; - state.projection = projection; - const panel = ( - - ); - await act(async () => { - renderer = create(panel); - }); - const stopButton = () => renderer.root.findByProps({ "aria-label": "Stop subagent Worker" }); - await act(async () => stopButton().props.onClick()); - expect(state.command).toHaveBeenCalledWith({ - environmentId: "test", - input: { threadId: "child" }, - }); - expect(state.navigate).not.toHaveBeenCalled(); - - state.command.mockClear(); - state.projection = { ...projection, subagents: [{ ...agent, origin: "provider_native" }] }; - for (const supported of [undefined, false, true]) { - state.configs.set("test", { - environment: { - capabilities: supported === undefined ? {} : { subagentInterrupt: supported }, - }, +it.each(["codex", "claudeAgent"])( + "stops only active app-owned %s subagents without opening their thread", + async (driver) => { + vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); + const parent = { + id: "parent", + lineage: { relationshipToParent: null }, + activeProviderThreadId: null, + }; + const child = { + id: "child", + title: "Worker", + lineage: { parentThreadId: "parent", relationshipToParent: "subagent" }, + }; + const agent = { + id: "agent", + childThreadId: "child", + origin: "app_owned", + driver, + providerInstanceId: "codex", + title: "Worker", + prompt: "Check the change", + model: "gpt-5.4", + status: "running", + progress: null, + result: null, + startedAt: DateTime.makeUnsafe("2026-09-16T12:00:00Z"), + completedAt: null, + updatedAt: DateTime.makeUnsafe("2026-09-16T12:00:00Z"), + }; + state.shells = [{ environmentId: "test", source: child }]; + const projection = { + thread: parent, + runs: [], + providerThreads: [], + providerSessions: [], + contextTransfers: [], + subagents: [agent], + }; + state.projection = projection; + const panel = ( + + ); + await act(async () => { + renderer = create(panel); }); - await act(async () => renderer.update(cloneElement(panel))); - expect( - renderer.root.findAll( - (node) => node.type === "button" && node.props["aria-label"] === "Stop subagent Worker", - ), - ).toHaveLength(supported === true ? 1 : 0); - } - await act(async () => stopButton().props.onClick()); - expect(state.command).toHaveBeenCalledWith({ - environmentId: "test", - input: { threadId: "parent", subagentId: "agent" }, - }); - - state.command.mockClear(); - state.projection = { - ...projection, - subagents: [{ ...agent, origin: "provider_native", status: "waiting" }], - }; - await act(async () => renderer.update(cloneElement(panel))); - await act(async () => stopButton().props.onClick()); - expect(state.command).toHaveBeenCalledWith({ - environmentId: "test", - input: { threadId: "parent", subagentId: "agent" }, - }); - expect(state.navigate).not.toHaveBeenCalled(); + const stopButton = () => renderer.root.findByProps({ "aria-label": "Stop subagent Worker" }); + await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenCalledWith({ + environmentId: "test", + input: { threadId: "child" }, + }); + expect(state.navigate).not.toHaveBeenCalled(); - for (const status of ["running", "waiting"] as const) { - state.shells = [ - { - environmentId: "test", - source: { - ...child, - activityRunStatus: status, - activityRunStartedAt: DateTime.makeUnsafe("2026-09-16T12:05:00Z"), + for (const status of ["starting", "running", "waiting"] as const) { + state.shells = [ + { + environmentId: "test", + source: { + ...child, + activityRunStatus: status, + activityRunStartedAt: DateTime.makeUnsafe("2026-09-16T12:05:00Z"), + }, }, - }, - ]; + ]; + state.projection = { + ...projection, + subagents: [{ ...agent, origin: "provider_native", status: "completed" }], + }; + await act(async () => renderer.update(cloneElement(panel))); + expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength( + 0, + ); + state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; + await act(async () => renderer.update(cloneElement(panel))); + await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenLastCalledWith({ + environmentId: "test", + input: { threadId: "child" }, + }); + } + state.shells = [{ environmentId: "test", source: child }]; + for (const status of ["completed", "failed", "interrupted"]) { + state.projection = { ...projection, subagents: [{ ...agent, status }] }; + await act(async () => renderer.update(cloneElement(panel))); + expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength( + 0, + ); + } + state.projection = { ...projection, subagents: [{ ...agent, startedAt: null }] }; + await act(async () => renderer.update(cloneElement(panel))); + expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); state.projection = { ...projection, - subagents: [{ ...agent, origin: "provider_native", status: "completed" }], + subagents: [{ ...agent, origin: "provider_native", driver: "claudeAgent" }], }; await act(async () => renderer.update(cloneElement(panel))); expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); - state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; - await act(async () => renderer.update(cloneElement(panel))); - await act(async () => stopButton().props.onClick()); - expect(state.command).toHaveBeenLastCalledWith({ - environmentId: "test", - input: { threadId: "child" }, - }); - } - state.shells = [{ environmentId: "test", source: child }]; - state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; - await act(async () => renderer.update(cloneElement(panel))); - expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); - state.projection = { - ...projection, - subagents: [{ ...agent, origin: "provider_native", driver: "claudeAgent" }], - }; - await act(async () => renderer.update(cloneElement(panel))); - expect(renderer.root.findAllByProps({ "aria-label": "Stop subagent Worker" })).toHaveLength(0); -}); + }, +); it("shows the matching child agent details and refreshes them when the agent settles", async () => { vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); @@ -340,7 +318,6 @@ it("shows readable models and only differing workspace details in agent tooltips ]; state.shells = [{ environmentId: "test", source: child }]; state.configs.set("test", { - environment: { capabilities: {} }, providers: [ { instanceId: "codex", diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx index 1400ffa11a19..85f4f7bb8a95 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.tsx @@ -24,12 +24,7 @@ import { canDetachThreadProviderSession, resolveLatestMergeBackRun, } from "@t3tools/client-runtime/state/thread-workflows"; -import type { - EnvironmentId, - NodeId, - OrchestrationV2ThreadShell, - ThreadId, -} from "@t3tools/contracts"; +import type { EnvironmentId, OrchestrationV2ThreadShell, ThreadId } from "@t3tools/contracts"; import { groupBy } from "effect/Array"; import * as DateTime from "effect/DateTime"; import { useNavigate } from "@tanstack/react-router"; @@ -203,9 +198,7 @@ export function ThreadRelationshipsPanel(props: { }) { const ref = scopeThreadRef(props.environmentId, props.threadId); const projection = useThreadProjection(ref)?.projection ?? null; - const config = useServerConfigs().get(props.environmentId); - const providers = config?.providers; - const supportsSubagentInterrupt = config?.environment.capabilities.subagentInterrupt === true; + const providers = useServerConfigs().get(props.environmentId)?.providers; const subagentsByThreadId = useMemo( () => new Map( @@ -215,7 +208,6 @@ export function ThreadRelationshipsPanel(props: { subagent.childThreadId, { ...projectedSubagentsToRuntime([subagent])[0]!, - subagentId: subagent.id, origin: subagent.origin, driver: subagent.driver, providerInstanceId: subagent.providerInstanceId, @@ -245,9 +237,8 @@ export function ThreadRelationshipsPanel(props: { const mergeBack = useAtomCommand(threadEnvironment.mergeBack); const stopSession = useAtomCommand(threadEnvironment.stopSession); const interruptTurn = useAtomCommand(threadEnvironment.interruptTurn); - const interruptSubagent = useAtomCommand(threadEnvironment.interruptSubagent); const [busyAction, setBusyAction] = useState<"merge" | "detach" | null>(null); - const [stoppingSubagentId, setStoppingSubagentId] = useState(null); + const [stoppingThreadId, setStoppingThreadId] = useState(null); const latestMergeBackRun = projection === null ? null : resolveLatestMergeBackRun(projection); const mergeTargetThreadId = resolveMergeBackTargetThreadId(projection); const relationshipRows = useMemo( @@ -323,24 +314,14 @@ export function ThreadRelationshipsPanel(props: { setBusyAction(null); }; - const stopSubagent = async ( - subagentId: NodeId, - origin: "app_owned" | "provider_native", - childThreadId: ThreadId, - ) => { - if (stoppingSubagentId !== null) return; - setStoppingSubagentId(subagentId); - const result = - origin === "app_owned" - ? await interruptTurn({ - environmentId: props.environmentId, - input: { threadId: childThreadId }, - }) - : await interruptSubagent({ - environmentId: props.environmentId, - input: { threadId: props.threadId, subagentId }, - }); - setStoppingSubagentId(null); + const stopSubagent = async (childThreadId: ThreadId) => { + if (stoppingThreadId !== null) return; + setStoppingThreadId(childThreadId); + const result = await interruptTurn({ + environmentId: props.environmentId, + input: { threadId: childThreadId }, + }); + setStoppingThreadId(null); if (result._tag === "Failure") { toastManager.add({ type: "error", title: "Could not stop subagent" }); } @@ -397,18 +378,14 @@ export function ThreadRelationshipsPanel(props: { ? BotIcon : GitForkIcon; const relationship = relationshipLabel(edge, props.threadId); - const storedAgent = - isSubagent && !isParent ? subagentsByThreadId.get(threadId) : undefined; - const agent = liveSubagent(storedAgent, node?.thread); + const agent = liveSubagent( + isSubagent && !isParent ? subagentsByThreadId.get(threadId) : undefined, + node?.thread, + ); const canStop = - agent?.startedAt && - ((agent.origin === "app_owned" && - ["pending", "running", "waiting"].includes(agent.status)) || - (supportsSubagentInterrupt && - agent.origin === "provider_native" && - agent.driver === "codex" && - (storedAgent?.status === "running" || storedAgent?.status === "waiting") && - (agent.status === "running" || agent.status === "waiting"))); + agent?.origin === "app_owned" && + agent.startedAt && + ["pending", "running", "waiting"].includes(agent.status); const threadTitle = relationshipThreadTitle({ title: node?.thread?.title ?? agent?.title ?? threadId, isSubagent, @@ -459,7 +436,7 @@ export function ThreadRelationshipsPanel(props: { {agent ? ( agent.startedAt ? ( @@ -555,7 +532,7 @@ export function ThreadRelationshipsPanel(props: { )} {canStop && agent ? ( -
    +
    - void stopSubagent(agent.subagentId, agent.origin, threadId) - } + disabled={stoppingThreadId !== null} + onClick={() => void stopSubagent(threadId)} /> } > - {stoppingSubagentId === agent.subagentId ? ( + {stoppingThreadId === threadId ? ( ) : ( diff --git a/packages/client-runtime/src/operations/commands.test.ts b/packages/client-runtime/src/operations/commands.test.ts index 8c9a547b9fb8..52e08fe6e74e 100644 --- a/packages/client-runtime/src/operations/commands.test.ts +++ b/packages/client-runtime/src/operations/commands.test.ts @@ -25,7 +25,6 @@ import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; -import * as Result from "effect/Result"; import * as SubscriptionRef from "effect/SubscriptionRef"; import { @@ -45,7 +44,6 @@ import { editQueuedRun, forkThreadFromRun, interruptThreadTurn, - interruptSubagent, mergeThreadBack, promoteQueuedRun, reorderActiveThread, @@ -80,7 +78,6 @@ const makeSupervisor = Effect.fn("TestEnvironmentCommands.makeSupervisor")(funct readonly projection?: OrchestrationV2ThreadProjection; readonly projectionRequests?: ThreadId[]; readonly advertiseServerResolvedCommandContext?: boolean; - readonly subagentInterrupt?: boolean; }) { const client = { [ORCHESTRATION_V2_WS_METHODS.dispatchCommand]: (command: OrchestrationV2Command) => @@ -128,9 +125,6 @@ const makeSupervisor = Effect.fn("TestEnvironmentCommands.makeSupervisor")(funct environment: { capabilities: { repositoryIdentity: true, - ...(input.subagentInterrupt === undefined - ? {} - : { subagentInterrupt: input.subagentInterrupt }), ...(input.advertiseServerResolvedCommandContext === false ? {} : { serverResolvedCommandContext: true }), @@ -154,35 +148,6 @@ const makeSupervisor = Effect.fn("TestEnvironmentCommands.makeSupervisor")(funct }); describe("V2 environment commands", () => { - it.effect.each([undefined, false, true])( - "dispatches a subagent interrupt only when support is advertised as %s", - (supported) => - Effect.gen(function* () { - const commands: OrchestrationV2Command[] = []; - const supervisor = yield* makeSupervisor({ - commands, - projects: [], - ...(supported === undefined ? {} : { subagentInterrupt: supported }), - }); - const commandId = CommandId.make("subagent-interrupt"); - const subagentId = NodeId.make("subagent-1"); - const result = yield* interruptSubagent({ - commandId, - threadId: v2ThreadId, - subagentId, - }).pipe( - Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), - Effect.result, - ); - expect(Result.isSuccess(result)).toBe(supported === true); - expect(commands).toEqual( - supported === true - ? [{ type: "subagent.interrupt", commandId, threadId: v2ThreadId, subagentId }] - : [], - ); - }).pipe(Effect.provide(layerTestCrypto)), - ); - it.effect("routes projects through the event-sourced project transport", () => Effect.gen(function* () { const projects: ProjectMutation[] = []; diff --git a/packages/client-runtime/src/operations/commands.ts b/packages/client-runtime/src/operations/commands.ts index c25ad3bb0675..84f4a2ad3cff 100644 --- a/packages/client-runtime/src/operations/commands.ts +++ b/packages/client-runtime/src/operations/commands.ts @@ -6,11 +6,9 @@ import { CheckpointScopeId, ORCHESTRATION_V2_WS_METHODS, OrchestrationV2CheckpointUnavailableError, - OrchestrationV2DispatchCommandError, WS_METHODS, type ChatAttachment, type MessageId, - type NodeId, type ModelSelection, type OrchestrationV2Command, type OrchestrationV2CreationSource, @@ -189,10 +187,6 @@ export interface InterruptThreadTurnInput extends ThreadCommandInput { readonly turnId?: string; } -export interface InterruptSubagentInput extends ThreadCommandInput { - readonly subagentId: NodeId; -} - export interface RespondToThreadApprovalInput extends ThreadCommandInput { readonly requestId: RuntimeRequestId; readonly decision: ProviderApprovalDecision; @@ -837,26 +831,6 @@ export const interruptThreadTurn = Effect.fn("EnvironmentCommands.interruptThrea }); }); -export const interruptSubagent = Effect.fn("EnvironmentCommands.interruptSubagent")(function* ( - input: InterruptSubagentInput, -) { - const config = yield* getInitialServerConfig(); - const commandId = yield* allocateCommandId(input); - if (config.environment.capabilities.subagentInterrupt !== true) { - return yield* new OrchestrationV2DispatchCommandError({ - commandId, - commandType: "subagent.interrupt", - message: "This server does not support stopping native subagents.", - }); - } - return yield* dispatch({ - type: "subagent.interrupt", - commandId, - threadId: input.threadId, - subagentId: input.subagentId, - }); -}); - export const respondToThreadApproval = Effect.fn("EnvironmentCommands.respondToThreadApproval")( function* (input: RespondToThreadApprovalInput) { return yield* dispatch({ diff --git a/packages/client-runtime/src/state/orchestrationV2Projection.test.ts b/packages/client-runtime/src/state/orchestrationV2Projection.test.ts index b6690c57c235..1a53afb8d7e8 100644 --- a/packages/client-runtime/src/state/orchestrationV2Projection.test.ts +++ b/packages/client-runtime/src/state/orchestrationV2Projection.test.ts @@ -105,19 +105,6 @@ const emptyProjection = { } as OrchestrationV2ThreadProjection; describe("applyOrchestrationV2ProjectionEvent", () => { - it("leaves subagent state unchanged when a stop is requested", () => { - const event = { - id: "event-subagent-stop", - type: "subagent.interrupt-requested", - threadId, - nodeId: NodeId.make("subagent-stop"), - occurredAt: DateTime.makeUnsafe("2026-06-20T01:00:00.000Z"), - payload: NodeId.make("subagent-stop"), - } as OrchestrationV2DomainEvent; - - expect(applyOrchestrationV2ProjectionEvent(emptyProjection, event)).toBe(emptyProjection); - }); - it("keeps live token usage when the terminal provider turn omits it", () => { const providerTurnId = ProviderTurnId.make("provider-turn-reducer"); const running = { diff --git a/packages/client-runtime/src/state/orchestrationV2Projection.ts b/packages/client-runtime/src/state/orchestrationV2Projection.ts index 8e406901b207..7e96a3039353 100644 --- a/packages/client-runtime/src/state/orchestrationV2Projection.ts +++ b/packages/client-runtime/src/state/orchestrationV2Projection.ts @@ -161,8 +161,6 @@ export function applyOrchestrationV2ProjectionEvent( const latestLocalTurnOrdinal = options?.latestLocalTurnOrdinal; const base = { ...projection, updatedAt: event.occurredAt }; switch (event.type) { - case "subagent.interrupt-requested": - return projection; case "thread.created": case "thread.archived": case "thread.unarchived": diff --git a/packages/client-runtime/src/state/threadCommands.ts b/packages/client-runtime/src/state/threadCommands.ts index 34c7063e9d37..f237f1ce943d 100644 --- a/packages/client-runtime/src/state/threadCommands.ts +++ b/packages/client-runtime/src/state/threadCommands.ts @@ -26,7 +26,6 @@ import { type DeleteThreadInput, type EditQueuedRunInput, type InterruptThreadTurnInput, - type InterruptSubagentInput, type MarkThreadUnreadInput, type ForkThreadFromRunInput, type MergeThreadBackInput, @@ -61,7 +60,6 @@ import { deleteThread, editQueuedRun, interruptThreadTurn, - interruptSubagent, forkThreadFromRun, markThreadUnread, mergeThreadBack, @@ -288,12 +286,6 @@ export function createThreadEnvironmentAtoms( scheduler, concurrency, }), - interruptSubagent: createEnvironmentCommand(runtime, { - label: "environment-data:commands:thread:interrupt-subagent", - execute: (input: InterruptSubagentInput) => interruptSubagent(input), - scheduler, - concurrency, - }), respondToApproval: createEnvironmentCommand(runtime, { label: "environment-data:commands:thread:respond-to-approval", execute: (input: RespondToThreadApprovalInput) => respondToThreadApproval(input), diff --git a/packages/contracts/src/environment.test.ts b/packages/contracts/src/environment.test.ts index 41c34af1341f..96e176bccba1 100644 --- a/packages/contracts/src/environment.test.ts +++ b/packages/contracts/src/environment.test.ts @@ -14,18 +14,6 @@ const descriptor = { } as const; describe("ExecutionEnvironmentDescriptor", () => { - it.each([undefined, false, true])("preserves subagent interrupt support as %s", (supported) => { - expect( - decodeDescriptor({ - ...descriptor, - capabilities: { - ...descriptor.capabilities, - ...(supported === undefined ? {} : { subagentInterrupt: supported }), - }, - }).capabilities.subagentInterrupt, - ).toBe(supported); - }); - it("decodes old, recognized and future manual installation descriptors", () => { expect(decodeDescriptor(descriptor).capabilities.serverInstallation).toBeUndefined(); for (const installation of [{ kind: "npx" }, { kind: "npm-global", prefix: "/opt/node" }]) { diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index 705f7bf1a46f..c9ac7482d4cc 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -98,7 +98,6 @@ export type ServerSelfUpdateCapability = typeof ServerSelfUpdateCapability.Type; export const ExecutionEnvironmentCapabilities = Schema.Struct({ repositoryIdentity: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), connectionProbe: Schema.optionalKey(Schema.Boolean), - subagentInterrupt: Schema.optionalKey(Schema.Boolean), /** Missing on older servers, which still accept inline image attachments. */ attachmentUploads: Schema.optionalKey(Schema.Boolean), /** Uploaded files may accompany question answers. */ diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 2d6cc32dbfdd..c26d5a9327e4 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -1706,11 +1706,6 @@ export const OrchestrationV2DomainEvent = Schema.Union([ type: Schema.Literal("subagent.updated"), payload: OrchestrationV2Subagent, }), - Schema.Struct({ - ...OrchestrationV2EventBase.fields, - type: Schema.Literal("subagent.interrupt-requested"), - payload: NodeId, - }), Schema.Struct({ ...OrchestrationV2EventBase.fields, type: Schema.Literals(["provider-session.attached", "provider-session.updated"]), @@ -2530,11 +2525,6 @@ export const OrchestrationV2DomainEventJson = Schema.Union([ type: Schema.Literal("subagent.updated"), payload: OrchestrationV2SubagentJson, }), - Schema.Struct({ - ...OrchestrationV2JsonEventBaseFields, - type: Schema.Literal("subagent.interrupt-requested"), - payload: NodeId, - }), Schema.Struct({ ...OrchestrationV2JsonEventBaseFields, type: Schema.Literals(["provider-session.attached", "provider-session.updated"]), @@ -2927,12 +2917,6 @@ export const OrchestrationV2Command = Schema.Union([ */ holdQueue: Schema.optional(Schema.Boolean), }), - Schema.Struct({ - type: Schema.Literal("subagent.interrupt"), - commandId: CommandId, - threadId: ThreadId, - subagentId: NodeId, - }), Schema.Struct({ type: Schema.Literal("queued-message.promote-to-steer"), commandId: CommandId, From ed876c3779793aee0c631bd42f1752a68e965547 Mon Sep 17 00:00:00 2001 From: Bil0000 Date: Tue, 6 Oct 2026 23:14:08 +0300 Subject: [PATCH 11/11] test(web): check each Lineage stop call independently --- .../components/chat/ThreadRelationshipsControl.agents.test.tsx | 2 ++ 1 file changed, 2 insertions(+) diff --git a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx index d347061ee496..6c7534d04884 100644 --- a/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx +++ b/apps/web/src/components/chat/ThreadRelationshipsControl.agents.test.tsx @@ -111,6 +111,7 @@ it.each(["codex", "claudeAgent"])( expect(state.navigate).not.toHaveBeenCalled(); for (const status of ["starting", "running", "waiting"] as const) { + state.command.mockClear(); state.shells = [ { environmentId: "test", @@ -132,6 +133,7 @@ it.each(["codex", "claudeAgent"])( state.projection = { ...projection, subagents: [{ ...agent, status: "completed" }] }; await act(async () => renderer.update(cloneElement(panel))); await act(async () => stopButton().props.onClick()); + expect(state.command).toHaveBeenCalledTimes(1); expect(state.command).toHaveBeenLastCalledWith({ environmentId: "test", input: { threadId: "child" },