diff --git a/packages/mcp-server/src/http-server.ts b/packages/mcp-server/src/http-server.ts index f23a2455..21fedd8c 100644 --- a/packages/mcp-server/src/http-server.ts +++ b/packages/mcp-server/src/http-server.ts @@ -90,11 +90,9 @@ interface SessionState { xml: string version: number lastUpdated: Date + createdAt?: Date // When this session's state was first created (heartbeat grace window) svg?: string // Cached SVG from last browser save syncRequested?: number // Timestamp when sync requested, cleared when browser responds - exportFormat?: "png" | "svg" // Set by MCP tool to request browser export - exportXml?: string // Single-page projection to load before a page-targeted export - exportData?: string // Base64/SVG data returned by browser after export } export const stateStore = new Map() @@ -115,40 +113,168 @@ export function setState(sessionId: string, xml: string, svg?: string): number { xml, version: newVersion, lastUpdated: new Date(), + createdAt: existing?.createdAt || new Date(), svg: svg || existing?.svg, // Preserve cached SVG if not provided syncRequested: undefined, // Clear sync request when browser pushes state - exportFormat: existing?.exportFormat, // Preserve pending export request - exportXml: existing?.exportXml, // Preserve pending projection - exportData: existing?.exportData, // Preserve export result }) log.debug(`State updated: session=${sessionId}, version=${newVersion}`) return newVersion } /** - * Ask the browser bridge to export the current diagram as png/svg. + * Export job queue (per session). * - * When `projectionXml` is given (a single-page ), the bridge loads it - * first, waits for draw.io's own load event, exports, then reloads the - * session's real document — so a page-targeted export never mutates the - * canonical session state and needs no fixed-delay guessing on the server. - * - * Returns false when the session is unknown. Callers should then poll - * `getState(sessionId)?.exportData` for the result. + * Replaces the old single-slot exportFormat/exportXml/exportData protocol. + * Concurrent MCP export calls each enqueue a job with a unique id; the + * browser bridge processes jobs strictly one at a time (head of queue) and + * reports each result/failure BY JOB ID, so: + * - concurrent exports queue up instead of racing on one state slot, + * - a late/stale response can never satisfy the wrong job, + * - a failed/timed-out job is dropped and the queue advances. */ -export function requestExport( +export interface ExportJob { + id: number + format: "png" | "svg" + projectionXml: string | null // Single-page for page-targeted exports +} + +interface JobQueue { + seq: number + queue: ExportJob[] + waiters: Map void }> +} + +const exportQueues = new Map() // sessionId -> queue state +const browserHeartbeats = new Map() // sessionId -> last /api/state poll +const HEARTBEAT_STALE_MS = 20_000 + +function getJobQueue(sessionId: string): JobQueue { + let jq = exportQueues.get(sessionId) + if (!jq) { + jq = { seq: 0, queue: [], waiters: new Map() } + exportQueues.set(sessionId, jq) + } + return jq +} + +export function touchBrowserHeartbeat(sessionId: string): void { + browserHeartbeats.set(sessionId, new Date()) +} + +export function getBrowserHeartbeat(sessionId: string): Date | undefined { + return browserHeartbeats.get(sessionId) +} + +/** + * True when the browser bridge is (or very recently was) polling us. Used to + * fail export enqueues fast instead of hanging until the per-job timeout when + * the preview tab is closed. A short grace window after session creation + * covers the page still loading right after start_session. + */ +export function isBrowserConnected(sessionId: string): boolean { + const hb = browserHeartbeats.get(sessionId) + if (hb && Date.now() - hb.getTime() <= HEARTBEAT_STALE_MS) return true + const st = stateStore.get(sessionId) + return st?.createdAt ? Date.now() - st.createdAt.getTime() < 30_000 : false +} + +/** Enqueue an export job. Returns the job, or null when the session is unknown. */ +export function enqueueExport( sessionId: string, format: "png" | "svg", projectionXml?: string, -): boolean { - const state = stateStore.get(sessionId) - if (!state) return false - state.exportData = undefined - state.exportXml = projectionXml - state.exportFormat = format +): ExportJob | null { + if (!stateStore.has(sessionId)) return null + const jq = getJobQueue(sessionId) + const job: ExportJob = { id: ++jq.seq, format, projectionXml: projectionXml || null } + jq.queue.push(job) + log.debug( + `Export job #${job.id} (${format}) enqueued for session=${sessionId} (queue=${jq.queue.length})`, + ) + return job +} + +/** The job the browser bridge should work on next (head of queue). */ +export function currentExportJob(sessionId: string): ExportJob | undefined { + const jq = exportQueues.get(sessionId) + return jq ? jq.queue[0] : undefined +} + +function settleJob(sessionId: string, jobId: number, data: string | null): boolean { + const jq = exportQueues.get(sessionId) + if (!jq) return false + const idx = jq.queue.findIndex((j) => j.id === jobId) + if (idx === -1) return false // stale/unknown id (already timed out) — ignore + jq.queue.splice(idx, 1) + const w = jq.waiters.get(jobId) + if (w) { + jq.waiters.delete(jobId) + w.resolve(data) + } return true } +/** Browser delivered image data for a job. Returns false for stale ids. */ +export function completeExportJob( + sessionId: string, + jobId: number, + data: string, +): boolean { + const ok = settleJob(sessionId, jobId, data) + if (ok) log.debug(`Export job #${jobId} completed for session=${sessionId}`) + return ok +} + +/** Browser gave up on a job (render timeout). Resolves its waiter with null. */ +export function failExportJob(sessionId: string, jobId: number): boolean { + const ok = settleJob(sessionId, jobId, null) + if (ok) + log.warn( + `Export job #${jobId} failed (browser-side timeout) for session=${sessionId}`, + ) + return ok +} + +/** + * Await one export job. Resolves with the data URL on success, null on + * timeout or browser disconnect. The job is removed from the queue on any + * outcome, so successors are never blocked by an abandoned head. + */ +export function waitForExportJob( + sessionId: string, + jobId: number, + timeoutMs: number, +): Promise { + const jq = getJobQueue(sessionId) + return new Promise((resolve) => { + let done = false + const finish = (val: string | null) => { + if (done) return + done = true + clearTimeout(timer) + clearInterval(hbCheck) + jq.waiters.delete(jobId) + if (val === null) { + const idx = jq.queue.findIndex((j) => j.id === jobId) + if (idx !== -1) jq.queue.splice(idx, 1) + } + resolve(val) + } + const timer = setTimeout(() => { + log.warn( + `Export job #${jobId} timed out after ${timeoutMs}ms for session=${sessionId}`, + ) + finish(null) + }, timeoutMs) + // If the browser page dies mid-queue (tab closed), fail fast instead + // of letting every queued job hang until its own timeout. + const hbCheck = setInterval(() => { + if (!isBrowserConnected(sessionId)) finish(null) + }, 2_000) + jq.waiters.set(jobId, { resolve: finish }) + }) +} + export function requestSync(sessionId: string): boolean { const state = stateStore.get(sessionId) if (state) { @@ -225,6 +351,15 @@ function cleanupExpiredSessions(): void { if (now - state.lastUpdated.getTime() > SESSION_TTL) { stateStore.delete(sessionId) clearHistory(sessionId) + // Fail any still-waiting export jobs and drop queue state. + const jq = exportQueues.get(sessionId) + if (jq) { + for (const job of [...jq.queue]) { + failExportJob(sessionId, job.id) + } + exportQueues.delete(sessionId) + } + browserHeartbeats.delete(sessionId) log.info(`Cleaned up expired session: ${sessionId}`) } } @@ -304,15 +439,24 @@ function handleStateApi( return } ensureSessionStateInitialized(sessionId) + touchBrowserHeartbeat(sessionId) const state = stateStore.get(sessionId) + // Expose the head of the export job queue. The browser works jobs + // strictly one at a time and reports results by job id (see POST). + const exportJobEntry = currentExportJob(sessionId) res.writeHead(200, { "Content-Type": "application/json" }) res.end( JSON.stringify({ xml: state?.xml || null, version: state?.version || 0, syncRequested: !!state?.syncRequested, - exportFormat: state?.exportFormat || null, - exportXml: state?.exportXml || null, + exportJob: exportJobEntry + ? { + id: exportJobEntry.id, + format: exportJobEntry.format, + xml: exportJobEntry.projectionXml || null, + } + : null, }), ) } else if (req.method === "POST") { @@ -326,16 +470,15 @@ function handleStateApi( return } - // Browser is returning export data (png/svg) - if (data.exportData !== undefined) { - const state = stateStore.get(sessionId) - if (state) { - state.exportData = data.exportData - state.exportFormat = undefined - state.exportXml = undefined - log.debug( - `Export data received for session=${sessionId}`, - ) + // Browser is returning a queued export job's result (or + // failure) BY JOB ID. Stale ids (job already timed out + // server-side) are ignored — this kills the old + // late-response-poisons-next-export race. + if (data.exportJobId !== undefined) { + if (data.exportFailed === true) { + failExportJob(sessionId, data.exportJobId) + } else if (data.exportData !== undefined) { + completeExportJob(sessionId, data.exportJobId, data.exportData) } res.writeHead(200, { "Content-Type": "application/json" }) res.end(JSON.stringify({ success: true })) @@ -702,7 +845,8 @@ function getHtmlPage(sessionId: string): string { let currentVersion = 0, isReady = false, pendingXml = null, lastXml = null; let pendingSvgExport = null; let pendingAiSvg = false; - let pendingMcpExport = null; // 'png' or 'svg' when MCP requested export + let pendingMcpExport = null; // { id, format } while a queued MCP export job is in flight + let mcpExportTimer = null; // fallback timer for the CURRENT job only (cleared on job switch) let projectionExportActive = false; // page-targeted export: showing a transient single-page projection let projectionRestoreXml = null; // the real document to reload once a projection export finishes @@ -724,19 +868,24 @@ function getHtmlPage(sessionId: string): string { // Fallback if export doesn't respond setTimeout(() => { if (pendingSvgExport === msg.xml) { pushState(msg.xml, ''); pendingSvgExport = null; } }, 2000); } else if (msg.event === 'export' && msg.data) { - // Handle MCP server export request (png/svg) + // Handle queued MCP export job (png/svg) // Verify the response matches the requested format to avoid capturing // unrelated exports (autosave SVG, sync XML) if (pendingMcpExport) { const d = msg.data; - const isPng = pendingMcpExport === 'png' && (d.startsWith('data:image/png') || (typeof d === 'string' && d.length > 100 && !d.startsWith('<'))); - const isSvg = pendingMcpExport === 'svg' && (d.startsWith('data:image/svg') || d.startsWith(' 100 && !d.startsWith('<'))); + const isSvg = pendingMcpExport.format === 'svg' && (d.startsWith('data:image/svg') || d.startsWith(' {}); // Page-targeted export: restore the user's real // multi-page document now that we have the image. @@ -848,45 +997,62 @@ function getHtmlPage(sessionId: string): string { currentVersion = s.version; loadDiagram(s.xml, true); } - // Handle export request from MCP server (png/svg). + // Handle ONE queued export job from the MCP server at a time + // (the head of the server-side queue, delivered as + // s.exportJob = { id, format, xml }). Concurrent MCP export + // calls queue server-side; we finish one job (result or + // failure, reported by id) and pick up the next on a + // following poll. // // Plain export: capture whatever tab is currently displayed. // - // Page-targeted export: the server sends a single-page - // projection in s.exportXml. We load it into the iframe, let + // Page-targeted export: the job carries a single-page + // projection in job.xml. We load it into the iframe, let // draw.io render it, export, then reload the user's real // document — all browser-side. The canonical session state is // never mutated, so there is no server-side restore race and no // dependence on poll timing. autosave is suppressed while the // projection is showing (see projectionExportActive guard). - if (s.exportFormat && !pendingMcpExport && isReady) { - pendingMcpExport = s.exportFormat; + if (s.exportJob && !pendingMcpExport && isReady) { + const job = s.exportJob; + pendingMcpExport = { id: job.id, format: job.format }; const fireExport = () => { - const exportOpts = pendingMcpExport === 'png' + const exportOpts = pendingMcpExport && pendingMcpExport.format === 'png' ? { action: 'export', format: 'png', scale: 2 } : { action: 'export', format: 'svg' }; iframe.contentWindow.postMessage(JSON.stringify(exportOpts), '*'); }; - if (s.exportXml) { + if (job.xml) { // Stash the real document so we can restore after export. projectionRestoreXml = lastXml; projectionExportActive = true; // Load the projection without touching lastXml/server state. - iframe.contentWindow.postMessage(JSON.stringify({ action: 'load', xml: s.exportXml, autosave: 0 }), '*'); + iframe.contentWindow.postMessage(JSON.stringify({ action: 'load', xml: job.xml, autosave: 0 }), '*'); // Let draw.io render the loaded page before exporting // (same proven settle delay as the AI-preview path). setTimeout(fireExport, 600); } else { fireExport(); } - // Timeout: reset if draw.io never responds, and restore the - // real document if a projection was left showing. - setTimeout(() => { - if (pendingMcpExport) { + // Fallback: if draw.io never responds for THIS job, report + // failure BY ID so the server fails that one job and the + // queue advances. The timer is tracked and cleared when a + // job completes, and only ever kills its own job — the old + // leaked 10s timer could silently swallow a successor. + if (mcpExportTimer) clearTimeout(mcpExportTimer); + mcpExportTimer = setTimeout(() => { + mcpExportTimer = null; + if (pendingMcpExport && pendingMcpExport.id === job.id) { + const dead = pendingMcpExport; pendingMcpExport = null; restoreFromProjection(); + fetch('/api/state', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ sessionId, exportJobId: dead.id, exportFailed: true }) + }).catch(() => {}); } - }, 10000); + }, 30000); } } catch {} } diff --git a/packages/mcp-server/src/index.ts b/packages/mcp-server/src/index.ts index 5a180951..f845f0bc 100644 --- a/packages/mcp-server/src/index.ts +++ b/packages/mcp-server/src/index.ts @@ -48,12 +48,14 @@ import { import { checkEditGate } from "./edit-gate.js" import { addHistory } from "./history.js" import { + enqueueExport, getState, - requestExport, + isBrowserConnected, requestSync, setState, shutdown, startHttpServer, + waitForExportJob, waitForSync, } from "./http-server.js" import { parseDrawioFileContent } from "./load-diagram.js" @@ -919,7 +921,8 @@ server.registerTool( "- .drawio with NO page selector: writes the full (all pages).\n" + "- .drawio with a page selector: writes a single-page containing only that page.\n" + "- .png / .svg with NO page selector: exports the currently active page in the browser.\n" + - "- .png / .svg with a page selector: temporarily loads a single-page projection of that page into the browser, captures the rendered image, then restores the full document. The user will see a brief tab-flicker (~1-2s) but the exported image is guaranteed to be the requested page.", + "- .png / .svg with a page selector: temporarily loads a single-page projection of that page into the browser, captures the rendered image, then restores the full document. The user will see a brief tab-flicker (~1-2s) but the exported image is guaranteed to be the requested page.\n" + + "Concurrent export calls are safe: they are queued and rendered one at a time in order; each call returns as soon as its own image is ready.", inputSchema: { ...pageSelectorSchema, path: z @@ -1079,33 +1082,55 @@ server.registerTool( projectionXml = projection.xml } - // Ask the browser to export (optionally via a page projection) and - // poll for the resulting image data. - requestExport( + // Queue-based export. Concurrent export calls each enqueue a job + // with a unique id; the browser bridge renders jobs strictly one + // at a time (head of queue) and reports each result BY JOB ID, + // which resolves exactly this call's promise. This replaces the + // old single-slot requestExport + poll loop that raced under + // concurrency (requests overwrote each other; the single response + // was consumed by whichever caller polled first; late responses + // poisoned the next export). + if (!isBrowserConnected(currentSession.id)) { + return { + content: [ + { + type: "text", + text: "Error: The browser preview page is not connected (no recent poll). Open or refresh the session URL (e.g. call start_session), then retry the export.", + }, + ], + isError: true, + } + } + const job = enqueueExport( currentSession.id, detectedFormat as "png" | "svg", projectionXml, ) + if (!job) { + return { + content: [ + { + type: "text", + text: "Error: Session state not found. Is the browser open?", + }, + ], + isError: true, + } + } - // A projection export does an extra load + render round-trip in the - // browser, so give it a longer window. Re-read the live store entry - // each tick: setState() (from a concurrent autosave or tool call) - // replaces the Map entry with a new object, so a captured reference - // would go stale and never observe the browser's exportData. - const timeoutMs = projectionXml ? 15000 : 10000 - const start = Date.now() - let exportData: string | undefined - while (Date.now() - start < timeoutMs) { - exportData = getState(currentSession.id)?.exportData - if (exportData) break - await new Promise((r) => setTimeout(r, 200)) - } - const live = getState(currentSession.id) - if (live) { - live.exportData = undefined - live.exportFormat = undefined - live.exportXml = undefined - } + // Generous per-job budget: a queued job also waits for its + // predecessors to render, so this must cover the whole + // queue-drain, not just one render. The browser bridge reports + // per-job failures much sooner (30s render timeout) and a dead + // browser page is detected via heartbeat, so this wall-clock cap + // is only a backstop. + const timeoutMs = + Number(process.env.DRAWIO_EXPORT_TIMEOUT_MS) || 120_000 + const exportData = await waitForExportJob( + currentSession.id, + job.id, + timeoutMs, + ) if (!exportData) { return { @@ -1113,8 +1138,8 @@ server.registerTool( { type: "text", text: projectionXml - ? "Error: Export timed out after loading the single-page projection. The browser may be closed or unresponsive." - : "Error: Export timed out. Make sure the browser tab is open and the diagram is loaded.", + ? `Error: Export job #${job.id} timed out or failed after loading the single-page projection. The browser may be closed or unresponsive.` + : `Error: Export job #${job.id} timed out or failed. Make sure the browser tab is open and the diagram is loaded.`, }, ], isError: true, diff --git a/packages/mcp-server/tests/export-queue.test.ts b/packages/mcp-server/tests/export-queue.test.ts new file mode 100644 index 00000000..fea6aa05 --- /dev/null +++ b/packages/mcp-server/tests/export-queue.test.ts @@ -0,0 +1,170 @@ +/** + * Export job queue test: drives the patched export protocol the way the real + * pieces interact — MCP tool calls (enqueue + waitForExportJob) on one side, + * the browser bridge (GET /api/state poll -> POST result by job id) on the + * other — against the real HTTP server. + * + * Covers the races the old single-slot protocol had: + * - concurrent exports overwriting each other's request, + * - the single response consumed by the wrong caller, + * - late/stale responses poisoning the next export, + * - a failed/timed-out head blocking the queue. + */ + +import { afterAll, beforeAll, describe, expect, it } from "vitest" +import { + completeExportJob, + currentExportJob, + enqueueExport, + failExportJob, + getState, + isBrowserConnected, + setState, + shutdown, + startHttpServer, + touchBrowserHeartbeat, + waitForExportJob, +} from "../src/http-server.js" + +const SESSION = "mcp-export-queue-test" + +let port = 0 +let base = "" + +type StateResponse = { + xml: string | null + version: number + syncRequested: boolean + exportJob: { id: number; format: string; xml: string | null } | null +} + +async function browserPoll(): Promise { + const r = await fetch(`${base}/api/state?sessionId=${SESSION}`) + expect(r.status).toBe(200) + return r.json() +} + +async function browserPost(body: Record) { + const r = await fetch(`${base}/api/state`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ sessionId: SESSION, ...body }), + }) + return r.json() +} + +beforeAll(async () => { + port = await startHttpServer(6070) + base = `http://127.0.0.1:${port}` + setState(SESSION, "") + touchBrowserHeartbeat(SESSION) +}) + +afterAll(() => { + shutdown() +}) + +describe("export job queue", () => { + it("serializes concurrent enqueues and routes results by job id", async () => { + const j1 = enqueueExport(SESSION, "png") + const j2 = enqueueExport(SESSION, "png", "pageB") + const j3 = enqueueExport(SESSION, "svg") + expect(j1 && j2 && j3).toBeTruthy() + expect([j1!.id, j2!.id, j3!.id]).toEqual([1, 2, 3]) + + const p1 = waitForExportJob(SESSION, j1!.id, 5000) + const p2 = waitForExportJob(SESSION, j2!.id, 5000) + const p3 = waitForExportJob(SESSION, j3!.id, 5000) + + // Browser sees only the head of the queue. + const s1 = await browserPoll() + expect(s1.exportJob).toBeTruthy() + expect(s1.exportJob!.id).toBe(1) + expect(s1.exportJob!.format).toBe("png") + expect(s1.exportJob!.xml).toBeNull() + + // Job 1 renders first; job 2's waiter must not resolve yet. + let resolved2 = false + void p2.then(() => { + resolved2 = true + }) + await browserPost({ exportJobId: 1, exportData: "data:image/png;base64,AAAAjob1" }) + expect(await p1).toBe("data:image/png;base64,AAAAjob1") + await new Promise((r) => setTimeout(r, 150)) + expect(resolved2).toBe(false) + + // Queue advanced; job 2 carries its own projection. + const s2 = await browserPoll() + expect(s2.exportJob!.id).toBe(2) + expect(s2.exportJob!.xml).toBe("pageB") + await browserPost({ exportJobId: 2, exportData: "data:image/png;base64,BBBBjob2" }) + expect(await p2).toBe("data:image/png;base64,BBBBjob2") + + const s3 = await browserPoll() + expect(s3.exportJob!.id).toBe(3) + expect(s3.exportJob!.format).toBe("svg") + await browserPost({ exportJobId: 3, exportData: "data:image/svg+xml;base64,CCCCjob3" }) + expect(await p3).toBe("data:image/svg+xml;base64,CCCCjob3") + + const s4 = await browserPoll() + expect(s4.exportJob).toBeNull() + }) + + it("ignores stale/late results for already-completed job ids", async () => { + const r = await browserPost({ exportJobId: 1, exportData: "data:image/png;base64,LATE" }) + expect(r.success).toBe(true) + expect(currentExportJob(SESSION)).toBeUndefined() + }) + + it("fails a browser-reported job fast and advances the queue", async () => { + const j4 = enqueueExport(SESSION, "png") + const j5 = enqueueExport(SESSION, "png") + const p4 = waitForExportJob(SESSION, j4!.id, 5000) + const p5 = waitForExportJob(SESSION, j5!.id, 5000) + + const s = await browserPoll() + expect(s.exportJob!.id).toBe(j4!.id) + await browserPost({ exportJobId: j4!.id, exportFailed: true }) + + expect(await p4).toBeNull() + const s2 = await browserPoll() + expect(s2.exportJob!.id).toBe(j5!.id) + await browserPost({ exportJobId: j5!.id, exportData: "data:image/png;base64,DDDDjob5" }) + expect(await p5).toBe("data:image/png;base64,DDDDjob5") + }) + + it("removes a timed-out head so successors are not blocked", async () => { + const j6 = enqueueExport(SESSION, "png") + const j7 = enqueueExport(SESSION, "png") + const t0 = Date.now() + expect(await waitForExportJob(SESSION, j6!.id, 300)).toBeNull() + expect(Date.now() - t0).toBeLessThan(1500) + expect(currentExportJob(SESSION)?.id).toBe(j7!.id) + + const p7 = waitForExportJob(SESSION, j7!.id, 5000) + expect(completeExportJob(SESSION, j7!.id, "data:image/png;base64,EEEEjob7")).toBe(true) + expect(await p7).toBe("data:image/png;base64,EEEEjob7") + }) + + it("returns null from enqueue for unknown sessions", () => { + expect(enqueueExport("mcp-nope", "png")).toBeNull() + expect(currentExportJob("mcp-nope")).toBeUndefined() + expect(completeExportJob("mcp-nope", 1, "x")).toBe(false) + expect(failExportJob("mcp-nope", 1)).toBe(false) + }) + + it("keeps non-export state paths working (autosave push)", async () => { + const push = await browserPost({ xml: "", svg: "" }) + expect(push.success).toBe(true) + expect(push.version).toBeGreaterThan(0) + expect(getState(SESSION)!.xml).toContain("pushed") + const s = await browserPoll() + expect(s.version).toBe(push.version) + expect(s.exportJob).toBeNull() + }) + + it("treats the session as connected while the heartbeat is fresh", () => { + touchBrowserHeartbeat(SESSION) + expect(isBrowserConnected(SESSION)).toBe(true) + }) +})