diff --git a/apps/server/src/mcp/McpHttpServer.test.ts b/apps/server/src/mcp/McpHttpServer.test.ts index ca3341be7f3..3dd85109169 100644 --- a/apps/server/src/mcp/McpHttpServer.test.ts +++ b/apps/server/src/mcp/McpHttpServer.test.ts @@ -11,6 +11,14 @@ import { HttpBody, HttpClient, HttpRouter, HttpServerResponse } from "effect/uns import * as McpHttpServer from "./McpHttpServer.ts"; import * as McpInvocationContext from "./McpInvocationContext.ts"; import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts"; +import { + ProjectionSnapshotQuery, + type ProjectionSnapshotQueryShape, +} from "../orchestration/Services/ProjectionSnapshotQuery.ts"; +import { + ProviderSessionDirectory, + type ProviderSessionDirectoryShape, +} from "../provider/Services/ProviderSessionDirectory.ts"; const environmentId = EnvironmentId.make("environment-mcp-test"); const threadId = ThreadId.make("thread-mcp-test"); @@ -38,6 +46,36 @@ const TestLayer = McpHttpServer.PreviewToolkitRegistrationLive.pipe( Layer.provideMerge(McpServer.McpServer.layer), Layer.provideMerge(PreviewAutomationBroker.layer.pipe(Layer.provide(NodeServices.layer))), ); +const unused = () => Effect.die("unused"); +const projectionSnapshotQueryStub: ProjectionSnapshotQueryShape = { + getCommandReadModel: unused, + getSnapshot: unused, + getShellSnapshot: unused, + getArchivedShellSnapshot: unused, + getSnapshotSequence: unused, + getCounts: unused, + getActiveProjectByWorkspaceRoot: unused, + getProjectShellById: unused, + getFirstActiveThreadIdByProjectId: unused, + getThreadCheckpointContext: unused, + getFullThreadDiffContext: unused, + getThreadShellById: unused, + getThreadDetailById: unused, + getThreadDetailSnapshot: unused, +}; +const providerSessionDirectoryStub: ProviderSessionDirectoryShape = { + upsert: unused, + getProvider: unused, + getBinding: unused, + listThreadIds: unused, + listBindings: unused, +}; +const CombinedToolkitTestLayer = McpHttpServer.ToolkitRegistrationLive.pipe( + Layer.provideMerge(McpServer.McpServer.layer), + Layer.provideMerge(PreviewAutomationBroker.layer.pipe(Layer.provide(NodeServices.layer))), + Layer.provideMerge(Layer.succeed(ProjectionSnapshotQuery, projectionSnapshotQueryStub)), + Layer.provideMerge(Layer.succeed(ProviderSessionDirectory, providerSessionDirectoryStub)), +); it("normalizes empty successful notification responses to accepted", () => { const notificationResponse = McpHttpServer.normalizeMcpHttpResponse( @@ -51,6 +89,21 @@ it("normalizes empty successful notification responses to accepted", () => { expect(resultResponse.status).toBe(200); }); +it.effect("registers cross-thread reading without replacing preview tools", () => + Effect.gen(function* () { + const server = yield* McpServer.McpServer; + const toolNames = new Set(server.tools.map(({ tool }) => tool.name)); + expect(toolNames.has("preview_status")).toBe(true); + expect(toolNames.has("preview_snapshot")).toBe(true); + expect(toolNames.has("read_thread")).toBe(true); + + const readThreadTool = server.tools.find(({ tool }) => tool.name === "read_thread"); + expect(readThreadTool?.tool.annotations?.readOnlyHint).toBe(true); + expect(readThreadTool?.tool.annotations?.idempotentHint).toBe(true); + expect(readThreadTool?.tool.annotations?.destructiveHint).toBe(false); + }).pipe(Effect.provide(CombinedToolkitTestLayer)), +); + it.effect("returns bounded structural preview snapshot failures", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/mcp/McpHttpServer.ts b/apps/server/src/mcp/McpHttpServer.ts index 6774731a73e..e18f0e5f78c 100644 --- a/apps/server/src/mcp/McpHttpServer.ts +++ b/apps/server/src/mcp/McpHttpServer.ts @@ -22,6 +22,8 @@ import { PreviewSnapshotToolkit, PreviewStandardToolkit, } from "./toolkits/preview/tools.ts"; +import { ThreadControlToolkitHandlersLive } from "./toolkits/thread-control/handlers.ts"; +import { ThreadControlToolkit } from "./toolkits/thread-control/tools.ts"; const unauthorized = HttpServerResponse.jsonUnsafe( { @@ -208,10 +210,19 @@ export const PreviewToolkitRegistrationLive = Layer.mergeAll( PreviewSnapshotRegistrationLive, ); +export const ThreadControlToolkitRegistrationLive = McpServer.toolkit(ThreadControlToolkit).pipe( + Layer.provide(ThreadControlToolkitHandlersLive), +); + +export const ToolkitRegistrationLive = Layer.mergeAll( + PreviewToolkitRegistrationLive, + ThreadControlToolkitRegistrationLive, +); + 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 = ToolkitRegistrationLive.pipe(Layer.provideMerge(McpTransportLive)); diff --git a/apps/server/src/mcp/McpInvocationContext.ts b/apps/server/src/mcp/McpInvocationContext.ts index b13bf2d312e..9944ce52ee2 100644 --- a/apps/server/src/mcp/McpInvocationContext.ts +++ b/apps/server/src/mcp/McpInvocationContext.ts @@ -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" | "thread-read"; export interface McpInvocationScope { readonly environmentId: EnvironmentId; @@ -25,7 +25,7 @@ export class McpInvocationContext extends Context.Service< >()("t3/mcp/McpInvocationContext") {} export const requireMcpCapability = Effect.fn("mcp.requireCapability")(function* ( - capability: McpCapability, + capability: "preview", ) { const invocation = yield* McpInvocationContext; if (!invocation.capabilities.has(capability)) { diff --git a/apps/server/src/mcp/McpSessionRegistry.test.ts b/apps/server/src/mcp/McpSessionRegistry.test.ts index a91d98febd8..813ef7651c8 100644 --- a/apps/server/src/mcp/McpSessionRegistry.test.ts +++ b/apps/server/src/mcp/McpSessionRegistry.test.ts @@ -47,6 +47,7 @@ it.effect("stores only a token hash, resolves the bearer token, and revokes by t const resolved = yield* registry.resolve(token); expect(resolved?.threadId).toBe(threadId); + expect(resolved?.capabilities).toEqual(new Set(["preview", "thread-read"])); yield* registry.revokeThread(threadId); expect(yield* registry.resolve(token)).toBeUndefined(); diff --git a/apps/server/src/mcp/McpSessionRegistry.ts b/apps/server/src/mcp/McpSessionRegistry.ts index 67c4f2f0ff0..9575d0355be 100644 --- a/apps/server/src/mcp/McpSessionRegistry.ts +++ b/apps/server/src/mcp/McpSessionRegistry.ts @@ -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", "thread-read"]), issuedAt, expiresAt, }; diff --git a/apps/server/src/mcp/toolkits/thread-control/handlers.test.ts b/apps/server/src/mcp/toolkits/thread-control/handlers.test.ts new file mode 100644 index 00000000000..ac45f314e31 --- /dev/null +++ b/apps/server/src/mcp/toolkits/thread-control/handlers.test.ts @@ -0,0 +1,209 @@ +import { expect, it } from "@effect/vitest"; +import { + EventId, + EnvironmentId, + MessageId, + ProjectId, + ProviderDriverKind, + ProviderInstanceId, + ThreadId, + type OrchestrationProjectShell, + type OrchestrationThread, +} from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; + +import * as McpInvocationContext from "../../McpInvocationContext.ts"; +import { + ProjectionSnapshotQuery, + type ProjectionSnapshotQueryShape, +} from "../../../orchestration/Services/ProjectionSnapshotQuery.ts"; +import { + ProviderSessionDirectory, + type ProviderRuntimeBindingWithMetadata, + type ProviderSessionDirectoryShape, +} from "../../../provider/Services/ProviderSessionDirectory.ts"; +import { readThread } from "./handlers.ts"; + +const NOW = "2026-01-01T00:00:00.000Z"; +const t3ThreadId = ThreadId.make("t3-thread-1"); +const providerThreadId = "019c-native-codex-thread"; +const projectId = ProjectId.make("project-1"); + +const project: OrchestrationProjectShell = { + id: projectId, + title: "T3 Code", + workspaceRoot: "/workspace/t3code", + defaultModelSelection: null, + scripts: [], + createdAt: NOW, + updatedAt: NOW, +}; + +const thread: OrchestrationThread = { + id: t3ThreadId, + projectId, + title: "Cross-thread implementation", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5.4", + }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + latestTurn: null, + createdAt: NOW, + updatedAt: NOW, + archivedAt: null, + settledOverride: null, + settledAt: null, + snoozedUntil: null, + snoozedAt: null, + deletedAt: null, + messages: [ + { + id: MessageId.make("message-1"), + role: "user", + text: "first", + turnId: null, + streaming: false, + createdAt: NOW, + updatedAt: NOW, + }, + { + id: MessageId.make("message-2"), + role: "assistant", + text: "second", + turnId: null, + streaming: false, + createdAt: NOW, + updatedAt: NOW, + }, + ], + proposedPlans: [], + activities: [ + { + id: EventId.make("activity-1"), + tone: "info", + kind: "status", + summary: "Working", + payload: {}, + turnId: null, + createdAt: NOW, + }, + ], + checkpoints: [], + session: { + threadId: t3ThreadId, + status: "ready", + providerName: "codex", + providerInstanceId: ProviderInstanceId.make("codex"), + providerThreadId, + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: NOW, + }, +}; + +const binding: ProviderRuntimeBindingWithMetadata = { + threadId: t3ThreadId, + provider: ProviderDriverKind.make("codex"), + providerInstanceId: ProviderInstanceId.make("codex"), + resumeCursor: { threadId: providerThreadId }, + lastSeenAt: NOW, +}; + +const invocation = McpInvocationContext.McpInvocationContext.of({ + environmentId: EnvironmentId.make("environment-1"), + threadId: ThreadId.make("requesting-thread"), + providerSessionId: "provider-session-1", + providerInstanceId: ProviderInstanceId.make("codex"), + capabilities: new Set(["preview", "thread-read"] as const), + issuedAt: 1, + expiresAt: Number.MAX_SAFE_INTEGER, +}); + +const unused = () => Effect.die("unused"); + +function makeSnapshotQuery(): ProjectionSnapshotQueryShape { + return { + getCommandReadModel: unused, + getSnapshot: unused, + getShellSnapshot: unused, + getArchivedShellSnapshot: unused, + getSnapshotSequence: unused, + getCounts: unused, + getActiveProjectByWorkspaceRoot: unused, + getFirstActiveThreadIdByProjectId: unused, + getThreadCheckpointContext: unused, + getFullThreadDiffContext: unused, + getThreadShellById: unused, + getThreadDetailSnapshot: unused, + getProjectShellById: (id) => + Effect.succeed(id === projectId ? Option.some(project) : Option.none()), + getThreadDetailById: (id) => + Effect.succeed(id === t3ThreadId ? Option.some(thread) : Option.none()), + }; +} + +const sessionDirectory: ProviderSessionDirectoryShape = { + upsert: unused, + getProvider: unused, + getBinding: unused, + listThreadIds: unused, + listBindings: () => Effect.succeed([binding]), +}; + +function runRead(input: Parameters[0]) { + return readThread(input).pipe( + Effect.provideService(McpInvocationContext.McpInvocationContext, invocation), + Effect.provideService(ProjectionSnapshotQuery, makeSnapshotQuery()), + Effect.provideService(ProviderSessionDirectory, sessionDirectory), + ); +} + +it.effect("reads a thread by its T3 ID without modifying it", () => + Effect.gen(function* () { + const before = structuredClone(thread); + const result = yield* runRead({ threadId: t3ThreadId }); + + expect(result.threadId).toBe(t3ThreadId); + expect(result.providerThreadId).toBe(providerThreadId); + expect(result.messages.map(({ text }) => text)).toEqual(["first", "second"]); + expect(thread).toEqual(before); + }), +); + +it.effect("resolves a copied Codex thread ID and bounds the returned history", () => + Effect.gen(function* () { + const result = yield* runRead({ + threadId: providerThreadId, + messageLimit: 1, + activityLimit: 1, + }); + + expect(result.threadId).toBe(t3ThreadId); + expect(result.messages.map(({ text }) => text)).toEqual(["second"]); + expect(result.totals.messages).toBe(2); + expect(result.truncated.messages).toBe(true); + }), +); + +it.effect("does not disclose a thread when cross-thread capability is absent", () => + Effect.gen(function* () { + const restrictedInvocation = { + ...invocation, + capabilities: new Set(["preview"] as const), + }; + const error = yield* readThread({ threadId: t3ThreadId }).pipe( + Effect.provideService(McpInvocationContext.McpInvocationContext, restrictedInvocation), + Effect.provideService(ProjectionSnapshotQuery, makeSnapshotQuery()), + Effect.provideService(ProviderSessionDirectory, sessionDirectory), + Effect.flip, + ); + + expect(error._tag).toBe("ThreadReadUnavailableError"); + }), +); diff --git a/apps/server/src/mcp/toolkits/thread-control/handlers.ts b/apps/server/src/mcp/toolkits/thread-control/handlers.ts new file mode 100644 index 00000000000..00c9b9c3e47 --- /dev/null +++ b/apps/server/src/mcp/toolkits/thread-control/handlers.ts @@ -0,0 +1,173 @@ +import { ThreadId } from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; + +import * as McpInvocationContext from "../../McpInvocationContext.ts"; +import { ProjectionSnapshotQuery } from "../../../orchestration/Services/ProjectionSnapshotQuery.ts"; +import type { ProviderRuntimeBindingWithMetadata } from "../../../provider/Services/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectory } from "../../../provider/Services/ProviderSessionDirectory.ts"; +import { + ThreadControlToolkit, + ThreadReadAmbiguousError, + ThreadReadFailedError, + ThreadReadNotFoundError, + ThreadReadUnavailableError, + type ReadThreadInput, + type ReadThreadResult, +} from "./tools.ts"; + +const DEFAULT_RESULT_LIMIT = 50; + +function isRecord(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} + +export function readProviderThreadId( + binding: Pick, +): string | undefined { + if (!isRecord(binding.resumeCursor)) { + return undefined; + } + if (binding.provider === "codex") { + return typeof binding.resumeCursor.threadId === "string" + ? binding.resumeCursor.threadId + : undefined; + } + if (binding.provider === "claude") { + if (typeof binding.resumeCursor.resume === "string") { + return binding.resumeCursor.resume; + } + return typeof binding.resumeCursor.threadId === "string" + ? binding.resumeCursor.threadId + : undefined; + } + return typeof binding.resumeCursor.sessionId === "string" + ? binding.resumeCursor.sessionId + : undefined; +} + +function mapReadFailure(operation: string) { + return (cause: unknown) => + Effect.logWarning("cross-thread read failed", { operation, cause }).pipe( + Effect.andThen(Effect.fail(new ThreadReadFailedError({ operation }))), + ); +} + +export const readThread = Effect.fn("ThreadControlToolkit.readThread")(function* ( + input: ReadThreadInput, +): Effect.fn.Return< + ReadThreadResult, + | ThreadReadUnavailableError + | ThreadReadNotFoundError + | ThreadReadAmbiguousError + | ThreadReadFailedError, + McpInvocationContext.McpInvocationContext | ProjectionSnapshotQuery | ProviderSessionDirectory +> { + const invocation = yield* McpInvocationContext.McpInvocationContext; + if (!invocation.capabilities.has("thread-read")) { + return yield* new ThreadReadUnavailableError({ threadId: input.threadId }); + } + + const snapshotQuery = yield* ProjectionSnapshotQuery; + const sessionDirectory = yield* ProviderSessionDirectory; + const directThreadId = ThreadId.make(input.threadId); + const directThread = yield* snapshotQuery + .getThreadDetailById(directThreadId) + .pipe(Effect.catch(mapReadFailure("resolve T3 thread ID"))); + + let targetThreadId = directThreadId; + let thread = Option.getOrUndefined(directThread); + let targetBinding: ProviderRuntimeBindingWithMetadata | undefined; + + const bindings = yield* sessionDirectory + .listBindings() + .pipe(Effect.catch(mapReadFailure("list provider bindings"))); + + if (!thread) { + const matches = bindings.filter((binding) => readProviderThreadId(binding) === input.threadId); + if (matches.length > 1) { + return yield* new ThreadReadAmbiguousError({ + threadId: input.threadId, + matchCount: matches.length, + }); + } + targetBinding = matches[0]; + if (!targetBinding) { + return yield* new ThreadReadNotFoundError({ threadId: input.threadId }); + } + targetThreadId = targetBinding.threadId; + const resolved = yield* snapshotQuery + .getThreadDetailById(targetThreadId) + .pipe(Effect.catch(mapReadFailure("read provider thread"))); + thread = Option.getOrUndefined(resolved); + } else { + targetBinding = bindings.find((binding) => binding.threadId === targetThreadId); + } + + if (!thread) { + return yield* new ThreadReadNotFoundError({ threadId: input.threadId }); + } + + const project = yield* snapshotQuery + .getProjectShellById(thread.projectId) + .pipe(Effect.catch(mapReadFailure("read thread project")), Effect.map(Option.getOrUndefined)); + if (!project) { + return yield* new ThreadReadNotFoundError({ threadId: input.threadId }); + } + + const messageLimit = input.messageLimit ?? DEFAULT_RESULT_LIMIT; + const activityLimit = input.activityLimit ?? DEFAULT_RESULT_LIMIT; + const messages = thread.messages.slice(-messageLimit).map((message) => ({ + id: message.id, + role: message.role, + text: message.text, + streaming: message.streaming, + createdAt: message.createdAt, + updatedAt: message.updatedAt, + })); + const activities = thread.activities.slice(-activityLimit).map((activity) => ({ + id: activity.id, + tone: activity.tone, + kind: activity.kind, + summary: activity.summary, + createdAt: activity.createdAt, + })); + const providerThreadId = + thread.session?.providerThreadId ?? + (targetBinding ? readProviderThreadId(targetBinding) : undefined); + + return { + threadId: thread.id, + ...(providerThreadId ? { providerThreadId } : {}), + title: thread.title, + project: { + id: project.id, + title: project.title, + workspaceRoot: project.workspaceRoot, + }, + provider: thread.session?.providerName ?? targetBinding?.provider ?? null, + status: thread.session?.status ?? "idle", + latestTurn: thread.latestTurn + ? { + id: thread.latestTurn.turnId, + state: thread.latestTurn.state, + startedAt: thread.latestTurn.startedAt, + completedAt: thread.latestTurn.completedAt, + } + : null, + messages, + activities, + totals: { + messages: thread.messages.length, + activities: thread.activities.length, + }, + truncated: { + messages: messages.length < thread.messages.length, + activities: activities.length < thread.activities.length, + }, + }; +}); + +export const ThreadControlToolkitHandlersLive = ThreadControlToolkit.toLayer({ + read_thread: readThread, +}); diff --git a/apps/server/src/mcp/toolkits/thread-control/tools.ts b/apps/server/src/mcp/toolkits/thread-control/tools.ts new file mode 100644 index 00000000000..d5e4cc72e81 --- /dev/null +++ b/apps/server/src/mcp/toolkits/thread-control/tools.ts @@ -0,0 +1,151 @@ +import * as Schema from "effect/Schema"; +import { Tool, Toolkit } from "effect/unstable/ai"; + +import * as McpInvocationContext from "../../McpInvocationContext.ts"; +import { ProjectionSnapshotQuery } from "../../../orchestration/Services/ProjectionSnapshotQuery.ts"; +import { ProviderSessionDirectory } from "../../../provider/Services/ProviderSessionDirectory.ts"; + +const BoundedResultLimit = Schema.Int.check(Schema.isBetween({ minimum: 1, maximum: 200 })); + +export const ReadThreadInput = Schema.Struct({ + threadId: Schema.String.annotate({ + description: + "A T3 Code thread ID or a provider-native thread ID copied from a T3 Code thread menu.", + }), + messageLimit: Schema.optionalKey( + BoundedResultLimit.annotate({ + description: + "Maximum number of newest conversation messages to return. Defaults to 50 and must be between 1 and 200.", + }), + ), + activityLimit: Schema.optionalKey( + BoundedResultLimit.annotate({ + description: + "Maximum number of newest activity summaries to return. Defaults to 50 and must be between 1 and 200.", + }), + ), +}); +export type ReadThreadInput = typeof ReadThreadInput.Type; + +const ThreadReadMessage = Schema.Struct({ + id: Schema.String, + role: Schema.String, + text: Schema.String, + streaming: Schema.Boolean, + createdAt: Schema.String, + updatedAt: Schema.String, +}); + +const ThreadReadActivity = Schema.Struct({ + id: Schema.String, + tone: Schema.String, + kind: Schema.String, + summary: Schema.String, + createdAt: Schema.String, +}); + +const ThreadReadLatestTurn = Schema.Struct({ + id: Schema.String, + state: Schema.String, + startedAt: Schema.NullOr(Schema.String), + completedAt: Schema.NullOr(Schema.String), +}); + +export const ReadThreadResult = Schema.Struct({ + threadId: Schema.String, + providerThreadId: Schema.optionalKey(Schema.String), + title: Schema.String, + project: Schema.Struct({ + id: Schema.String, + title: Schema.String, + workspaceRoot: Schema.String, + }), + provider: Schema.NullOr(Schema.String), + status: Schema.String, + latestTurn: Schema.NullOr(ThreadReadLatestTurn), + messages: Schema.Array(ThreadReadMessage), + activities: Schema.Array(ThreadReadActivity), + totals: Schema.Struct({ + messages: Schema.Int, + activities: Schema.Int, + }), + truncated: Schema.Struct({ + messages: Schema.Boolean, + activities: Schema.Boolean, + }), +}); +export type ReadThreadResult = typeof ReadThreadResult.Type; + +export class ThreadReadUnavailableError extends Schema.TaggedErrorClass()( + "ThreadReadUnavailableError", + { + threadId: Schema.String, + }, +) { + override get message(): string { + return "This MCP credential does not grant cross-thread read access."; + } +} + +export class ThreadReadNotFoundError extends Schema.TaggedErrorClass()( + "ThreadReadNotFoundError", + { + threadId: Schema.String, + }, +) { + override get message(): string { + return `No readable T3 Code thread was found for '${this.threadId}'.`; + } +} + +export class ThreadReadAmbiguousError extends Schema.TaggedErrorClass()( + "ThreadReadAmbiguousError", + { + threadId: Schema.String, + matchCount: Schema.Int, + }, +) { + override get message(): string { + return `Provider thread ID '${this.threadId}' matches multiple T3 Code threads; use a T3 thread ID instead.`; + } +} + +export class ThreadReadFailedError extends Schema.TaggedErrorClass()( + "ThreadReadFailedError", + { + operation: Schema.String, + }, +) { + override get message(): string { + return `T3 Code could not read the requested thread during ${this.operation}.`; + } +} + +export const ThreadReadError = Schema.Union([ + ThreadReadUnavailableError, + ThreadReadNotFoundError, + ThreadReadAmbiguousError, + ThreadReadFailedError, +]); +export type ThreadReadError = typeof ThreadReadError.Type; + +const dependencies = [ + McpInvocationContext.McpInvocationContext, + ProjectionSnapshotQuery, + ProviderSessionDirectory, +]; + +export const ReadThreadTool = Tool.make("read_thread", { + description: + "Read another T3 Code conversation by its T3 thread ID or copied provider-native thread ID. Returns current status plus bounded recent messages and activity summaries without resuming, steering, or modifying the target thread.", + parameters: ReadThreadInput, + success: ReadThreadResult, + failure: ThreadReadError, + dependencies, +}) + .annotate(Tool.Title, "Read T3 Code thread") + .annotate(Tool.Readonly, true) + .annotate(Tool.Destructive, false) + .annotate(Tool.Idempotent, true); + +export const ThreadControlToolkit = Toolkit.make(ReadThreadTool); diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 1f24a4a0200..14cd60d0ae2 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -1068,6 +1068,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti status: event.payload.session.status, providerName: event.payload.session.providerName, providerInstanceId: event.payload.session.providerInstanceId ?? null, + providerThreadId: event.payload.session.providerThreadId ?? null, runtimeMode: event.payload.session.runtimeMode, activeTurnId: event.payload.session.activeTurnId, lastError: event.payload.session.lastError, diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index d4a24a209ad..69fc443e366 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -361,6 +361,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { threadId: ThreadId.make("thread-1"), status: "running", providerName: "codex", + providerThreadId: "provider-thread-1", runtimeMode: "approval-required", activeTurnId: asTurnId("turn-1"), lastError: null, @@ -430,6 +431,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { threadId: ThreadId.make("thread-1"), status: "running", providerName: "codex", + providerThreadId: "provider-thread-1", runtimeMode: "approval-required", activeTurnId: asTurnId("turn-1"), lastError: null, diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index 3d05bef4bdf..1172426b411 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -217,6 +217,7 @@ function mapSessionRow( status: row.status, providerName: row.providerName, ...(row.providerInstanceId !== null ? { providerInstanceId: row.providerInstanceId } : {}), + ...(row.providerThreadId ? { providerThreadId: row.providerThreadId } : {}), runtimeMode: row.runtimeMode, activeTurnId: row.activeTurnId, lastError: row.lastError, @@ -858,6 +859,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { status, provider_name AS "providerName", provider_instance_id AS "providerInstanceId", + provider_thread_id AS "providerThreadId", runtime_mode AS "runtimeMode", active_turn_id AS "activeTurnId", last_error AS "lastError", @@ -1165,6 +1167,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ...(row.providerInstanceId !== null ? { providerInstanceId: row.providerInstanceId } : {}), + ...(row.providerThreadId ? { providerThreadId: row.providerThreadId } : {}), runtimeMode: row.runtimeMode, activeTurnId: row.activeTurnId, lastError: row.lastError, diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index b6bff8c766a..b363b13e0cb 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -442,6 +442,9 @@ const make = Effect.gen(function* () { status: "starting", providerName: activeSession?.provider ?? preferredProvider, providerInstanceId: activeSession?.providerInstanceId ?? desiredInstanceId, + ...(thread.session?.providerThreadId + ? { providerThreadId: thread.session.providerThreadId } + : {}), runtimeMode: desiredRuntimeMode, activeTurnId: null, lastError: null, @@ -526,6 +529,9 @@ const make = Effect.gen(function* () { : mapProviderSessionStatusToOrchestrationStatus(session.status), providerName: session.provider, providerInstanceId: session.providerInstanceId, + ...(thread.session?.providerThreadId + ? { providerThreadId: thread.session.providerThreadId } + : {}), runtimeMode: desiredRuntimeMode, // Provider turn ids are not orchestration turn ids. activeTurnId: null, @@ -1017,6 +1023,9 @@ const make = Effect.gen(function* () { ...(thread.session?.providerInstanceId !== undefined ? { providerInstanceId: thread.session.providerInstanceId } : {}), + ...(thread.session?.providerThreadId + ? { providerThreadId: thread.session.providerThreadId } + : {}), runtimeMode: thread.session?.runtimeMode ?? DEFAULT_RUNTIME_MODE, activeTurnId: null, lastError: thread.session?.lastError ?? null, diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 74ece50cd31..1c7a48d16e4 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -243,10 +243,13 @@ describe("ProviderRuntimeIngestion", () => { Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(NodeServices.layer), ); - runtime = ManagedRuntime.make(layer); - const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); - const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); - const ingestion = await runtime.runPromise(Effect.service(ProviderRuntimeIngestionService)); + const activeRuntime = ManagedRuntime.make(layer); + runtime = activeRuntime; + const engine = await activeRuntime.runPromise(Effect.service(OrchestrationEngineService)); + const snapshotQuery = await activeRuntime.runPromise(Effect.service(ProjectionSnapshotQuery)); + const ingestion = await activeRuntime.runPromise( + Effect.service(ProviderRuntimeIngestionService), + ); scope = await Effect.runPromise(Scope.make("sequential")); await Effect.runPromise(ingestion.start().pipe(Scope.provide(scope))); const drain = () => Effect.runPromise(ingestion.drain); @@ -313,6 +316,7 @@ describe("ProviderRuntimeIngestion", () => { return { engine, readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), + readEvents: () => activeRuntime.runPromise(Stream.runCollect(engine.readEvents(0))), emit: provider.emit, setProviderSession: provider.setSession, drain, @@ -2724,7 +2728,19 @@ describe("ProviderRuntimeIngestion", () => { provider: ProviderDriverKind.make("codex"), createdAt: now, threadId: asThreadId("thread-1"), + payload: { + providerThreadId: "019c-native-codex-thread", + }, }); + await harness.drain(); + const events = Array.from(await harness.readEvents()); + expect( + events.some( + (event) => + event.type === "thread.session-set" && + event.payload.session.providerThreadId === "019c-native-codex-thread", + ), + ).toBe(true); harness.emit({ type: "item.started", eventId: asEventId("evt-tool-started"), diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index a8a51b30260..39a91eabeb0 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1412,6 +1412,10 @@ const make = Effect.gen(function* () { : (thread.session?.lastError ?? null); if (shouldApplyThreadLifecycle) { + const providerThreadId = + event.type === "thread.started" + ? (event.payload.providerThreadId ?? thread.session?.providerThreadId) + : thread.session?.providerThreadId; if (event.type === "turn.started" && acceptedTurnStartedSourcePlan !== null) { yield* markSourceProposedPlanImplemented( acceptedTurnStartedSourcePlan.sourceThreadId, @@ -1443,6 +1447,7 @@ const make = Effect.gen(function* () { ...(event.providerInstanceId !== undefined ? { providerInstanceId: event.providerInstanceId } : {}), + ...(providerThreadId ? { providerThreadId } : {}), runtimeMode: thread.session?.runtimeMode ?? "full-access", activeTurnId: nextActiveTurnId, lastError, @@ -1693,6 +1698,9 @@ const make = Effect.gen(function* () { ...(event.providerInstanceId !== undefined ? { providerInstanceId: event.providerInstanceId } : {}), + ...(thread.session?.providerThreadId + ? { providerThreadId: thread.session.providerThreadId } + : {}), runtimeMode: thread.session?.runtimeMode ?? "full-access", activeTurnId: eventTurnId ?? null, lastError: runtimeErrorMessage, diff --git a/apps/server/src/persistence/Layers/ProjectionThreadSessions.ts b/apps/server/src/persistence/Layers/ProjectionThreadSessions.ts index dcb750983a0..8e864ff7079 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadSessions.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadSessions.ts @@ -25,6 +25,7 @@ const makeProjectionThreadSessionRepository = Effect.gen(function* () { status, provider_name, provider_instance_id, + provider_thread_id, runtime_mode, active_turn_id, last_error, @@ -35,6 +36,7 @@ const makeProjectionThreadSessionRepository = Effect.gen(function* () { ${row.status}, ${row.providerName}, ${row.providerInstanceId}, + ${row.providerThreadId ?? null}, ${row.runtimeMode}, ${row.activeTurnId}, ${row.lastError}, @@ -45,6 +47,7 @@ const makeProjectionThreadSessionRepository = Effect.gen(function* () { status = excluded.status, provider_name = excluded.provider_name, provider_instance_id = excluded.provider_instance_id, + provider_thread_id = excluded.provider_thread_id, runtime_mode = excluded.runtime_mode, active_turn_id = excluded.active_turn_id, last_error = excluded.last_error, @@ -62,6 +65,7 @@ const makeProjectionThreadSessionRepository = Effect.gen(function* () { status, provider_name AS "providerName", provider_instance_id AS "providerInstanceId", + provider_thread_id AS "providerThreadId", runtime_mode AS "runtimeMode", active_turn_id AS "activeTurnId", last_error AS "lastError", diff --git a/apps/server/src/persistence/Services/ProjectionThreadSessions.ts b/apps/server/src/persistence/Services/ProjectionThreadSessions.ts index 7cecac33eb6..d11feb03503 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadSessions.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadSessions.ts @@ -26,6 +26,7 @@ export const ProjectionThreadSession = Schema.Struct({ status: OrchestrationSessionStatus, providerName: Schema.NullOr(Schema.String), providerInstanceId: Schema.NullOr(ProviderInstanceId), + providerThreadId: Schema.NullOr(Schema.String), runtimeMode: RuntimeMode, activeTurnId: Schema.NullOr(TurnId), lastError: Schema.NullOr(Schema.String), diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index a1d95eaa734..2e7f8ddbc73 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -198,6 +198,7 @@ import { getProjectOrderKey, selectProjectGroupingSettings, } from "../logicalProject"; +import { providerThreadCopyAction } from "../threadReferences"; import type { SidebarThreadSummary } from "../types"; import { buildPhysicalToLogicalProjectKeyMap, @@ -1143,6 +1144,27 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec ); }, }); + const { copyToClipboard: copyProviderThreadIdToClipboard } = useCopyToClipboard<{ + threadId: string; + label: string; + }>({ + onCopy: (ctx) => { + toastManager.add({ + type: "success", + title: `${ctx.label} copied`, + description: ctx.threadId, + }); + }, + onError: (error) => { + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Failed to copy provider thread ID", + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + }, + }); const { copyToClipboard: copyPathToClipboard } = useCopyToClipboard<{ path: string; }>({ @@ -2112,6 +2134,7 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec ); const threadWorkspacePath = thread.worktreePath ?? threadProject?.workspaceRoot ?? project.workspaceRoot ?? null; + const providerCopyAction = providerThreadCopyAction(thread.session); const clicked = await api.contextMenu.show( [ ...(thread.branch @@ -2121,6 +2144,9 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec { id: "mark-unread", label: "Mark unread" }, { id: "copy-path", label: "Copy Path" }, { id: "copy-thread-id", label: "Copy Thread ID" }, + ...(providerCopyAction + ? [{ id: providerCopyAction.id, label: providerCopyAction.label }] + : []), { id: "delete", label: "Delete", destructive: true, icon: "trash" }, ], position, @@ -2177,6 +2203,13 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec copyThreadIdToClipboard(thread.id, { threadId: thread.id }); return; } + if (clicked === "copy-provider-thread-id" && providerCopyAction) { + copyProviderThreadIdToClipboard(providerCopyAction.value, { + threadId: providerCopyAction.value, + label: providerCopyAction.label, + }); + return; + } if (clicked !== "delete") return; if (appSettingsConfirmThreadDelete) { const confirmed = await api.dialogs.confirm( @@ -2204,6 +2237,7 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec [ appSettingsConfirmThreadDelete, copyPathToClipboard, + copyProviderThreadIdToClipboard, copyThreadIdToClipboard, deleteThread, handleNewThread, diff --git a/apps/web/src/components/SidebarV2.tsx b/apps/web/src/components/SidebarV2.tsx index 8c5891ebe7e..1a65677bc9f 100644 --- a/apps/web/src/components/SidebarV2.tsx +++ b/apps/web/src/components/SidebarV2.tsx @@ -85,6 +85,7 @@ import { useThreadActions } from "../hooks/useThreadActions"; import { useHandleNewThread } from "../hooks/useHandleNewThread"; import { openCommandPalette } from "../commandPaletteBus"; import { startNewThreadFromContext } from "../lib/chatThreadActions"; +import { providerThreadCopyAction } from "../threadReferences"; import { useClientSettings, useUpdateClientSettings } from "../hooks/useSettings"; import { useCopyToClipboard } from "../hooks/useCopyToClipboard"; import { useNowMinute } from "../hooks/useNowMinute"; @@ -1036,6 +1037,45 @@ export default function SidebarV2() { ); }, }); + const { copyToClipboard: copyThreadId } = useCopyToClipboard<{ threadId: string }>({ + onCopy: ({ threadId }) => { + toastManager.add({ + type: "success", + title: "Thread ID copied", + description: threadId, + }); + }, + onError: (error) => { + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Failed to copy thread ID", + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + }, + }); + const { copyToClipboard: copyProviderThreadId } = useCopyToClipboard<{ + threadId: string; + label: string; + }>({ + onCopy: ({ threadId, label }) => { + toastManager.add({ + type: "success", + title: `${label} copied`, + description: threadId, + }); + }, + onError: (error) => { + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Failed to copy provider thread ID", + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + }, + }); const [projectActionsTarget, setProjectActionsTarget] = useState( null, ); @@ -1980,6 +2020,7 @@ export default function SidebarV2() { serverConfigs.get(thread.environmentId)?.environment.capabilities.threadSnooze === true; const isSettled = settledThreadKeysRef.current.has(threadKey); const isSnoozed = snoozedThreadKeysRef.current.has(threadKey); + const providerCopyAction = providerThreadCopyAction(thread.session); // Presets resolve at menu-open time (same as the popover). const snoozePresets = resolveSnoozePresets(new Date()); const clicked = await settlePromise(() => @@ -2017,6 +2058,10 @@ export default function SidebarV2() { : []), { id: "rename", label: "Rename thread" }, { id: "mark-unread", label: "Mark unread" }, + { id: "copy-thread-id", label: "Copy Thread ID" }, + ...(providerCopyAction + ? [{ id: providerCopyAction.id, label: providerCopyAction.label }] + : []), { id: "delete", label: "Delete", destructive: true, icon: "trash" }, ], position, @@ -2069,6 +2114,17 @@ export default function SidebarV2() { case "mark-unread": markThreadUnread(threadKey, thread.latestTurn?.completedAt); return; + case "copy-thread-id": + copyThreadId(thread.id, { threadId: thread.id }); + return; + case "copy-provider-thread-id": + if (providerCopyAction) { + copyProviderThreadId(providerCopyAction.value, { + threadId: providerCopyAction.value, + label: providerCopyAction.label, + }); + } + return; case "delete": { if (confirmThreadDelete) { const confirmed = await settlePromise(() => @@ -2106,6 +2162,8 @@ export default function SidebarV2() { attemptUnsettle, attemptUnsnooze, confirmThreadDelete, + copyProviderThreadId, + copyThreadId, deleteThread, handleMultiSelectContextMenu, markThreadUnread, diff --git a/apps/web/src/threadReferences.test.ts b/apps/web/src/threadReferences.test.ts new file mode 100644 index 00000000000..99ee0e14e39 --- /dev/null +++ b/apps/web/src/threadReferences.test.ts @@ -0,0 +1,47 @@ +import { describe, expect, it } from "vite-plus/test"; + +import { providerThreadCopyAction } from "./threadReferences"; + +const baseSession = { + threadId: "t3-thread-id", + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: "2026-07-26T00:00:00.000Z", +} as const; + +describe("providerThreadCopyAction", () => { + it("keeps the existing T3 ID separate and labels Codex IDs explicitly", () => { + expect( + providerThreadCopyAction({ + ...baseSession, + providerThreadId: "codex-thread-id", + }), + ).toEqual({ + id: "copy-provider-thread-id", + label: "Copy Codex Thread ID", + value: "codex-thread-id", + }); + }); + + it("uses a provider-neutral label outside Codex", () => { + expect( + providerThreadCopyAction({ + ...baseSession, + providerName: "claude", + providerThreadId: "claude-session-id", + }), + ).toEqual({ + id: "copy-provider-thread-id", + label: "Copy Provider Thread ID", + value: "claude-session-id", + }); + }); + + it("does not add an action for legacy sessions without a provider ID", () => { + expect(providerThreadCopyAction(baseSession)).toBeNull(); + expect(providerThreadCopyAction(null)).toBeNull(); + }); +}); diff --git a/apps/web/src/threadReferences.ts b/apps/web/src/threadReferences.ts new file mode 100644 index 00000000000..8e515c7e3ff --- /dev/null +++ b/apps/web/src/threadReferences.ts @@ -0,0 +1,22 @@ +import type { OrchestrationSession } from "@t3tools/contracts"; + +export interface ProviderThreadCopyAction { + readonly id: "copy-provider-thread-id"; + readonly label: string; + readonly value: string; +} + +export function providerThreadCopyAction( + session: Pick | null, +): ProviderThreadCopyAction | null { + const providerThreadId = session?.providerThreadId?.trim(); + const providerName = session?.providerName; + if (!providerThreadId) { + return null; + } + return { + id: "copy-provider-thread-id", + label: providerName === "codex" ? "Copy Codex Thread ID" : "Copy Provider Thread ID", + value: providerThreadId, + }; +} diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 84b7a8fa07f..759fc55af17 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -274,6 +274,7 @@ export const OrchestrationSession = Schema.Struct({ status: OrchestrationSessionStatus, providerName: Schema.NullOr(TrimmedNonEmptyString), providerInstanceId: Schema.optional(ProviderInstanceId), + providerThreadId: Schema.optional(TrimmedNonEmptyString), runtimeMode: RuntimeMode.pipe(Schema.withDecodingDefault(Effect.succeed(DEFAULT_RUNTIME_MODE))), activeTurnId: Schema.NullOr(TurnId), lastError: Schema.NullOr(TrimmedNonEmptyString),