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
11 changes: 10 additions & 1 deletion apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ import {
PreviewSnapshotToolkit,
PreviewStandardToolkit,
} from "./toolkits/preview/tools.ts";
import { ThreadToolkitHandlersLayerLive } from "./toolkits/threads/handlers.ts";
import { ThreadToolkit } from "./toolkits/threads/tools.ts";

const unauthorized = HttpServerResponse.jsonUnsafe(
{
Expand Down Expand Up @@ -208,10 +210,17 @@ export const PreviewToolkitRegistrationLive = Layer.mergeAll(
PreviewSnapshotRegistrationLive,
);

export const ThreadToolkitRegistrationLive = McpServer.toolkit(ThreadToolkit).pipe(
Layer.provide(ThreadToolkitHandlersLayerLive),
);

const McpTransportLive = McpServer.layerHttp({
name: "T3 Code",
version: packageJson.version,
path: "/mcp",
}).pipe(Layer.provide(McpAuthMiddlewareLive));

export const layer = PreviewToolkitRegistrationLive.pipe(Layer.provideMerge(McpTransportLive));
export const layer = Layer.mergeAll(
PreviewToolkitRegistrationLive,
ThreadToolkitRegistrationLive,
).pipe(Layer.provideMerge(McpTransportLive));
23 changes: 23 additions & 0 deletions apps/server/src/mcp/McpInvocationContext.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,3 +37,26 @@ it.effect("reports the scoped credential context when preview capability is unav
expect(error.message).toBe("MCP credential does not grant the preview capability.");
});
});

it.effect("reports a missing thread coordination capability", () => {
const invocation: McpInvocationContext.McpInvocationScope = {
environmentId: EnvironmentId.make("environment-1"),
threadId: ThreadId.make("thread-1"),
providerSessionId: "provider-session-1",
providerInstanceId: ProviderInstanceId.make("codex"),
capabilities: new Set(["preview"]),
issuedAt: 1,
expiresAt: 2,
};

return Effect.gen(function* () {
const error = yield* McpInvocationContext.requireMcpCapability("threads").pipe(
Effect.provideService(McpInvocationContext.McpInvocationContext, invocation),
Effect.flip,
);

expect(error).toBeInstanceOf(PreviewAutomationUnavailableError);
expect(error.capability).toBe("threads");
expect(error.message).toBe("MCP credential does not grant the threads capability.");
});
});
2 changes: 1 addition & 1 deletion apps/server/src/mcp/McpInvocationContext.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import {
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";

export type McpCapability = "preview";
export type McpCapability = "preview" | "threads";

export interface McpInvocationScope {
readonly environmentId: EnvironmentId;
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/mcp/McpSessionRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,7 @@ const makeWithOptions = Effect.fn("McpSessionRegistry.make")(function* (
threadId: ThreadId.make(request.threadId),
providerSessionId,
providerInstanceId: ProviderInstanceId.make(request.providerInstanceId),
capabilities: new Set(["preview"]),
capabilities: new Set(["preview", "threads"]),
issuedAt,
expiresAt,
};
Expand Down
249 changes: 249 additions & 0 deletions apps/server/src/mcp/toolkits/threads/handlers.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,249 @@
import {
CommandId,
MessageId,
ThreadId,
type OrchestrationThread,
type OrchestrationThreadShell,
} from "@t3tools/contracts";
import * as Cause from "effect/Cause";
import * as Crypto from "effect/Crypto";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";

import * as McpInvocationContext from "../../McpInvocationContext.ts";
import { OrchestrationEngineService } from "../../../orchestration/Services/OrchestrationEngine.ts";
import { ProjectionSnapshotQuery } from "../../../orchestration/Services/ProjectionSnapshotQuery.ts";
import { ThreadToolkit, ThreadToolError } from "./tools.ts";

const asToolFailure = (operation: string) => (cause: Cause.Cause<unknown>) =>
Effect.fail(
new ThreadToolError({
operation,
detail: Cause.pretty(cause),
}),
);

const runThreadOperation = <A, E, R>(
operation: string,
effect: Effect.Effect<A, E, R>,
): Effect.Effect<A, ThreadToolError, R> => effect.pipe(Effect.catchCause(asToolFailure(operation)));

const makeHandlers = Effect.gen(function* () {
const invocation = yield* McpInvocationContext.requireMcpCapability("threads");
const engine = yield* OrchestrationEngineService;
const query = yield* ProjectionSnapshotQuery;
const crypto = yield* Crypto.Crypto;

const newId = Effect.fn("ThreadToolkit.newId")(function* () {
return yield* crypto.randomUUIDv4.pipe(Effect.orDie);
});

const currentThread = Effect.fn("ThreadToolkit.currentThread")(function* () {
const thread = yield* query.getThreadDetailById(invocation.threadId);
if (Option.isNone(thread)) {
return yield* new ThreadToolError({
operation: "thread_current",
detail: `Current thread '${invocation.threadId}' was not found.`,
});
}
return thread.value;
});

const projectThread = Effect.fn("ThreadToolkit.projectThread")(function* (threadId: ThreadId) {
const [current, target] = yield* Effect.all([
currentThread(),
query.getThreadDetailById(threadId),
]);
if (Option.isNone(target) || target.value.projectId !== current.projectId) {
return yield* new ThreadToolError({
operation: "thread_target",
detail: `Thread '${threadId}' is not an active thread in the current project.`,
});
}
return target.value;
});

const startTurn = Effect.fn("ThreadToolkit.startTurn")(function* (
thread: OrchestrationThread,
prompt: string,
) {
const [commandUuid, messageUuid, createdAt] = yield* Effect.all([
newId(),
newId(),
DateTime.now.pipe(Effect.map(DateTime.formatIso)),
]);
yield* engine.dispatch({
type: "thread.turn.start",
commandId: CommandId.make(commandUuid),
threadId: thread.id,
message: {
messageId: MessageId.make(messageUuid),
role: "user",
text: prompt,
attachments: [],
},
modelSelection: thread.modelSelection,
runtimeMode: thread.runtimeMode,
interactionMode: thread.interactionMode,
createdAt,
});
});

const resultFor = (thread: Pick<OrchestrationThread, "id" | "title">, turnStarted: boolean) => ({
threadId: thread.id,
title: thread.title,
turnStarted,
});

return {
thread_list: (input: { readonly includeArchived?: boolean | undefined }) =>
Effect.gen(function* () {
const current = yield* currentThread();
const snapshot = yield* query.getShellSnapshot();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Medium threads/handlers.ts:103

thread_list ignores includeArchived: true — archived threads are always excluded from the result, even when the caller explicitly requests them. The handler only calls query.getShellSnapshot(), which does not return archived threads, so the subsequent input.includeArchived === true || thread.archivedAt === null filter can never surface archived entries because they were never in the snapshot to begin with. To honor includeArchived, the handler must also query (and merge) getArchivedShellSnapshot() when archived entries are requested.

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/mcp/toolkits/threads/handlers.ts around line 103:

`thread_list` ignores `includeArchived: true` — archived threads are always excluded from the result, even when the caller explicitly requests them. The handler only calls `query.getShellSnapshot()`, which does not return archived threads, so the subsequent `input.includeArchived === true || thread.archivedAt === null` filter can never surface archived entries because they were never in the snapshot to begin with. To honor `includeArchived`, the handler must also query (and merge) `getArchivedShellSnapshot()` when archived entries are requested.

return snapshot.threads
.filter(
(thread) =>
thread.projectId === current.projectId &&
(input.includeArchived === true || thread.archivedAt === null),
)
.map((thread: OrchestrationThreadShell) => ({
threadId: thread.id,
title: thread.title,
status: thread.session?.status ?? "idle",
createdAt: thread.createdAt,
...(thread.forkedFrom?.threadId
? { forkedFromThreadId: thread.forkedFrom.threadId }
: {}),
...(thread.forkedFrom?.turnId ? { forkedFromTurnId: thread.forkedFrom.turnId } : {}),
}));
}),
thread_create: (input: {
readonly title?: string | undefined;
readonly prompt?: string | undefined;
}) =>
Effect.gen(function* () {
const current = yield* currentThread();
const [threadUuid, commandUuid, createdAt] = yield* Effect.all([
newId(),
newId(),
DateTime.now.pipe(Effect.map(DateTime.formatIso)),
]);
const threadId = ThreadId.make(threadUuid);
const title = input.title ?? "New thread";
yield* engine.dispatch({
type: "thread.create",
commandId: CommandId.make(commandUuid),
threadId,
projectId: current.projectId,
title,
modelSelection: current.modelSelection,
runtimeMode: current.runtimeMode,
interactionMode: current.interactionMode,
branch: current.branch,
worktreePath: current.worktreePath,
createdAt,
});
const created = { ...current, id: threadId, title };
if (input.prompt !== undefined) yield* startTurn(created, input.prompt);
return resultFor(created, input.prompt !== undefined);
}),
thread_fork: (input: {
readonly sourceThreadId?: ThreadId | undefined;
readonly sourceTurnId?: import("@t3tools/contracts").TurnId | undefined;
readonly title?: string | undefined;
readonly prompt?: string | undefined;
}) =>
Effect.gen(function* () {
const source = yield* projectThread(input.sourceThreadId ?? invocation.threadId);
const [threadUuid, commandUuid, createdAt] = yield* Effect.all([
newId(),
newId(),
DateTime.now.pipe(Effect.map(DateTime.formatIso)),
]);
const threadId = ThreadId.make(threadUuid);
const title = input.title ?? `${source.title} (fork)`;
yield* engine.dispatch({
type: "thread.fork",
commandId: CommandId.make(commandUuid),
threadId,
sourceThreadId: source.id,
...(input.sourceTurnId !== undefined ? { sourceTurnId: input.sourceTurnId } : {}),
title,
createdAt,
});
const forked = {
...source,
id: threadId,
title,
forkedFrom: {
threadId: source.id,
turnId: input.sourceTurnId ?? source.latestTurn?.turnId ?? null,
},
session: null,
};
if (input.prompt !== undefined) yield* startTurn(forked, input.prompt);
return resultFor(forked, input.prompt !== undefined);
}),
thread_send: (input: { readonly threadId: ThreadId; readonly prompt: string }) =>
Effect.gen(function* () {
const target = yield* projectThread(input.threadId);
yield* startTurn(target, input.prompt);
return resultFor(target, true);
}),
thread_archive: (input: { readonly threadId: ThreadId }) =>
Effect.gen(function* () {
if (input.threadId === invocation.threadId) {
return yield* new ThreadToolError({
operation: "thread_archive",
detail: "The current agent thread cannot archive itself.",
});
}
yield* projectThread(input.threadId);
yield* engine.dispatch({
type: "thread.archive",
commandId: CommandId.make(yield* newId()),
threadId: input.threadId,
});
return null;
}),
};
});

export const ThreadToolkitHandlers = {
thread_list: (input: { readonly includeArchived?: boolean | undefined }) =>
runThreadOperation(
"thread_list",
makeHandlers.pipe(Effect.flatMap((handlers) => handlers.thread_list(input))),
),
thread_create: (input: {
readonly title?: string | undefined;
readonly prompt?: string | undefined;
}) =>
runThreadOperation(
"thread_create",
makeHandlers.pipe(Effect.flatMap((handlers) => handlers.thread_create(input))),
),
thread_fork: (input: {
readonly sourceThreadId?: ThreadId | undefined;
readonly sourceTurnId?: import("@t3tools/contracts").TurnId | undefined;
readonly title?: string | undefined;
readonly prompt?: string | undefined;
}) =>
runThreadOperation(
"thread_fork",
makeHandlers.pipe(Effect.flatMap((handlers) => handlers.thread_fork(input))),
),
thread_send: (input: { readonly threadId: ThreadId; readonly prompt: string }) =>
runThreadOperation(
"thread_send",
makeHandlers.pipe(Effect.flatMap((handlers) => handlers.thread_send(input))),
),
thread_archive: (input: { readonly threadId: ThreadId }) =>
runThreadOperation(
"thread_archive",
makeHandlers.pipe(Effect.flatMap((handlers) => handlers.thread_archive(input))),
),
} satisfies Parameters<typeof ThreadToolkit.toLayer>[0];

export const ThreadToolkitHandlersLayerLive = ThreadToolkit.toLayer(ThreadToolkitHandlers);
26 changes: 26 additions & 0 deletions apps/server/src/mcp/toolkits/threads/tools.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
import { expect, it } from "@effect/vitest";
import { Tool } from "effect/unstable/ai";

import { ThreadToolkit } from "./tools.ts";

it("exports provider-compatible thread management schemas", () => {
expect(Object.keys(ThreadToolkit.tools).sort()).toEqual([
"thread_archive",
"thread_create",
"thread_fork",
"thread_list",
"thread_send",
]);

for (const tool of Object.values(ThreadToolkit.tools)) {
const schema = Tool.getJsonSchema(tool) as {
readonly type?: unknown;
readonly anyOf?: unknown;
readonly oneOf?: unknown;
};
expect(tool.description?.length ?? 0).toBeGreaterThan(40);
expect(schema.type, `${tool.name} must export a top-level object schema`).toBe("object");
expect(schema.anyOf, `${tool.name} must not export a root anyOf`).toBeUndefined();
expect(schema.oneOf, `${tool.name} must not export a root oneOf`).toBeUndefined();
}
});
Loading
Loading