Skip to content
Merged
1 change: 1 addition & 0 deletions apps/server/src/auth/RpcAuthorization.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,7 @@ describe("RPC authorization scopes", () => {
WS_METHODS.previewReportStatus,
WS_METHODS.previewAdjust,
WS_METHODS.previewClearProfile,
WS_METHODS.previewReportProfiles,
]) {
expect(requiredScopeForRpcMethod(method)).toBe(AuthPreviewOperateScope);
}
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,7 @@ export const RPC_REQUIRED_SCOPES = {
[WS_METHODS.previewClose]: AuthPreviewOperateScope,
[WS_METHODS.previewList]: AuthOrchestrationReadScope,
[WS_METHODS.previewClearProfile]: AuthPreviewOperateScope,
[WS_METHODS.previewReportProfiles]: AuthPreviewOperateScope,
[WS_METHODS.previewReportStatus]: AuthPreviewOperateScope,
[WS_METHODS.subscribePreviewEvents]: AuthOrchestrationReadScope,
[WS_METHODS.subscribeDiscoveredLocalServers]: AuthOrchestrationReadScope,
Expand Down
23 changes: 22 additions & 1 deletion apps/server/src/mcp/McpDeviceToolkit.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { expect, it } from "@effect/vitest";
import * as NodeServices from "@effect/platform-node/NodeServices";
import {
DeviceHostUnavailableError,
DeviceId,
EnvironmentId,
ProviderInstanceId,
ThreadId,
Expand Down Expand Up @@ -93,7 +94,21 @@ const layerDeviceServiceMock = Layer.mock(DeviceService.DeviceService)({
platform: input.platform,
openedAt: "2026-09-08T00:00:00.000Z",
}),
sessionsForThread: () => Effect.succeed([]),
// UDID-1 is open in the test thread; UDID-2 exists but belongs to another thread.
sessionsForThread: (id) =>
Effect.succeed(
id === threadId
? [
{
threadId,
hostId: "local",
deviceId: DeviceId.make("UDID-1"),
platform: "ios" as const,
openedAt: "2026-09-08T00:00:00.000Z",
},
]
: [],
),
screenshot: () => Effect.succeed({ device, png }),
close: () => Effect.void,
agentCli: Effect.succeed("/cli"),
Expand Down Expand Up @@ -135,6 +150,12 @@ it.effect("registers the device tools and returns the screenshot as image conten
screenshot: { mimeType: "image/png", width: 1206, height: 2622 },
});

const foreign = yield* server
.callTool({ name: "device_screenshot", arguments: { deviceId: "UDID-2" } })
.pipe(callWith(["device"]), Effect.provideService(McpSchema.McpServerClient, client));
expect(foreign.isError).toBe(true);
expect(foreign.content.map((entry) => entry.type)).toEqual(["text"]);

const denied = yield* server
.callTool({ name: "device_list", arguments: {} })
.pipe(callWith(["preview"]), Effect.provideService(McpSchema.McpServerClient, client));
Expand Down
23 changes: 11 additions & 12 deletions apps/server/src/mcp/McpHttpServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,7 @@ it.effect("tells the agent how to fall back when no desktop app can run the snap
);

it.effect.each([
{ mode: "default", input: {}, images: true },
{ mode: "default", input: {}, images: false },
{ mode: "explicit image", input: { includeImage: true }, images: true },
{ mode: "text only", input: { includeImage: false }, images: false },
])("returns fresh $mode snapshots on repeated MCP calls", ({ input, images }) =>
Expand Down Expand Up @@ -351,12 +351,7 @@ it.effect.each([
Effect.provideService(McpInvocationContext.McpInvocationContext, invocation),
Effect.provideService(McpSchema.McpServerClient, client),
);
expect(nextDefault.content.map((content) => content.type)).toEqual([
"text",
"text",
"text",
"image",
]);
expect(nextDefault.content.map((content) => content.type)).toEqual(["text", "text", "text"]);
expect(nextDefault.structuredContent).toMatchObject({ title: "Snapshot 7", screenshot });
expect(nextDefault.structuredContent).not.toHaveProperty("accessibilityTree");
expect(requests).toBe(7);
Expand Down Expand Up @@ -416,11 +411,12 @@ it.effect("saves the snapshot PNG on request and reports its path", () =>
const path = yield* Path.Path;
const inputs = yield* serveSnapshots("mcp-save-client", snapshotResult);

const snapshot = yield* callSnapshot({ save: true });
const snapshot = yield* callSnapshot({ save: true, includeImage: true });

expect(snapshot.isError).toBe(false);
// The browser never receives the server-only `save` flag.
expect(inputs).toEqual([{}]);
expect(snapshot.content.map((content) => content.type)).toContain("image");
const structured = snapshot.structuredContent as { readonly screenshotPath?: string };
const screenshotPath = structured.screenshotPath;
expect(typeof screenshotPath).toBe("string");
Expand All @@ -436,7 +432,7 @@ it.effect("saves the snapshot PNG on request and reports its path", () =>
expect(unsaved.structuredContent).not.toHaveProperty("screenshotPath");

// A save without the image skips the page dump.
const pathOnly = yield* callSnapshot({ save: true, includeImage: false });
const pathOnly = yield* callSnapshot({ save: true });
const saved = pathOnly.structuredContent as { readonly screenshotPath: string };
expect(saved).toEqual({ url: snapshotResult.url, screenshotPath: expect.any(String) });
expect(Buffer.from(yield* fileSystem.readFile(saved.screenshotPath)).toString()).toBe("png");
Expand Down Expand Up @@ -851,8 +847,8 @@ it.effect("registers annotated tools and preserves authenticated request context
expect(statusTool?.tool.annotations?.destructiveHint).toBe(false);

const snapshotTool = server.tools.find(({ tool }) => tool.name === "preview_snapshot");
expect(snapshotTool?.tool.annotations?.readOnlyHint).toBe(true);
expect(snapshotTool?.tool.annotations?.idempotentHint).toBe(true);
expect(snapshotTool?.tool.annotations?.readOnlyHint).toBe(false);
expect(snapshotTool?.tool.annotations?.idempotentHint).toBe(false);
expect(snapshotTool?.tool.annotations?.openWorldHint).toBe(true);

const clickTool = server.tools.find(({ tool }) => tool.name === "preview_click");
Expand Down Expand Up @@ -891,7 +887,10 @@ it.effect("registers annotated tools and preserves authenticated request context
expect(malformed._tag).toBe("InvalidParams");

const snapshot = yield* server
.callTool({ name: "preview_snapshot", arguments: { tabId: alternateTabId } })
.callTool({
name: "preview_snapshot",
arguments: { tabId: alternateTabId, includeImage: true },
})
.pipe(
Effect.provideService(McpInvocationContext.McpInvocationContext, invocation),
Effect.provideService(McpSchema.McpServerClient, client),
Expand Down
15 changes: 9 additions & 6 deletions apps/server/src/mcp/McpHttpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -503,7 +503,10 @@ const registerPreviewSnapshot = Effect.fn("McpHttpServer.registerPreviewSnapshot
const png = new Uint8Array(Buffer.from(screenshot.data, "base64"));
const screenshotPath =
payload?.save === true ? yield* saveScreenshot(snapshot.url, png) : undefined;
if (screenshotPath !== undefined && payload?.includeImage === false) {
// Images stay out of tool history unless asked for: providers replay them on every
// later request, and some reject inline images outright.
const includeImage = payload?.includeImage === true;
if (screenshotPath !== undefined && !includeImage) {
// The agent only wants a file to show the user. The url keeps the site icon on the tool row.
const saved = {
url: cutText(snapshot.url, MAX_SNAPSHOT_IDENTIFIER_CHARS),
Expand Down Expand Up @@ -548,9 +551,9 @@ const registerPreviewSnapshot = Effect.fn("McpHttpServer.registerPreviewSnapshot
text: `Snapshot text was bounded. Omitted: ${bounded.omitted.join("; ")}.`,
},
]),
...(payload?.includeImage === false
? []
: [{ type: "image" as const, data: png, mimeType: screenshot.mimeType }]),
...(includeImage
? [{ type: "image" as const, data: png, mimeType: screenshot.mimeType }]
: []),
],
});
}),
Expand Down Expand Up @@ -815,7 +818,7 @@ const layerEnvironmentRegistration = toolkitRegistration(

const layerProjectRegistration = toolkitRegistration(ProjectToolkit, ProjectHandlers.layer);

const layerAttachmentRegistration = toolkitRegistration(
export const layerAttachmentToolkit = toolkitRegistration(
AttachmentToolkit,
AttachmentHandlers.layer,
);
Expand Down Expand Up @@ -851,7 +854,7 @@ export const layer = Layer.mergeAll(
layerPreviewToolkit,
layerOrchestratorToolkit,
layerThreadToolkit,
layerAttachmentRegistration,
layerAttachmentToolkit,
layerProjectRegistration,
layerEnvironmentRegistration,
layerPreviewControlsRegistration,
Expand Down
150 changes: 99 additions & 51 deletions apps/server/src/mcp/toolkits/attachment/handlers.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,12 @@
import { type ChatAttachment, MessageId, OrchestratorMcpFailure } from "@t3tools/contracts";
import {
ATTACHMENT_UPLOAD_URL_TTL_MS,
type ChatAttachment,
MessageId,
OrchestratorMcpFailure,
} from "@t3tools/contracts";
import * as Clock from "effect/Clock";
import * as Effect from "effect/Effect";
import * as McpInvocationContext from "../../McpInvocationContext.ts";
import * as Upload from "../../../assets/AttachmentUpload.ts";
import * as Claims from "../../../orchestration-v2/AttachmentClaims.ts";
import * as ThreadMessageIntake from "../../../orchestration-v2/ThreadMessageIntake.ts";
Expand Down Expand Up @@ -27,53 +34,94 @@ export function resolveAttachmentReferences(
});
}

export const layer = McpToolAccess.toLayer(AttachmentToolkit, {
t3_attachment_prepare_upload: McpToolAccess.writes((input) =>
Upload.issueAttachmentUploadUrl(input.upload).pipe(Effect.mapError(unavailable)),
),
t3_attachment_discard: McpToolAccess.writes((input) =>
Upload.deletePendingAttachment(input.attachmentId).pipe(Effect.as({})),
),
t3_thread_send_attachments: McpToolAccess.writesThreads(
(input) => [input.threadId],
(input) =>
Effect.gen(function* () {
const { caller, projection } = yield* readThread(input.threadId, ["messages"]);
if (projection.thread.archivedAt !== null)
return yield* new OrchestratorMcpFailure({
code: "invalid_request",
message: "Unarchive the target thread before sending attachments.",
});
const attachments = yield* resolveAttachmentReferences(
input.attachments,
projection.messages.flatMap((message) => message.attachments),
);
const commandId = yield* newCommandId();
const messageId = MessageId.make(commandId);
const result = yield* ThreadMessageIntake.sendToThread({
projectId: projection.thread.projectId,
threadId: projection.thread.id,
commandId,
messageId,
...(caller === undefined ? {} : { senderThreadId: caller.id }),
text: input.message ?? "",
attachments,
mode: "auto",
createdBy: "agent",
creationSource: "mcp",
}).pipe(
Effect.mapError((error) =>
error._tag === "AttachmentClaimError"
? new OrchestratorMcpFailure({ code: "orchestration_error", message: error.message })
: unavailable(),
),
);
return {
threadId: projection.thread.id,
messageId,
runId: result.run.id,
status: result.run.status,
};
}),
),
});
/**
* As long as the pending file can live. The sweep counts its 24 hours from when
* the upload finishes, which can be up to the upload URL's lifetime after issue.
*/
const UPLOAD_OWNER_TTL_MS = 24 * 60 * 60 * 1000 + ATTACHMENT_UPLOAD_URL_TTL_MS;

/** A thread owns its uploads across provider sessions; an outside client per MCP session. */
const uploadOwner = McpInvocationContext.McpInvocationContext.pipe(
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Effect.map((scope) => scope.thread?.threadId ?? scope.requestNamespace),
);

export const layer = McpToolAccess.toLayer(
AttachmentToolkit,
Effect.sync(() => {
// Which caller prepared each pending upload, so only that caller can discard it.
const uploadOwners = new Map<string, { readonly owner: string; readonly issuedAt: number }>();
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
return {
t3_attachment_prepare_upload: McpToolAccess.writes((input) =>
Effect.gen(function* () {
const result = yield* Upload.issueAttachmentUploadUrl(input.upload).pipe(
Effect.mapError(unavailable),
);
const now = yield* Clock.currentTimeMillis;
for (const [id, entry] of uploadOwners) {
if (now - entry.issuedAt > UPLOAD_OWNER_TTL_MS) uploadOwners.delete(id);
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
}
uploadOwners.set(result.attachmentId, { owner: yield* uploadOwner, issuedAt: now });
return result;
}),
),
t3_attachment_discard: McpToolAccess.writes((input) =>
Effect.gen(function* () {
if (uploadOwners.get(input.attachmentId)?.owner !== (yield* uploadOwner)) {
return yield* new OrchestratorMcpFailure({
code: "invalid_request",
message: "Only the caller that prepared a pending upload can discard it.",
});
}
yield* Upload.deletePendingAttachment(input.attachmentId);
uploadOwners.delete(input.attachmentId);
return {};
}),
),
t3_thread_send_attachments: McpToolAccess.writesThreads(
(input) => [input.threadId],
(input) =>
Effect.gen(function* () {
const { caller, projection } = yield* readThread(input.threadId, ["messages"]);
if (projection.thread.archivedAt !== null)
return yield* new OrchestratorMcpFailure({
code: "invalid_request",
message: "Unarchive the target thread before sending attachments.",
});
const attachments = yield* resolveAttachmentReferences(
input.attachments,
projection.messages.flatMap((message) => message.attachments),
);
const commandId = yield* newCommandId();
const messageId = MessageId.make(commandId);
const result = yield* ThreadMessageIntake.sendToThread({
projectId: projection.thread.projectId,
threadId: projection.thread.id,
commandId,
messageId,
...(caller === undefined ? {} : { senderThreadId: caller.id }),
text: input.message ?? "",
attachments,
mode: "auto",
createdBy: "agent",
creationSource: "mcp",
}).pipe(
Effect.mapError((error) =>
error._tag === "AttachmentClaimError"
? new OrchestratorMcpFailure({
code: "orchestration_error",
message: error.message,
})
: unavailable(),
),
);
return {
threadId: projection.thread.id,
messageId,
runId: result.run.id,
status: result.run.status,
};
}),
),
};
}),
);
Loading
Loading