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
147 changes: 147 additions & 0 deletions apps/server/src/mcp/AttachmentMcpService.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,147 @@
import {
AttachmentMcpFailure,
IsoDateTime,
PROVIDER_SEND_TURN_SUPPORTED_IMAGE_MIME_TYPES,
type AttachmentMcpDiscardUploadInput,
type AttachmentMcpDiscardUploadResult,
type AttachmentMcpPrepareUploadInput,
type AttachmentMcpPrepareUploadResult,
} from "@t3tools/contracts";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";

import {
deletePendingAttachment,
issueAttachmentUploadUrl,
validateAttachmentUploadToken,
} from "../assets/AttachmentUpload.ts";
import {
parseThreadSegmentFromAttachmentId,
PENDING_ATTACHMENT_THREAD_SEGMENT,
} from "../attachmentStore.ts";
import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import * as ServerConfig from "../config.ts";
import type { McpInvocationScope } from "./McpInvocationContext.ts";

export class AttachmentMcpService extends Context.Service<
AttachmentMcpService,
{
readonly prepareUpload: (
scope: McpInvocationScope,
input: AttachmentMcpPrepareUploadInput,
) => Effect.Effect<AttachmentMcpPrepareUploadResult, AttachmentMcpFailure>;
readonly discardUpload: (
scope: McpInvocationScope,
input: AttachmentMcpDiscardUploadInput,
) => Effect.Effect<AttachmentMcpDiscardUploadResult, AttachmentMcpFailure>;
}
>()("t3/mcp/AttachmentMcpService") {}

function failure(code: AttachmentMcpFailure["code"], message: string): AttachmentMcpFailure {
return new AttachmentMcpFailure({ code, message });
}

const requireCapability = (scope: McpInvocationScope) =>
scope.capabilities.has("orchestration")
? Effect.void
: Effect.fail(
failure(
"capability_denied",
"This MCP credential does not grant orchestration capabilities.",
),
);

export const make = Effect.gen(function* () {
const config = yield* ServerConfig.ServerConfig;
const secretStore = yield* ServerSecretStore.ServerSecretStore;
const fileSystem = yield* FileSystem.FileSystem;
const issueUpload = (input: Parameters<typeof issueAttachmentUploadUrl>[0]) =>
issueAttachmentUploadUrl(input).pipe(
Effect.provideService(ServerConfig.ServerConfig, config),
Effect.provideService(ServerSecretStore.ServerSecretStore, secretStore),
);
const discardPending = (attachmentId: string) =>
deletePendingAttachment(attachmentId).pipe(
Effect.provideService(ServerConfig.ServerConfig, config),
Effect.provideService(FileSystem.FileSystem, fileSystem),
);
const validateUploadToken = (token: string) =>
validateAttachmentUploadToken(token).pipe(
Effect.provideService(ServerSecretStore.ServerSecretStore, secretStore),
);

return AttachmentMcpService.of({
prepareUpload: (scope, input) =>
Effect.gen(function* () {
yield* requireCapability(scope);
const uploadInput =
input.type === "file"
? {
type: "file" as const,
name: input.name,
mimeType: input.mimeType,
sizeBytes: input.sizeBytes,
}
: {
type: "image" as const,
name: input.name,
mimeType:
input.mimeType as (typeof PROVIDER_SEND_TURN_SUPPORTED_IMAGE_MIME_TYPES)[number],
sizeBytes: input.sizeBytes,
};
const issued = yield* issueUpload(uploadInput).pipe(
Effect.mapError(() =>
failure("upload_error", "Unable to prepare a signed attachment upload."),
),
);
const attachment = {
id: issued.attachmentId,
type: input.type ?? "image",
name: input.name,
mimeType: input.mimeType,
sizeBytes: input.sizeBytes,
} as const;
return {
attachmentId: issued.attachmentId,
attachment,
type: input.type ?? "image",
name: input.name,
mimeType: input.mimeType,
sizeBytes: input.sizeBytes,
upload: {
method: "POST",
relativeUrl: issued.relativeUrl,
expiresAt: IsoDateTime.make(DateTime.formatIso(DateTime.makeUnsafe(issued.expiresAt))),
},
};
}),
discardUpload: (scope, input) =>
Effect.gen(function* () {
yield* requireCapability(scope);
if (
parseThreadSegmentFromAttachmentId(input.attachmentId) !==
PENDING_ATTACHMENT_THREAD_SEGMENT
) {
return yield* failure(
"invalid_attachment",
"Only a pending upload can be discarded; thread-owned attachments are immutable here.",
);
}
const token = input.uploadRelativeUrl.split("/").at(-1) ?? "";
const claims = yield* validateUploadToken(token);
if (claims?.attachmentId !== input.attachmentId) {
return yield* failure(
"invalid_attachment",
"The signed upload URL does not authorize this pending attachment id.",
);
}
yield* discardPending(input.attachmentId);
Comment thread
juliusmarminge marked this conversation as resolved.
return { attachmentId: input.attachmentId, discarded: true };
}),
});
});

export const layer = Layer.effect(AttachmentMcpService, make);
9 changes: 9 additions & 0 deletions apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import { McpProtocol, McpSchema, McpServer, Tool } from "effect/unstable/ai";
import { HttpRouter, HttpServerRequest, HttpServerResponse } from "effect/unstable/http";

import packageJson from "../../package.json" with { type: "json" };
import * as AttachmentMcpService from "./AttachmentMcpService.ts";
import * as McpInvocationContext from "./McpInvocationContext.ts";
import * as OrchestratorMcpService from "./OrchestratorMcpService.ts";
import * as ThreadMetadataMcpService from "./ThreadMetadataMcpService.ts";
Expand All @@ -29,6 +30,8 @@ import {
import { WorktreeToolkitHandlersLive } from "./toolkits/worktree/handlers.ts";
import { WorktreeToolkit } from "./toolkits/worktree/tools.ts";
import * as WorktreeMcpService from "./WorktreeMcpService.ts";
import { AttachmentToolkitHandlersLive } from "./toolkits/attachment/handlers.ts";
import { AttachmentToolkit } from "./toolkits/attachment/tools.ts";

const unauthorized = HttpServerResponse.jsonUnsafe(
{
Expand Down Expand Up @@ -233,6 +236,11 @@ export const WorktreeToolkitRegistrationLive = McpServer.toolkit(WorktreeToolkit
Layer.provide(WorktreeMcpService.layer),
);

export const AttachmentToolkitRegistrationLive = McpServer.toolkit(AttachmentToolkit).pipe(
Layer.provide(AttachmentToolkitHandlersLive),
Layer.provide(AttachmentMcpService.layer),
);

const McpTransportLive = McpServer.layerHttp({
name: "T3 Code",
version: packageJson.version,
Expand All @@ -244,4 +252,5 @@ export const layer = Layer.mergeAll(
PreviewToolkitRegistrationLive,
OrchestratorToolkitRegistrationLive,
WorktreeToolkitRegistrationLive,
AttachmentToolkitRegistrationLive,
).pipe(Layer.provideMerge(McpTransportLive));
142 changes: 142 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpAttachments.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
// @effect-diagnostics nodeBuiltinImport:off
import * as NodeFS from "node:fs";
import * as NodePath from "node:path";

import * as NodeServices from "@effect/platform-node/NodeServices";
import { expect, it } from "@effect/vitest";
import {
ChatAttachmentId,
EnvironmentId,
ProjectId,
ProviderInstanceId,
type ServerProvider,
ThreadId,
type OrchestrationV2ThreadProjection,
} from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Ref from "effect/Ref";

import { createPendingAttachmentId } from "../attachmentStore.ts";
import * as ServerConfig from "../config.ts";
import {
ThreadManagementPostDispatchProjectionError,
ThreadManagementService,
} from "../orchestration-v2/ThreadManagementService.ts";
import { ProviderRegistry } from "../provider/Services/ProviderRegistry.ts";
import { ScheduledTaskService } from "../scheduledTasks/ScheduledTaskService.ts";
import type { McpInvocationScope } from "./McpInvocationContext.ts";
import * as OrchestratorMcpService from "./OrchestratorMcpService.ts";

it.effect(
"retains fresh accepted claims and releases unused replay claims after projection errors",
() =>
Effect.gen(function* () {
const projectId = ProjectId.make("project:mcp-attachment-cleanup");
const threadId = ThreadId.make("thread:mcp-attachment-cleanup");
const projection = {
thread: {
id: threadId,
projectId,
runtimeMode: "full-access",
interactionMode: "default",
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5.6-sol",
},
deletedAt: null,
},
messages: [],
} as unknown as OrchestrationV2ThreadProjection;
const dispatchReplayed = yield* Ref.make(false);
const configLayer = ServerConfig.layerTest(process.cwd(), {
prefix: "t3-mcp-attachment-cleanup-",
}).pipe(Layer.provide(NodeServices.layer));
const dependencies = Layer.mergeAll(
NodeServices.layer,
configLayer,
Layer.mock(ThreadManagementService)({
getCommandReceipt: () => Effect.succeed(Option.none()),
getThreadProjection: () => Effect.succeed(projection),
sendToThread: (input) =>
Ref.get(dispatchReplayed).pipe(
Effect.flatMap((replayed) =>
Effect.fail(
new ThreadManagementPostDispatchProjectionError({
projectId,
threadId,
messageId: input.messageId,
dispatchReplayed: replayed,
cause: new Error("projection unavailable after accepted dispatch"),
}),
),
),
),
}),
Layer.mock(ProviderRegistry)({
getProviders: Effect.succeed([
{
instanceId: ProviderInstanceId.make("codex"),
driver: "codex",
} as unknown as ServerProvider,
]),
}),
Layer.mock(ScheduledTaskService)({}),
);
const testLayer = OrchestratorMcpService.layer.pipe(Layer.provideMerge(dependencies));

yield* Effect.gen(function* () {
const service = yield* OrchestratorMcpService.OrchestratorMcpService;
const config = yield* ServerConfig.ServerConfig;
const scope: McpInvocationScope = {
environmentId: EnvironmentId.make("environment:mcp-attachment-cleanup"),
threadId,
providerSessionId: "provider-session:mcp-attachment-cleanup",
providerInstanceId: ProviderInstanceId.make("codex"),
capabilities: new Set(["orchestration"]),
issuedAt: 1,
};
const stage = (name: string) => {
const id = createPendingAttachmentId();
if (id === null) throw new Error("Expected a pending attachment id.");
NodeFS.writeFileSync(
NodePath.join(config.attachmentsDir, `${id}.png`),
Buffer.from([1, 2, 3, 4]),
);
return {
type: "image" as const,
id: ChatAttachmentId.make(id),
name,
mimeType: "image/png",
sizeBytes: 4,
};
};
const claimedFiles = () =>
NodeFS.readdirSync(config.attachmentsDir).filter(
(entry) => !entry.startsWith("pending-"),
);

const freshError = yield* service
.sendToThread(scope, {
threadId,
attachments: [stage("fresh.png")],
clientRequestId: "fresh-accepted-claim",
})
.pipe(Effect.flip);
expect(freshError.code).toBe("orchestration_error");
expect(claimedFiles()).toHaveLength(1);

yield* Ref.set(dispatchReplayed, true);
const replayError = yield* service
.sendToThread(scope, {
threadId,
attachments: [stage("replay.png")],
clientRequestId: "replayed-accepted-claim",
})
.pipe(Effect.flip);
expect(replayError.code).toBe("orchestration_error");
expect(claimedFiles()).toHaveLength(1);
}).pipe(Effect.provide(testLayer));
}),
);
14 changes: 14 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpService.activity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,11 +12,13 @@ import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as NodeCrypto from "@effect/platform-node/NodeCrypto";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { expect, it } from "vite-plus/test";

import { ProviderRegistry } from "../provider/Services/ProviderRegistry.ts";
import { ScheduledTaskService } from "../scheduledTasks/ScheduledTaskService.ts";
import { ThreadManagementService } from "../orchestration-v2/ThreadManagementService.ts";
import * as ServerConfig from "../config.ts";
import type * as McpInvocationContext from "./McpInvocationContext.ts";
import {
layer as orchestratorMcpServiceLayer,
Expand Down Expand Up @@ -136,6 +138,10 @@ it("readThread prefers activity-run status over a newer cancelled queued run", a
list: () => Effect.succeed({ tasks: [] }),
} satisfies Partial<ScheduledTaskService["Service"]>),
NodeCrypto.layer,
NodeServices.layer,
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-mcp-orchestrator-activity-",
}).pipe(Layer.provide(NodeServices.layer)),
),
),
);
Expand Down Expand Up @@ -185,6 +191,10 @@ it("readThread prefers waiting activity status over a newer cancelled queued run
list: () => Effect.succeed({ tasks: [] }),
} satisfies Partial<ScheduledTaskService["Service"]>),
NodeCrypto.layer,
NodeServices.layer,
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-mcp-orchestrator-activity-",
}).pipe(Layer.provide(NodeServices.layer)),
),
),
);
Expand Down Expand Up @@ -291,6 +301,10 @@ it("taskStatus returns task.providerInstanceId rather than the driver kind", asy
list: () => Effect.succeed({ tasks: [] }),
} satisfies Partial<ScheduledTaskService["Service"]>),
NodeCrypto.layer,
NodeServices.layer,
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-mcp-orchestrator-activity-",
}).pipe(Layer.provide(NodeServices.layer)),
),
),
);
Expand Down
Loading
Loading