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
14 changes: 14 additions & 0 deletions .github/workflows/deploy-kit.yml
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,17 @@ on:
branches: [main]
paths:
- "packages/kit/**"
# kit.openiap.dev/mcp is served by kit's Fly binary importing
# @hyodotdev/openiap-mcp-server/web straight from source, so an
# MCP-server change must redeploy kit or it never ships (issue #287).
- "packages/mcp-server/**"
- ".github/workflows/deploy-kit.yml"
- "bun.lock"
- "package.json"
pull_request:
paths:
- "packages/kit/**"
- "packages/mcp-server/**"
Comment on lines 19 to +22

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.

📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win

Mirror the dependency filters for pull requests.

The pull_request.paths list omits bun.lock, package.json, and .github/workflows/deploy-kit.yml. A pull request that changes only one of these files skips this workflow, including the MCP server suite. Add the same dependency and workflow paths used by push.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In @.github/workflows/deploy-kit.yml around lines 19 - 22, Update the
pull_request.paths configuration in deploy-kit.yml to mirror the push path
filters by adding bun.lock, package.json, and .github/workflows/deploy-kit.yml,
while preserving the existing packages/kit and packages/mcp-server entries.

- ".github/workflows/deploy-kit.yml"
- "bun.lock"
- "package.json"
Expand Down Expand Up @@ -56,6 +61,15 @@ jobs:
- name: Run tests (convex + server unit tests)
run: bun run test

- name: Lint + test MCP server (served by kit's /mcp route)
# kit's Fly binary imports @hyodotdev/openiap-mcp-server/web from
# source, so its regressions ship with kit deploys. This workflow
# is the only CI that runs the MCP server's own suite.
working-directory: packages/mcp-server
run: |
bun run lint
bun run test

- name: Vite build
env:
VITE_KIT_CONVEX_URL: https://placeholder-build-1.convex.cloud
Expand Down
68 changes: 65 additions & 3 deletions packages/mcp-server/src/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,11 @@ import {
IAPKIT_MCP_SERVER_NAME,
IAPKIT_MCP_SERVER_VERSION,
} from "./mcp.js";
import {
buildSessionId,
currentMachineId,
routeUnknownSession,
} from "./session-routing.js";

const DEFAULT_MCP_PATH = "/mcp";
const DEFAULT_PORT = 3939;
Expand Down Expand Up @@ -48,6 +53,13 @@ export interface RemoteMcpHttpServerOptions {
allowedOrigins?: string[];
/** Logger for lifecycle and request failures. Defaults to console. */
logger?: Pick<Console, "error" | "info">;
/**
* Identity of this process for session affinity. Defaults to
* FLY_MACHINE_ID; session ids are prefixed with it so a follow-up
* request landing on a sibling machine can be replayed to the owner
* (GitHub issue #287). Undefined disables replay routing.
*/
machineId?: string;
}

/** Runtime handle for an IAPKit remote MCP HTTP server. */
Expand All @@ -72,6 +84,7 @@ export function createRemoteMcpHttpServer(
const allowedOrigins =
options.allowedOrigins ??
parseAllowedOrigins(process.env.IAPKIT_MCP_ALLOWED_ORIGINS);
const machineId = options.machineId ?? currentMachineId();

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.

🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Validate configured machine IDs in both handlers.

Both handlers accept an explicit machineId without applying the validation used for FLY_MACHINE_ID. A value containing . creates an ID that resolves to a different replay target.

  • packages/mcp-server/src/http.ts#L87-L87: normalize options.machineId before storing it.
  • packages/mcp-server/src/web.ts#L48-L48: use the same normalizer before storing options.machineId.
📍 Affects 2 files
  • packages/mcp-server/src/http.ts#L87-L87 (this comment)
  • packages/mcp-server/src/web.ts#L48-L48
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@packages/mcp-server/src/http.ts` at line 87, Normalize the configured machine
ID before storing it in both handlers: update the machine ID assignment in
packages/mcp-server/src/http.ts at lines 87-87 and
packages/mcp-server/src/web.ts at lines 48-48 to apply the existing normalizer
to options.machineId while preserving currentMachineId() fallback behavior.

const transports = new Map<string, StreamableHTTPServerTransport>();

const server = createServer(async (req, res) => {
Expand Down Expand Up @@ -134,12 +147,13 @@ export function createRemoteMcpHttpServer(
res,
transports,
logger,
machineId,
);
return;
}

if (req.method === "GET" || req.method === "DELETE") {
await handleExistingMcpSession(req, res, transports);
await handleExistingMcpSession(req, res, transports, machineId);
return;
}

Expand Down Expand Up @@ -223,6 +237,7 @@ async function handleMcpPost(
res: ServerResponse,
transports: Map<string, StreamableHTTPServerTransport>,
logger: Pick<Console, "error" | "info">,
machineId: string | undefined,
): Promise<void> {
const sessionId = headerString(req.headers["mcp-session-id"]);
const body = await readJsonBody(req);
Expand All @@ -233,7 +248,12 @@ async function handleMcpPost(
return;
}

if (sessionId || !isInitializeRequest(body)) {
if (sessionId) {
writeUnknownSessionResponse(req, res, sessionId, machineId);
return;
}

if (!isInitializeRequest(body)) {
writeJsonRpcError(
res,
400,
Expand All @@ -245,7 +265,7 @@ async function handleMcpPost(

let transport!: StreamableHTTPServerTransport;
transport = new StreamableHTTPServerTransport({
sessionIdGenerator: () => randomUUID(),
sessionIdGenerator: () => buildSessionId(machineId, randomUUID()),
onsessioninitialized: (initializedSessionId) => {
transports.set(initializedSessionId, transport);
logger.info(`IAPKit MCP session initialized: ${initializedSessionId}`);
Expand All @@ -269,18 +289,60 @@ async function handleExistingMcpSession(
req: IncomingMessage,
res: ServerResponse,
transports: Map<string, StreamableHTTPServerTransport>,
machineId: string | undefined,
): Promise<void> {
const sessionId = headerString(req.headers["mcp-session-id"]);
const transport = sessionId ? transports.get(sessionId) : undefined;

if (!transport) {
if (sessionId) {
writeUnknownSessionResponse(req, res, sessionId, machineId);
return;
}
writeJsonRpcError(res, 400, -32000, "Invalid or missing mcp-session-id");
return;
}

await transport.handleRequest(req as AuthenticatedRequest, res);
}

/**
* Answers a request whose session id isn't in this process's transport
* map: replay it to the machine that minted the id when possible,
* otherwise 404 so a spec-compliant client transparently re-initializes.
* (The previous 400 "initialize first" reply broke that recovery path —
* GitHub issue #287.)
*/
function writeUnknownSessionResponse(
req: IncomingMessage,
res: ServerResponse,
sessionId: string,
machineId: string | undefined,
): void {
const routing = routeUnknownSession({
sessionId,
machineId,
alreadyReplayed: req.headers["fly-replay-src"] !== undefined,
});

if (routing.action === "replay") {
// Fly's proxy intercepts any response carrying `fly-replay` and
// re-sends the original request to the named machine; the client
// never sees this interim response.
res
.writeHead(204, { "fly-replay": `instance=${routing.targetMachineId}` })
.end();
return;
}

writeJsonRpcError(
res,
404,
-32001,
"Session not found — initialize a new MCP session.",
);
}

function attachAuthInfo(req: AuthenticatedRequest): void {
const bearerToken = parseBearerToken(headerString(req.headers.authorization));
if (!bearerToken) return;
Expand Down
73 changes: 73 additions & 0 deletions packages/mcp-server/src/session-routing.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
// MCP session ids are held in per-process memory (the transport object
// itself is stateful — an SSE stream can't be serialized into a shared
// store), so a session created on one Fly machine is invisible to its
// siblings. Fix (GitHub issue #287): embed the creating machine's id in
// the session id, and when a request lands on the wrong machine, answer
// with a `fly-replay` header so Fly's proxy re-routes the original
// request to the owner. Off Fly (no FLY_MACHINE_ID) session ids stay
// plain UUIDs and routing always resolves to `not-found`.

/**
* Fly machine ids are lowercase hex today, but only shape-check them:
* the prefix is attacker-controlled (it arrives inside the client's
* `mcp-session-id` header), so the pattern also guards the value we
* echo back inside the `fly-replay` response header.
*/
const MACHINE_ID_PATTERN = /^[A-Za-z0-9]{1,32}$/;

const SESSION_MACHINE_SEPARATOR = ".";

/** Reads the Fly machine identity, or undefined when not running on Fly. */
export function currentMachineId(
env: Record<string, string | undefined> = process.env,
): string | undefined {
const raw = env.FLY_MACHINE_ID;
return raw && MACHINE_ID_PATTERN.test(raw) ? raw : undefined;
}

/** Builds a session id that carries the creating machine's identity. */
export function buildSessionId(
machineId: string | undefined,
uuid: string,
): string {
return machineId ? `${machineId}${SESSION_MACHINE_SEPARATOR}${uuid}` : uuid;
}

/** Routing decision for a session id this process doesn't recognize. */
export type UnknownSessionRouting =
| { action: "replay"; targetMachineId: string }
| { action: "not-found" };

/**
* Decides what to do with a session id that isn't in the local
* transport map.
*
* @param options.sessionId Session id from the `mcp-session-id` header.
* @param options.machineId This process's machine id (undefined off Fly).
* @param options.alreadyReplayed True when the request carries
* `fly-replay-src`, i.e. it was already replayed once — never replay
* again or two stale machines could bounce a request forever.
* @returns `replay` toward the owning machine, or `not-found` (the
* caller answers 404 so the client re-initializes per the MCP spec).
*/
export function routeUnknownSession(options: {
sessionId: string;
machineId: string | undefined;
alreadyReplayed: boolean;
}): UnknownSessionRouting {
if (!options.machineId || options.alreadyReplayed) {
return { action: "not-found" };
}

const separatorIndex = options.sessionId.indexOf(SESSION_MACHINE_SEPARATOR);
if (separatorIndex <= 0) return { action: "not-found" };

const prefix = options.sessionId.slice(0, separatorIndex);
if (!MACHINE_ID_PATTERN.test(prefix) || prefix === options.machineId) {
// Malformed prefix, or the session was minted by this very machine
// (map lost to a restart/deploy) — replaying to ourselves would loop.
return { action: "not-found" };
}

return { action: "replay", targetMachineId: prefix };
}
63 changes: 61 additions & 2 deletions packages/mcp-server/src/web.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,11 @@ import {
isPublishableApiKey,
} from "./auth.js";
import { createIapKitMcpServer } from "./mcp.js";
import {
buildSessionId,
currentMachineId,
routeUnknownSession,
} from "./session-routing.js";

const MAX_MCP_BODY_BYTES = 1024 * 1024;
const MCP_BODY_TOO_LARGE_ERROR = "MCP request body is too large";
Expand All @@ -24,6 +29,13 @@ const DEFAULT_ALLOWED_ORIGINS = [
export interface IapKitWebMcpHandlerOptions {
allowedOrigins?: string[];
logger?: Pick<Console, "error" | "info">;
/**
* Identity of this process for session affinity. Defaults to
* FLY_MACHINE_ID; session ids are prefixed with it so a follow-up
* request landing on a sibling machine can be replayed to the owner
* (GitHub issue #287). Undefined disables replay routing.
*/
machineId?: string;
}

export function createIapKitWebMcpHandler(
Expand All @@ -33,6 +45,7 @@ export function createIapKitWebMcpHandler(
const allowedOrigins =
options.allowedOrigins ??
parseAllowedOrigins(process.env.IAPKIT_MCP_ALLOWED_ORIGINS);
const machineId = options.machineId ?? currentMachineId();
const transports = new Map<
string,
WebStandardStreamableHTTPServerTransport
Expand Down Expand Up @@ -69,6 +82,7 @@ export function createIapKitWebMcpHandler(
transports,
logger,
authInfo,
machineId,
);
return withCors(request, response, allowedOrigins);
}
Expand All @@ -78,6 +92,7 @@ export function createIapKitWebMcpHandler(
request,
transports,
authInfo,
machineId,
);
return withCors(request, response, allowedOrigins);
}
Expand Down Expand Up @@ -117,6 +132,7 @@ async function handlePost(
transports: Map<string, WebStandardStreamableHTTPServerTransport>,
logger: Pick<Console, "error" | "info">,
authInfo: AuthInfo | undefined,
machineId: string | undefined,
): Promise<Response> {
const sessionId = request.headers.get("mcp-session-id") ?? undefined;
const body = await readJsonBody(request);
Expand All @@ -129,7 +145,11 @@ async function handlePost(
});
}

if (sessionId || !isInitializeRequest(body)) {
if (sessionId) {
return unknownSessionResponse(request, sessionId, machineId);
}

if (!isInitializeRequest(body)) {
return jsonRpcError(
400,
-32000,
Expand All @@ -139,7 +159,7 @@ async function handlePost(

let transport!: WebStandardStreamableHTTPServerTransport;
transport = new WebStandardStreamableHTTPServerTransport({
sessionIdGenerator: () => randomUUID(),
sessionIdGenerator: () => buildSessionId(machineId, randomUUID()),
onsessioninitialized: (initializedSessionId) => {
transports.set(initializedSessionId, transport);
logger.info(`IAPKit MCP session initialized: ${initializedSessionId}`);
Expand Down Expand Up @@ -167,17 +187,56 @@ async function handleExistingSession(
request: Request,
transports: Map<string, WebStandardStreamableHTTPServerTransport>,
authInfo: AuthInfo | undefined,
machineId: string | undefined,
): Promise<Response> {
const sessionId = request.headers.get("mcp-session-id") ?? undefined;
const transport = sessionId ? transports.get(sessionId) : undefined;

if (!transport) {
if (sessionId) {
return unknownSessionResponse(request, sessionId, machineId);
}
return jsonRpcError(400, -32000, "Invalid or missing mcp-session-id");
}

return transport.handleRequest(request, { authInfo });
}

/**
* Answers a request whose session id isn't in this process's transport
* map: replay it to the machine that minted the id when possible,
* otherwise 404 so a spec-compliant client transparently re-initializes.
* (The previous 400 "initialize first" reply broke that recovery path —
* GitHub issue #287.)
*/
function unknownSessionResponse(
request: Request,
sessionId: string,
machineId: string | undefined,
): Response {
const routing = routeUnknownSession({
sessionId,
machineId,
alreadyReplayed: request.headers.has("fly-replay-src"),
});

if (routing.action === "replay") {
// Fly's proxy intercepts any response carrying `fly-replay` and
// re-sends the original request to the named machine; the client
// never sees this interim response.
return new Response(null, {
status: 204,
headers: { "fly-replay": `instance=${routing.targetMachineId}` },
});
}

return jsonRpcError(
404,
-32001,
"Session not found — initialize a new MCP session.",
);
}

function authInfoFromRequest(request: Request): AuthInfo | undefined {
const token = parseBearerToken(request.headers.get("authorization"));
if (!token) return undefined;
Expand Down
Loading