Skip to content
Closed
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
77 changes: 77 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ import {
type OrchestratorMcpTaskCancelResult,
type OrchestratorMcpUpdateScheduledTaskInput,
type OrchestratorMcpThreadDetail,
type OrchestratorMcpThreadDeferOrganizationInput,
type OrchestratorMcpThreadDeferOrganizationResult,
type OrchestratorMcpThreadInterruptInput,
type OrchestratorMcpThreadInterruptResult,
type OrchestratorMcpThreadDeleteInput,
Expand Down Expand Up @@ -158,6 +160,10 @@ export interface OrchestratorMcpServiceShape {
scope: McpInvocationScope,
input: OrchestratorMcpThreadOrganizeInput,
) => Effect.Effect<OrchestratorMcpThreadOrganizeResult, OrchestratorMcpFailure>;
readonly deferThreadOrganization: (
scope: McpInvocationScope,
input: OrchestratorMcpThreadDeferOrganizationInput,
) => Effect.Effect<OrchestratorMcpThreadDeferOrganizationResult, OrchestratorMcpFailure>;
readonly deleteThread: (
scope: McpInvocationScope,
input: OrchestratorMcpThreadDeleteInput,
Expand Down Expand Up @@ -680,6 +686,23 @@ function organizationState(
};
}

function deferredOrganizationResult(
projection: OrchestrationV2ThreadProjection,
): OrchestratorMcpThreadDeferOrganizationResult {
const intent = projection.thread.deferredOrganization;
return {
threadId: projection.thread.id,
intent:
intent == null
? null
: {
runId: intent.runId,
action: intent.action,
requestedAt: DateTime.formatIso(intent.requestedAt),
},
};
}

function organizationCommand(input: {
readonly action: OrchestratorMcpThreadOrganizeAction;
readonly commandId: CommandId;
Expand Down Expand Up @@ -1833,6 +1856,60 @@ const make = Effect.gen(function* () {
);
return { outcomes } satisfies OrchestratorMcpThreadOrganizeResult;
}),
deferThreadOrganization: (scope, input) =>
Effect.gen(function* () {
yield* requireCapability(scope);
const parent = yield* loadProjection(scope.threadId);
if (input.operation === "read") return deferredOrganizationResult(parent);

const key = yield* requestKey(input.clientRequestId);
const command: OrchestrationV2Command =
input.operation === "cancel"
? {
type: "thread.organization.defer.cancel",
commandId: stableCommandId({
scope,
requestKey: key,
operation: "thread-defer-organization-cancel",
}),
threadId: scope.threadId,
}
: yield* Effect.gen(function* () {
const parentRun = latestActiveRun(parent);
if (
parentRun === undefined ||
parentRun.rootNodeId === null ||
parentRun.providerInstanceId !== scope.providerInstanceId
) {
return yield* failure(
"parent_not_active",
"Deferred organization requires an active run owned by this MCP provider session.",
);
}
return {
type: "thread.organization.defer",
commandId: stableCommandId({
scope,
requestKey: key,
operation: `thread-defer-organization-${input.action}`,
}),
threadId: scope.threadId,
runId: parentRun.id,
action: input.action!,
} satisfies OrchestrationV2Command;
});
yield* threadManagement
.dispatch(command)
.pipe(
Effect.mapError((error) =>
failure(
"orchestration_error",
`Unable to ${input.operation} deferred organization for thread ${scope.threadId}: ${errorMessage(error)}`,
),
),
);
return deferredOrganizationResult(yield* loadProjection(scope.threadId));
}),
deleteThread: (scope, input) =>
Effect.gen(function* () {
const threadId = input.threadId ?? scope.threadId;
Expand Down
32 changes: 32 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import {
OrchestratorMcpTaskCancelResult,
OrchestratorMcpThreadInterruptResult,
OrchestratorMcpThreadDeleteResult,
OrchestratorMcpThreadDeferOrganizationResult,
OrchestratorMcpThreadListResult,
OrchestratorMcpThreadReadResult,
OrchestratorMcpThreadOrganizeResult,
Expand Down Expand Up @@ -100,6 +101,9 @@ const decodeThreadInterruptResult = Schema.decodeUnknownEffect(
OrchestratorMcpThreadInterruptResult,
);
const decodeThreadDeleteResult = Schema.decodeUnknownEffect(OrchestratorMcpThreadDeleteResult);
const decodeThreadDeferOrganizationResult = Schema.decodeUnknownEffect(
OrchestratorMcpThreadDeferOrganizationResult,
);
const decodeThreadListResult = Schema.decodeUnknownEffect(OrchestratorMcpThreadListResult);
const decodeThreadReadResult = Schema.decodeUnknownEffect(OrchestratorMcpThreadReadResult);
const decodeThreadOrganizeResult = Schema.decodeUnknownEffect(OrchestratorMcpThreadOrganizeResult);
Expand Down Expand Up @@ -1230,6 +1234,10 @@ describe("orchestrator MCP toolkit", () => {
({ tool }) => tool.name === "t3_thread_organize",
);
expect(threadOrganizeTool?.tool.annotations?.destructiveHint).toBe(true);
const threadDeferOrganizationTool = server.tools.find(
({ tool }) => tool.name === "t3_thread_defer_organization",
);
expect(threadDeferOrganizationTool?.tool.annotations?.destructiveHint).toBe(true);
const threadDeleteTool = server.tools.find(
({ tool }) => tool.name === "t3_thread_delete",
);
Expand Down Expand Up @@ -1279,6 +1287,30 @@ describe("orchestrator MCP toolkit", () => {
]),
});

const scheduledOrganization = yield* decodeThreadDeferOrganizationResult(
(yield* invoke("t3_thread_defer_organization", {
operation: "schedule",
action: "settle",
clientRequestId: "defer-parent-settlement",
})).structuredContent,
).pipe(Effect.orDie);
expect(scheduledOrganization).toMatchObject({
threadId: parentThreadId,
intent: { runId: parentRun.id, action: "settle", requestedAt: expect.any(String) },
});
const readOrganization = yield* decodeThreadDeferOrganizationResult(
(yield* invoke("t3_thread_defer_organization", { operation: "read" }))
.structuredContent,
).pipe(Effect.orDie);
expect(readOrganization).toEqual(scheduledOrganization);
const cancelledOrganization = yield* decodeThreadDeferOrganizationResult(
(yield* invoke("t3_thread_defer_organization", {
operation: "cancel",
clientRequestId: "cancel-parent-settlement",
})).structuredContent,
).pipe(Effect.orDie);
expect(cancelledOrganization).toEqual({ threadId: parentThreadId, intent: null });

const scheduleTool = server.tools.find(({ tool }) => tool.name === "schedule_task");
expect(scheduleTool?.tool.annotations?.destructiveHint).toBe(true);
const scheduleCall = yield* invoke("schedule_task", {
Expand Down
6 changes: 6 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,12 @@ const handlers = {
const service = yield* OrchestratorMcpService;
return yield* service.organizeThreads(scope, input);
}),
t3_thread_defer_organization: (input) =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext;
const service = yield* OrchestratorMcpService;
return yield* service.deferThreadOrganization(scope, input);
}),
t3_thread_delete: (input) =>
Effect.gen(function* () {
const scope = yield* McpInvocationContext;
Expand Down
7 changes: 7 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/tools.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import {
CreateThreadsTool,
DelegateTaskTool,
ScheduleTaskTool,
ThreadDeferOrganizationTool,
ThreadUpdateTool,
} from "./tools.ts";

Expand Down Expand Up @@ -61,6 +62,12 @@ describe("orchestrator MCP tool guidance", () => {
assert.include(ScheduleTaskTool.description ?? "", "nextRunAt");
});

it("describes safe, calling-run-bound deferred organization", () => {
assert.include(ThreadDeferOrganizationTool.description ?? "", "THIS calling thread");
assert.include(ThreadDeferOrganizationTool.description ?? "", "current run");
assert.include(ThreadDeferOrganizationTool.description ?? "", "approval-blocked");
});

it("publishes thread metadata actions from an object-root schema", () => {
const schema = Tool.getJsonSchema(ThreadUpdateTool) as {
readonly type?: unknown;
Expand Down
15 changes: 15 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ import {
OrchestratorMcpThreadInterruptResult,
OrchestratorMcpThreadDeleteInput,
OrchestratorMcpThreadDeleteResult,
OrchestratorMcpThreadDeferOrganizationInput,
OrchestratorMcpThreadDeferOrganizationResult,
OrchestratorMcpThreadListInput,
OrchestratorMcpThreadListResult,
OrchestratorMcpThreadReadInput,
Expand Down Expand Up @@ -226,6 +228,18 @@ export const ThreadOrganizeTool = Tool.make("t3_thread_organize", {
.annotate(Tool.Title, "Organize T3 threads")
.annotate(Tool.Destructive, true);

export const ThreadDeferOrganizationTool = Tool.make("t3_thread_defer_organization", {
description:
"Schedule settlement or archival of THIS calling thread after the current run completes safely. The intent is durable and applies only when that run completes successfully with no newer, queued, active, approval-blocked, or title-regeneration work. Otherwise it is discarded. Use operation='read' to inspect the current intent or operation='cancel' to remove it. clientRequestId makes schedule and cancel retries idempotent.",
parameters: OrchestratorMcpThreadDeferOrganizationInput,
success: OrchestratorMcpThreadDeferOrganizationResult,
failure: OrchestratorMcpFailure,
failureMode: "return",
dependencies,
})
.annotate(Tool.Title, "Defer thread organization")
.annotate(Tool.Destructive, true);

export const ThreadDeleteTool = Tool.make("t3_thread_delete", {
description:
"Permanently delete one T3 thread in the calling project, defaulting to this thread. This cancels its active work and removes it from thread listings. clientRequestId makes retries idempotent.",
Expand Down Expand Up @@ -292,6 +306,7 @@ export const OrchestratorToolkit = Toolkit.make(
ThreadReadTool,
ThreadUpdateTool,
ThreadOrganizeTool,
ThreadDeferOrganizationTool,
ThreadDeleteTool,
ThreadSendTool,
ThreadWaitTool,
Expand Down
11 changes: 10 additions & 1 deletion apps/server/src/mcp/toolkits/worktree/registration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,10 @@ const ToolsListPayload = Schema.fromJsonString(
tools: Schema.Array(
Schema.Struct({
name: Schema.String,
inputSchema: Schema.Struct({ type: Schema.optional(Schema.String) }),
inputSchema: Schema.Struct({
type: Schema.optional(Schema.String),
properties: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)),
}),
annotations: Schema.optional(
Schema.Struct({
readOnlyHint: Schema.optional(Schema.Boolean),
Expand Down Expand Up @@ -114,6 +117,7 @@ it.effect("production mcp layer lists worktree tools over http", () =>
// than replacing them.
expect(toolNames).toContain("preview_status");
expect(toolNames).toContain("delegate_task");
expect(toolNames).toContain("t3_thread_defer_organization");

// The handoff tool mutates thread state, reaches the network (origin
// fetch), and runs project setup scripts, so its MCP hints must not
Expand All @@ -132,6 +136,11 @@ it.effect("production mcp layer lists worktree tools over http", () =>
for (const tool of tools) {
expect(tool.inputSchema.type, `inputSchema.type of ${tool.name}`).toBe("object");
}
const deferredOrganization = tools.find(
(tool) => tool.name === "t3_thread_defer_organization",
);
expect(deferredOrganization?.inputSchema.properties).toHaveProperty("operation");
expect(deferredOrganization?.inputSchema.properties).toHaveProperty("action");
}),
).pipe(Effect.provide(Layer.mergeAll(NodeHttpServer.layerTest, NodeServices.layer))),
);
Loading
Loading