From 3042834b167ba33ae942983cdf76f3613c521bd8 Mon Sep 17 00:00:00 2001 From: MoatazNoaman Date: Tue, 6 Oct 2026 10:33:43 +0300 Subject: [PATCH 1/3] fix(sync): a malformed websocket frame can no longer crash the server parseClientMessage was a bare JSON.parse, called inside the async handleMessage with no try/catch. Every Node/Bun transport invokes handleMessage fire-and-forget (`void ...`), so one invalid frame from any unauthenticated peer (e.g. "x", or {"type":"ModifyQuerySet"}) became an unhandled rejection and exited the whole process. - parseClientMessage now shape-checks every client message and throws ProtocolError for anything malformed. - handleMessage answers a ProtocolError with FatalError and closes only that session. - The concile dev/serve (node ws + Bun) and Vite embed transports catch any rejection from handleMessage and close the offending connection. --- .changeset/sync-malformed-frames.md | 11 +++ packages/cli/src/server.ts | 16 ++- packages/sync/src/handler.ts | 13 ++- packages/sync/src/index.ts | 1 + packages/sync/src/protocol.ts | 91 ++++++++++++++++- packages/sync/test/malformed-frame.test.ts | 110 +++++++++++++++++++++ packages/vite/src/embed.ts | 8 +- 7 files changed, 245 insertions(+), 5 deletions(-) create mode 100644 .changeset/sync-malformed-frames.md create mode 100644 packages/sync/test/malformed-frame.test.ts diff --git a/.changeset/sync-malformed-frames.md b/.changeset/sync-malformed-frames.md new file mode 100644 index 00000000..1c1603f1 --- /dev/null +++ b/.changeset/sync-malformed-frames.md @@ -0,0 +1,11 @@ +--- +"@concile/sync": patch +"@concile/cli": patch +"@concile/vite": patch +--- + +A malformed WebSocket frame no longer crashes the server. `parseClientMessage` now validates the +shape of every inbound client message and throws a `ProtocolError` for anything else; the handler +answers with a `FatalError` and closes only that session. The Node and Bun transports in +`concile dev`/`serve` and the Vite embed also catch any rejection from `handleMessage` and close +the offending connection, instead of leaving an unhandled rejection that exits the process. diff --git a/packages/cli/src/server.ts b/packages/cli/src/server.ts index cc0dfdf1..93b71cc9 100644 --- a/packages/cli/src/server.ts +++ b/packages/cli/src/server.ts @@ -291,7 +291,14 @@ async function startNodeServer(runtime: EmbeddedRuntime, options: DevServerOptio }, }; runtime.handler.connect(sessionId, syncSocket); - ws.on("message", (data: Buffer) => void runtime.handler.handleMessage(sessionId, data.toString("utf8"))); + ws.on("message", (data: Buffer) => { + // Fire-and-forget: a rejection here would otherwise be unhandled and exit the process, so a + // failing frame costs only its own connection. + runtime.handler.handleMessage(sessionId, data.toString("utf8")).catch((e: unknown) => { + console.error("[concile] sync message failed; closing connection:", e); + ws.close(1011); + }); + }); ws.on("close", () => runtime.handler.disconnect(sessionId)); ws.on("error", () => runtime.handler.disconnect(sessionId)); }); @@ -432,7 +439,12 @@ async function startBunServer(runtime: EmbeddedRuntime, options: DevServerOption runtime.handler.connect(ws.data.sessionId, syncSocket); }, message(ws, message) { - void runtime.handler.handleMessage(ws.data.sessionId, typeof message === "string" ? message : new TextDecoder().decode(message)); + const raw = typeof message === "string" ? message : new TextDecoder().decode(message); + // See the node `ws` path above: never leave this rejection unhandled. + runtime.handler.handleMessage(ws.data.sessionId, raw).catch((e: unknown) => { + console.error("[concile] sync message failed; closing connection:", e); + ws.close(); + }); }, close(ws) { bunPongCallbacks.delete(ws.data.sessionId); diff --git a/packages/sync/src/handler.ts b/packages/sync/src/handler.ts index 104a5d3d..59efe95c 100644 --- a/packages/sync/src/handler.ts +++ b/packages/sync/src/handler.ts @@ -15,6 +15,7 @@ import type { DiffableRange, DiffablePage } from "@concile/executor"; import { encodeServerMessage, parseClientMessage, + ProtocolError, INITIAL_VERSION, type ClientMessage, type ServerMessage, @@ -483,7 +484,17 @@ export class SyncProtocolHandler { const session = this.sessions.get(sessionId); if (!session) throw new Error(`unknown session: ${sessionId}`); session.hb.noteActivity(); // any inbound frame is liveness credit - const msg: ClientMessage = parseClientMessage(raw); + let msg: ClientMessage; + try { + msg = parseClientMessage(raw); + } catch (e) { + if (!(e instanceof ProtocolError)) throw e; + // A malformed frame is that peer's protocol violation: tell it why and drop only its session. + // Rethrowing would surface as an unhandled rejection in the fire-and-forget transports. + this.send(session, { type: "FatalError", message: e.message }); + this.reap(sessionId); + return; + } switch (msg.type) { case "Connect": return this.handleConnect(session, msg); diff --git a/packages/sync/src/index.ts b/packages/sync/src/index.ts index d58969a0..9c05b5fd 100644 --- a/packages/sync/src/index.ts +++ b/packages/sync/src/index.ts @@ -20,6 +20,7 @@ export { compareStateVersion, isContiguous, parseClientMessage, + ProtocolError, encodeServerMessage, } from "./protocol"; diff --git a/packages/sync/src/protocol.ts b/packages/sync/src/protocol.ts index cfdd6b12..cd9fb934 100644 --- a/packages/sync/src/protocol.ts +++ b/packages/sync/src/protocol.ts @@ -204,8 +204,97 @@ export type ServerMessage = | { type: "FatalError"; message: string } | { type: "Ping" }; +/** + * A client frame that isn't valid JSON or doesn't match a {@link ClientMessage} shape. The raw frame + * comes straight off an unauthenticated socket, so the handler treats this as a protocol violation + * by that one peer (FatalError + close), never as an engine failure. + */ +export class ProtocolError extends Error { + constructor(message: string) { + super(message); + this.name = "ProtocolError"; + } +} + +type Fields = Record; + +const isRecord = (v: unknown): v is Fields => typeof v === "object" && v !== null && !Array.isArray(v); +const isFiniteNumber = (v: unknown): v is number => typeof v === "number" && Number.isFinite(v); + +function assertShape(ok: boolean, what: string): void { + if (!ok) throw new ProtocolError(`malformed client message: ${what}`); +} + +function assertOptional(m: Fields, key: string, check: (v: unknown) => boolean, what: string): void { + if (m[key] !== undefined) assertShape(check(m[key]), what); +} + +const isMutationRef = (v: unknown): boolean => + isRecord(v) && typeof v.clientId === "string" && isFiniteNumber(v.seq); + +const isMutationRefs = (v: unknown): boolean => Array.isArray(v) && v.every(isMutationRef); + +const isQueryRequest = (v: unknown): boolean => + isRecord(v) && + isFiniteNumber(v.queryId) && + typeof v.udfPath === "string" && + (v.resultHash === undefined || typeof v.resultHash === "string") && + (v.sinceTs === undefined || isFiniteNumber(v.sinceTs)); + +const isMutationEntry = (v: unknown): boolean => + isRecord(v) && + typeof v.requestId === "string" && + typeof v.udfPath === "string" && + (v.clientId === undefined || typeof v.clientId === "string") && + (v.seq === undefined || isFiniteNumber(v.seq)); + +/** + * Parse and shape-check one inbound client frame. Only the fields a handler dereferences + * unconditionally are checked; `args`/`event` payloads stay opaque `JSONValue`s, validated later by + * the function's own validators. Throws {@link ProtocolError} on anything else. + */ export function parseClientMessage(raw: string): ClientMessage { - return JSON.parse(raw) as ClientMessage; + let m: unknown; + try { + m = JSON.parse(raw); + } catch { + throw new ProtocolError("malformed client message: invalid JSON"); + } + assertShape(isRecord(m), "expected an object"); + const msg = m as Fields; + switch (msg.type) { + case "Connect": + assertOptional(msg, "sessionId", (v) => typeof v === "string", "Connect.sessionId"); + assertOptional(msg, "clientId", (v) => typeof v === "string", "Connect.clientId"); + assertOptional(msg, "held", isMutationRefs, "Connect.held"); + assertOptional(msg, "ackedThrough", isMutationRefs, "Connect.ackedThrough"); + break; + case "ModifyQuerySet": + assertShape(Array.isArray(msg.add) && msg.add.every(isQueryRequest), "ModifyQuerySet.add"); + assertShape(Array.isArray(msg.remove) && msg.remove.every(isFiniteNumber), "ModifyQuerySet.remove"); + break; + case "Mutation": + assertShape(isMutationEntry(msg), "Mutation"); + break; + case "MutationBatch": + assertShape(Array.isArray(msg.entries) && msg.entries.every(isMutationEntry), "MutationBatch.entries"); + break; + case "Action": + assertShape(typeof msg.requestId === "string" && typeof msg.udfPath === "string", "Action"); + break; + case "EphemeralPublish": + assertShape(typeof msg.topic === "string", "EphemeralPublish"); + break; + case "SetAuth": + assertShape(typeof msg.token === "string" || msg.token === null, "SetAuth.token"); + break; + case "SetAdminAuth": + assertShape(typeof msg.key === "string", "SetAdminAuth.key"); + break; + default: + throw new ProtocolError("malformed client message: unknown type"); + } + return msg as unknown as ClientMessage; } export function encodeServerMessage(msg: ServerMessage): string { diff --git a/packages/sync/test/malformed-frame.test.ts b/packages/sync/test/malformed-frame.test.ts new file mode 100644 index 00000000..a8662c63 --- /dev/null +++ b/packages/sync/test/malformed-frame.test.ts @@ -0,0 +1,110 @@ +/** + * Inbound frames come straight off an unauthenticated socket. A malformed one must cost only the + * sending session (FatalError + close), never reject out of `handleMessage`: every Node/Bun + * transport calls it fire-and-forget, so a rejection there is an unhandled rejection, which exits + * the whole server process. + */ +import { describe, it, expect, vi } from "vitest"; +import { + SyncProtocolHandler, + parseClientMessage, + ProtocolError, + type SyncUdfExecutor, + type ServerMessage, +} from "../src/index"; + +const exec: SyncUdfExecutor = { + async runQuery(path) { + return { value: `user:${path}` as never, tables: ["t"], readRanges: [], globalTables: [] }; + }, + async runMutation() { + return { value: "ok" as never, tables: ["t"], writeRanges: [], commitTs: 1 }; + }, + async runAdminQuery(path) { + return { value: `admin:${path}` as never, tables: ["t"], readRanges: [], globalTables: [] }; + }, + async runAction(path) { + return { value: `acted:${path}` as never }; + }, +}; + +function sock() { + const sent: ServerMessage[] = []; + return { sent, send: (d: string) => sent.push(JSON.parse(d) as ServerMessage), bufferedAmount: 0, close: vi.fn() }; +} + +const MALFORMED = [ + "x", + "", + "null", + "42", + "[]", + JSON.stringify({}), + JSON.stringify({ type: "Nope" }), + JSON.stringify({ type: "ModifyQuerySet" }), + JSON.stringify({ type: "ModifyQuerySet", add: [null], remove: [] }), + JSON.stringify({ type: "ModifyQuerySet", add: [{ queryId: "1", udfPath: "a:b", args: {} }], remove: [] }), + JSON.stringify({ type: "ModifyQuerySet", add: [], remove: "1" }), + JSON.stringify({ type: "Mutation", udfPath: "a:b", args: {} }), + JSON.stringify({ type: "MutationBatch", entries: {} }), + JSON.stringify({ type: "MutationBatch", entries: [{ requestId: 1, udfPath: "a:b", args: {} }] }), + JSON.stringify({ type: "Action", requestId: "r" }), + JSON.stringify({ type: "Connect", sessionId: 7 }), + JSON.stringify({ type: "Connect", sessionId: "s", held: [{ clientId: "c" }] }), + JSON.stringify({ type: "SetAuth" }), + JSON.stringify({ type: "SetAdminAuth", key: 1 }), + JSON.stringify({ type: "EphemeralPublish", event: {} }), +]; + +describe("malformed client frames", () => { + it.each(MALFORMED)("%j: FatalError + close for that session, and handleMessage resolves", async (raw) => { + const h = new SyncProtocolHandler(exec, { autoNotifyOnMutation: false }); + const bad = sock(); + const good = sock(); + h.connect("bad", bad as never); + h.connect("good", good as never); + + await expect(h.handleMessage("bad", raw)).resolves.toBeUndefined(); + + expect(bad.sent).toEqual([{ type: "FatalError", message: expect.stringMatching(/^malformed client message/) }]); + expect(bad.close).toHaveBeenCalledTimes(1); + await expect(h.handleMessage("bad", JSON.stringify({ type: "SetAuth", token: null }))).rejects.toThrow(/unknown session/); + + // Other sessions are untouched. + await h.handleMessage("good", JSON.stringify({ type: "Mutation", requestId: "r1", udfPath: "app:mut", args: {} })); + expect(good.close).not.toHaveBeenCalled(); + expect(good.sent.find((m) => m.type === "MutationResponse")).toMatchObject({ requestId: "r1", success: true }); + }); +}); + +describe("parseClientMessage", () => { + it("throws ProtocolError (not SyntaxError/TypeError) for malformed frames", () => { + for (const raw of MALFORMED) expect(() => parseClientMessage(raw)).toThrow(ProtocolError); + }); + + it("accepts every well-formed message shape the client sends", () => { + const valid = [ + { type: "Connect", sessionId: "s" }, + { type: "Connect", supportsQueryDiff: true }, + { + type: "Connect", + sessionId: "s", + clientId: "c", + held: [{ clientId: "c", seq: 1 }], + ackedThrough: [{ clientId: "c", seq: 0 }], + supportsQueryDiff: true, + }, + { type: "ModifyQuerySet", add: [{ queryId: 1, udfPath: "a:b", args: {} }], remove: [2] }, + { type: "ModifyQuerySet", add: [{ queryId: 1, udfPath: "a:b", args: {}, resultHash: "h", sinceTs: 5 }], remove: [] }, + { type: "Mutation", requestId: "r", udfPath: "a:b", args: { x: 1 } }, + { type: "Mutation", requestId: "r", udfPath: "a:b", args: {}, clientId: "c", seq: 3 }, + { type: "MutationBatch", entries: [{ requestId: "r", udfPath: "a:b", args: {}, clientId: "c", seq: 1 }] }, + { type: "Action", requestId: "r", udfPath: "a:b", args: {} }, + { type: "EphemeralPublish", topic: "t", event: { x: 1 } }, + { type: "SetAuth", token: "tok" }, + { type: "SetAuth", token: null }, + { type: "SetAdminAuth", key: "k" }, + ]; + for (const msg of valid) expect(parseClientMessage(JSON.stringify(msg))).toEqual(msg); + }); +}); diff --git a/packages/vite/src/embed.ts b/packages/vite/src/embed.ts index 9dd5ffb9..bb494cd6 100644 --- a/packages/vite/src/embed.ts +++ b/packages/vite/src/embed.ts @@ -210,7 +210,13 @@ export function embedPlugin(options: ConcileVitePluginOptions): Plugin { ping: (onPong: () => void) => { ws.once("pong", onPong); ws.ping(); }, }; runtime.handler.connect(sessionId, syncSocket); - ws.on("message", (data: Buffer) => void runtime.handler.handleMessage(sessionId, data.toString("utf8"))); + ws.on("message", (data: Buffer) => { + // An unhandled rejection here would take down the Vite dev server; drop only this socket. + runtime.handler.handleMessage(sessionId, data.toString("utf8")).catch((e: unknown) => { + console.error("[concile] sync message failed; closing connection:", e); + ws.close(1011); + }); + }); ws.on("close", () => runtime.handler.disconnect(sessionId)); ws.on("error", () => runtime.handler.disconnect(sessionId)); }); From 41b49316d65f0c94a837303d69e32474785cdb3e Mon Sep 17 00:00:00 2001 From: MoatazNoaman Date: Tue, 6 Oct 2026 10:33:43 +0300 Subject: [PATCH 2/3] fix(storage): stop serving uploaded files as active content handleServe echoed the uploader-supplied content type back with no nosniff and no Content-Disposition. With the default FS blobstore the bytes are streamed from the app's own origin, so an upload labelled text/html or image/svg+xml ran script with the app's origin (and, in dev, the dashboard's). Streamed files now always carry X-Content-Type-Options: nosniff, and any type outside an inline-safe allowlist (raster images, audio, video, text/plain, PDF) is served as an attachment. --- .changeset/storage-serve-headers.md | 9 ++++++ packages/storage/src/http.ts | 37 +++++++++++++++++++++-- packages/storage/test/http.test.ts | 46 +++++++++++++++++++++++++++++ 3 files changed, 90 insertions(+), 2 deletions(-) create mode 100644 .changeset/storage-serve-headers.md diff --git a/.changeset/storage-serve-headers.md b/.changeset/storage-serve-headers.md new file mode 100644 index 00000000..c0823dc3 --- /dev/null +++ b/.changeset/storage-serve-headers.md @@ -0,0 +1,9 @@ +--- +"@concile/storage": patch +--- + +Files streamed from `/api/storage/:id` can no longer run script on the app's origin. Every +response now sends `X-Content-Type-Options: nosniff`, and any type outside a small inline-safe +allowlist (raster images, audio, video, `text/plain`, PDF) is served with +`Content-Disposition: attachment`. Previously an uploader-chosen `text/html` or `image/svg+xml` +content type was echoed back and rendered inline. `fetch()`, `` and `