Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 7 additions & 11 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -767,17 +767,13 @@ describe("orchestrator MCP toolkit", () => {
const invoke = (name: string, args: Record<string, unknown>) =>
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",
);

Expand Down
7 changes: 6 additions & 1 deletion apps/server/src/mcp/toolkits/thread/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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":
Expand Down
7 changes: 5 additions & 2 deletions apps/server/src/mcp/toolkits/thread/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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([
Expand All @@ -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: [
Expand Down
37 changes: 37 additions & 0 deletions apps/server/src/orchestration-v2/ThreadManagementService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string> = [];
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"] : []);
}),
);
62 changes: 62 additions & 0 deletions apps/server/src/orchestration-v2/ThreadManagementService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -326,6 +326,30 @@ export interface ThreadManagementServiceShape {
readonly waitForThread: (
input: ThreadManagementWaitInput,
) => Effect.Effect<ThreadManagementWaitResult, ThreadManagementError>;
/**
* 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<void, ThreadManagementFailure>;
/**
* 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<ThreadManagementInterruptResult, ThreadManagementFailure>;
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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"]);
Expand Down Expand Up @@ -829,6 +889,8 @@ const make = Effect.gen(function* () {
listProjectThreads,
sendToThread,
waitForThread,
settleAfterRun,
settleThread,
interruptThread,
stopDelegatedTasks,
getThreadEventSequence: orchestrator.getThreadEventSequence,
Expand Down
68 changes: 68 additions & 0 deletions apps/server/src/orchestration-v2/runtimeLayer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/relay/AgentAwarenessRelay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading