diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index 0dbb2f571683..cfca51b629c8 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -767,17 +767,13 @@ describe("orchestrator MCP toolkit", () => { const invoke = (name: string, args: Record) => invokeAs(invocation, name, args); - const refusedSettle = yield* invoke("t3_thread_organize", { action: "settle" }); - expect(refusedSettle.isError).toBe(true); - expect(refusedSettle.structuredContent).toBeUndefined(); - expect(declaredFailure(refusedSettle)).toEqual({ - _tag: "OrchestratorMcpFailure", - code: "orchestration_error", - message: `Thread ${parentThreadId} has active or blocked work and cannot be settled.`, - }); - const afterRefusedSettle = yield* orchestrator.getThreadProjection(parentThreadId); - expect(afterRefusedSettle.thread.settledOverride).not.toBe("settled"); - expect(afterRefusedSettle.runs.find((run) => run.id === parentRun?.id)?.status).toBe( + // Settling would stop the session, so the agent's own turn keeps running. + const deferredSettle = yield* invoke("t3_thread_organize", { action: "settle" }); + expect(deferredSettle.isError).toBe(false); + expect(deferredSettle.structuredContent).toEqual({ settlesWhenTurnEnds: true }); + const afterDeferredSettle = yield* orchestrator.getThreadProjection(parentThreadId); + expect(afterDeferredSettle.thread.settledOverride).not.toBe("settled"); + expect(afterDeferredSettle.runs.find((run) => run.id === parentRun?.id)?.status).toBe( "running", ); diff --git a/apps/server/src/mcp/toolkits/thread/handlers.ts b/apps/server/src/mcp/toolkits/thread/handlers.ts index 2a06c58c1b43..f54065c0824e 100644 --- a/apps/server/src/mcp/toolkits/thread/handlers.ts +++ b/apps/server/src/mcp/toolkits/thread/handlers.ts @@ -292,8 +292,13 @@ export const layer = McpToolAccess.toLayer(ThreadToolkit, { ), t3_thread_organize: writesThread((input) => Effect.gen(function* () { - const { threads, projection } = yield* readThread(input.threadId); + const { threads, projection, caller } = yield* readThread(input.threadId); const common = { commandId: yield* newCommandId(), threadId: projection.thread.id }; + if (input.action === "settle") { + return yield* threads + .settleThread({ ...common, byOwnAgent: caller?.id === projection.thread.id }) + .pipe(Effect.mapError(dispatchFailure)); + } let command: OrchestrationV2Command; switch (input.action) { case "snooze": diff --git a/apps/server/src/mcp/toolkits/thread/tools.ts b/apps/server/src/mcp/toolkits/thread/tools.ts index 0958e6f3add9..24cd6af444ba 100644 --- a/apps/server/src/mcp/toolkits/thread/tools.ts +++ b/apps/server/src/mcp/toolkits/thread/tools.ts @@ -30,7 +30,7 @@ import * as McpInvocationContext from "../../McpInvocationContext.ts"; const ThreadOrganizeTool = Tool.make("t3_thread_organize", { description: - "Pin, snooze, settle, archive, or mark a thread unread. Omit threadId for this thread. snooze requires snoozedUntil. Existing thread lifecycle rules apply; this does not schedule a future action.", + "Pin, snooze, settle, archive, or mark a thread unread. Omit threadId for this thread. snooze requires snoozedUntil. Existing thread lifecycle rules apply. Settling this thread takes effect when your turn completes, returning settlesWhenTurnEnds=true; a turn that fails or is interrupted, or a queued message, leaves it active.", parameters: Schema.Struct({ threadId: Schema.optional(ThreadId), action: Schema.Literals([ @@ -46,7 +46,10 @@ const ThreadOrganizeTool = Tool.make("t3_thread_organize", { ]), snoozedUntil: Schema.optional(IsoDateTime), }), - success: OrchestrationV2DispatchCommandResult, + success: Schema.Union([ + OrchestrationV2DispatchCommandResult, + Schema.Struct({ settlesWhenTurnEnds: Schema.Literal(true) }), + ]), failure: OrchestratorMcpFailure, failureMode: "return" as const, dependencies: [ diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.test.ts b/apps/server/src/orchestration-v2/ThreadManagementService.test.ts index aa2d823227c4..94e17bdc567e 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.test.ts @@ -473,3 +473,40 @@ it.effect("waitForThread reads the run again only when the run updates", () => expect(reads).toBe(2); }), ); + +it.effect.each([ + { status: "completed" as const, settles: true }, + { status: "failed" as const, settles: false }, + { status: "interrupted" as const, settles: false }, +])("settleAfterRun settles only when the run $status", ({ status, settles }) => + Effect.gen(function* () { + const projectId = ProjectId.make("project:thread-management:settle-after-run"); + const threadId = ThreadId.make("thread:thread-management:settle-after-run"); + const runId = RunId.make("run:thread-management:settle-after-run"); + const dispatched: Array = []; + const layerTest = ThreadManagementService.layer.pipe( + Layer.provide( + Layer.mock(Orchestrator.OrchestratorV2)({ + getThreadEventSequence: () => Effect.succeed(0), + getThreadRecords: () => + Effect.succeed({ + thread: { id: threadId, projectId, deletedAt: null }, + runs: [{ id: runId, status }], + } as unknown as OrchestrationV2ThreadProjection), + dispatch: (command) => + Effect.sync(() => { + dispatched.push(command.type); + return { sequence: 1, storedEvents: [] } as never; + }), + }), + ), + ); + const service = yield* ThreadManagementService.ThreadManagementService.pipe( + Effect.provide(layerTest), + ); + + yield* service.settleAfterRun({ projectId, threadId, runId }); + + expect(dispatched).toEqual(settles ? ["thread.settle"] : []); + }), +); diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.ts b/apps/server/src/orchestration-v2/ThreadManagementService.ts index b48beb08875c..e59fce35bcbd 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.ts @@ -326,6 +326,30 @@ export interface ThreadManagementServiceShape { readonly waitForThread: ( input: ThreadManagementWaitInput, ) => Effect.Effect; + /** + * Waits for `runId` to end, then settles the thread if the run completed. + * Settling stops the provider session, so an agent that settles its own + * thread uses this to wait for its turn to end. A run that ends any other + * way leaves the thread active, as does anything `thread.settle` refuses + * then, such as a queued message. + */ + readonly settleAfterRun: (input: { + readonly projectId: ProjectId; + readonly threadId: ThreadId; + readonly runId: RunId; + }) => Effect.Effect; + /** + * Settles a thread. When the thread's own agent asks, it is mid-turn, so + * the thread settles through `settleAfterRun` once that turn ends. + */ + readonly settleThread: (input: { + readonly threadId: ThreadId; + readonly commandId: CommandId; + readonly byOwnAgent: boolean; + }) => Effect.Effect< + { readonly sequence: number } | { readonly settlesWhenTurnEnds: true }, + Orchestrator.OrchestratorV2Error + >; readonly interruptThread: ( input: ThreadManagementInterruptInput, ) => Effect.Effect; @@ -402,8 +426,11 @@ function latestSteerableRun( .toSorted((left, right) => right.ordinal - left.ordinal)[0]; } +const SETTLE_AFTER_RUN_WAIT_MS = 24 * 60 * 60 * 1_000; + const make = Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; + const layerScope = yield* Effect.scope; const legacyImporter = yield* LegacyV1ThreadImporter.LegacyV1ThreadImporter; const ensureLegacyTranscript = Effect.fn( @@ -724,6 +751,39 @@ const make = Effect.gen(function* () { return { threadId: input.threadId, run, timedOut: !isTerminalRunStatus(run.status) }; }); + const settleAfterRun: ThreadManagementServiceShape["settleAfterRun"] = Effect.fn( + "orchestrationV2.threadManagement.settleAfterRun", + )(function* (input) { + // A turn that outlives this wait leaves its thread active. + const { run } = yield* waitForThread({ ...input, timeoutMs: SETTLE_AFTER_RUN_WAIT_MS }); + if (run?.status !== "completed") return; + yield* dispatch({ + type: "thread.settle", + commandId: CommandId.make(`server:settle-after-run:${input.runId}`), + threadId: input.threadId, + }); + }); + + const settleThread: ThreadManagementServiceShape["settleThread"] = Effect.fn( + "orchestrationV2.threadManagement.settleThread", + )(function* (input) { + const shell = input.byOwnAgent ? yield* orchestrator.getThreadShell(input.threadId) : null; + if (shell != null && shell.activeRunId !== null) { + yield* settleAfterRun({ + projectId: shell.projectId, + threadId: shell.id, + runId: shell.activeRunId, + }).pipe(Effect.ignoreCause({ log: true }), Effect.forkIn(layerScope)); + return { settlesWhenTurnEnds: true } as const; + } + const result = yield* dispatch({ + type: "thread.settle", + commandId: input.commandId, + threadId: input.threadId, + }); + return { sequence: result.sequence }; + }); + const interruptThread: ThreadManagementServiceShape["interruptThread"] = (input) => Effect.gen(function* () { const target = yield* getProjectThreadRecords(input, ["runs", "providerTurns"]); @@ -829,6 +889,8 @@ const make = Effect.gen(function* () { listProjectThreads, sendToThread, waitForThread, + settleAfterRun, + settleThread, interruptThread, stopDelegatedTasks, getThreadEventSequence: orchestrator.getThreadEventSequence, diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 131f05034051..79411e6fcc55 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -3373,6 +3373,74 @@ it.layer(layerTest)("OrchestrationV2LayerLive lifecycle", (it) => { }), ); + it.effect("settles a thread its own agent settled once the turn completes", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const eventSink = yield* EventSink.EventSinkV2; + const threadManagement = yield* ThreadManagementService.ThreadManagementService; + const projectId = ProjectId.make("runtime-layer-settle-after-run-project"); + const threadId = ThreadId.make("runtime-layer-settle-after-run-thread"); + + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("runtime-layer-settle-after-run-create"), + threadId, + projectId, + title: "Settle after run", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: "/tmp/runtime-layer-settle-after-run", + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("runtime-layer-settle-after-run-message"), + threadId, + messageId: MessageId.make("runtime-layer-settle-after-run-message"), + text: "Fix it and then settle this thread.", + attachments: [], + modelSelection, + dispatchMode: { type: "start_immediately" }, + }); + const run = (yield* orchestrator.getThreadProjection(threadId)).runs[0]; + if (run === undefined) return yield* Effect.die(new Error("Run missing.")); + + const settled = yield* orchestrator + .streamStoredEventsFrom({ threadId, afterSequence: 0, eventType: "thread.settled" }) + .pipe(Stream.runHead, Effect.forkChild); + const result = yield* threadManagement.settleThread({ + threadId, + commandId: CommandId.make("runtime-layer-settle-after-run-settle"), + byOwnAgent: true, + }); + assert.deepEqual(result, { settlesWhenTurnEnds: true }); + assert.isNull((yield* orchestrator.getThreadProjection(threadId)).thread.settledOverride); + const now = yield* DateTime.now; + yield* eventSink.write({ + commandId: CommandId.make("runtime-layer-settle-after-run-completed"), + events: [ + { + id: EventId.make("runtime-layer-settle-after-run-completed"), + type: "run.updated", + threadId, + runId: run.id, + occurredAt: now, + payload: { ...run, status: "completed", startedAt: now, completedAt: now }, + }, + ], + }); + yield* Fiber.join(settled); + + const projection = yield* orchestrator.getThreadProjection(threadId); + assert.equal(projection.thread.settledOverride, "settled"); + }), + ); + it.effect("settles past held automatic runs but not held user messages", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index ee37b9d5a5ae..aa64f048ad80 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -206,6 +206,8 @@ const makeTestRelay = Effect.fnUntraced(function* ( listProjectThreads: unused, sendToThread: unused, waitForThread: unused, + settleAfterRun: unused, + settleThread: unused, interruptThread: unused, stopDelegatedTasks: unused, getThreadEventSequence: unused,