Skip to content
Open
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
8 changes: 8 additions & 0 deletions .changeset/paginate-cursor-clamp.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
---
"@concile/query-engine": patch
---

`paginate` now clamps the client-supplied cursor to the query's own index range. A forged cursor
(or one from a different query) could previously move the scan's start or end outside the range
the query's `eq`/range constraints define, returning rows the query should never see. An
out-of-range cursor now yields the first page or an empty one.
9 changes: 9 additions & 0 deletions .changeset/storage-serve-headers.md
Original file line number Diff line number Diff line change
@@ -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()`, `<img>` and `<video>` are unaffected.
11 changes: 11 additions & 0 deletions .changeset/sync-malformed-frames.md
Original file line number Diff line number Diff line change
@@ -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.
16 changes: 14 additions & 2 deletions packages/cli/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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));
});
Expand Down Expand Up @@ -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);
Expand Down
25 changes: 20 additions & 5 deletions packages/query-engine/src/query-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,25 @@ function base64ToBytes(b64: string): Uint8Array {
return out;
}

/**
* The part of `base` after `cursor` (asc) or before it (desc). The cursor is raw key bytes handed
* back by the client, so it is only ever allowed to NARROW `base`: a cursor outside the query's own
* range (forged, or from a different query) must not move a bound past it, or a page could return
* rows the query's `eq`/range constraints exclude. An out-of-range cursor yields the first page
* (cursor before the range) or an empty one (cursor past it).
*/
function resumeInterval(base: IndexInterval, cursor: Uint8Array, order: ScanOrder): IndexInterval {
if (order === "asc") {
const after = keySuccessor(cursor);
const start = compareKeyBytes(after, base.start) > 0 ? after : base.start;
if (base.end !== null && compareKeyBytes(start, base.end) >= 0) return { start: base.end, end: base.end };
return { start, end: base.end };
}
const end = base.end === null || compareKeyBytes(cursor, base.end) < 0 ? cursor : base.end;
if (compareKeyBytes(end, base.start) <= 0) return { start: base.start, end: base.start };
return { start: base.start, end };
}

/** Stable string form of index-key bytes, for keying the overlay merge map. */
function hexKey(b: Uint8Array): string {
let s = "";
Expand Down Expand Up @@ -170,11 +189,7 @@ export class QueryRuntime {
const tableId = encodeStorageTableId(query.index.tableNumber);

// Resume from the cursor: keys strictly after it (asc) or strictly before it (desc).
let interval: IndexInterval = base;
if (opts.cursor) {
const k = base64ToBytes(opts.cursor);
interval = order === "asc" ? { start: keySuccessor(k), end: base.end } : { start: base.start, end: k };
}
const interval = opts.cursor ? resumeInterval(base, base64ToBytes(opts.cursor), order) : base;

// Read-your-own-writes overlay (see `collect`). `maxScan`/`scanCapped` don't apply here — the
// overlay path scans the whole remaining interval to merge staged writes, and runs only inside a
Expand Down
105 changes: 105 additions & 0 deletions packages/query-engine/test/paginate-cursor.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
import { describe, it, expect } from "vitest";
import { SqliteDocStore, NodeSqliteAdapter } from "@concile/docstore-sqlite";
import {
newDocumentId,
encodeInternalDocumentId,
encodeStorageIndexId,
type InternalDocumentId,
} from "@concile/id-codec";
import type { DocumentValue, IndexWrite } from "@concile/docstore";
import { QueryRuntime, computeIndexUpdates, type IndexSpec, type Query } from "../src/index";

const TABLE = 7002;

const byOwner: IndexSpec = {
table: "notes",
tableNumber: TABLE,
index: "by_owner",
fields: ["owner"],
indexId: encodeStorageIndexId(TABLE, "by_owner"),
};

function makeDoc(id: InternalDocumentId, creation: number, extra: Record<string, unknown>): DocumentValue {
return { _id: encodeInternalDocumentId(id), _creationTime: creation, ...extra } as DocumentValue;
}

async function seed(): Promise<{ store: SqliteDocStore; qr: QueryRuntime }> {
const store = new SqliteDocStore(new NodeSqliteAdapter());
await store.setupSchema();
let ts = 0n;
for (const owner of ["alice", "alice", "bob", "bob", "carol", "carol"]) {
ts++;
const id = newDocumentId(TABLE);
const doc = makeDoc(id, Number(ts), { owner });
const indexWrites: IndexWrite[] = computeIndexUpdates([byOwner], null, doc, id).map((update) => ({ ts, update }));
await store.write([{ ts, id, prev_ts: null, value: { id, value: doc } }], indexWrites, "Error");
}
return { store, qr: new QueryRuntime(store) };
}

const bobs = (order: "asc" | "desc"): Query => ({
index: byOwner,
range: [{ field: "owner", operator: "eq", value: "bob" }],
order,
});

const owners = (page: DocumentValue[]): unknown[] => page.map((d) => (d as Record<string, unknown>).owner);

// Cursors are raw index-key bytes the client hands back; these are ones no real page ever minted.
const BEFORE_ALL = btoa("\x00");
const AFTER_ALL = btoa("\xff");

describe("paginate cursor stays inside the query's range", () => {
it("asc: a cursor before the range restarts at the range, never reading rows below it", async () => {
const { store, qr } = await seed();
const res = await qr.paginate(bobs("asc"), await store.maxTimestamp(), { cursor: BEFORE_ALL, pageSize: 10 });
expect(owners(res.page)).toEqual(["bob", "bob"]);
expect(res.hasMore).toBe(false);
});

it("asc: a cursor past the range yields an empty final page", async () => {
const { store, qr } = await seed();
const res = await qr.paginate(bobs("asc"), await store.maxTimestamp(), { cursor: AFTER_ALL, pageSize: 10 });
expect(res.page).toEqual([]);
expect(res.hasMore).toBe(false);
});

it("desc: a cursor past the range restarts at the range, never reading rows above it", async () => {
const { store, qr } = await seed();
const res = await qr.paginate(bobs("desc"), await store.maxTimestamp(), { cursor: AFTER_ALL, pageSize: 10 });
expect(owners(res.page)).toEqual(["bob", "bob"]);
expect(res.hasMore).toBe(false);
});

it("desc: a cursor before the range yields an empty final page", async () => {
const { store, qr } = await seed();
const res = await qr.paginate(bobs("desc"), await store.maxTimestamp(), { cursor: BEFORE_ALL, pageSize: 10 });
expect(res.page).toEqual([]);
expect(res.hasMore).toBe(false);
});

it("a cursor minted by another owner's query cannot reach that owner's rows", async () => {
const { store, qr } = await seed();
const ts = await store.maxTimestamp();
const alices: Query = { ...bobs("asc"), range: [{ field: "owner", operator: "eq", value: "alice" }] };
const first = await qr.paginate(alices, ts, { pageSize: 1 });
expect(first.nextCursor).not.toBeNull();

const res = await qr.paginate(bobs("asc"), ts, { cursor: first.nextCursor, pageSize: 10 });
expect(owners(res.page)).toEqual(["bob", "bob"]);
});

it("genuine cursors still page through the range in both orders", async () => {
const { store, qr } = await seed();
const ts = await store.maxTimestamp();
for (const order of ["asc", "desc"] as const) {
const p1 = await qr.paginate(bobs(order), ts, { pageSize: 1 });
expect(owners(p1.page)).toEqual(["bob"]);
expect(p1.hasMore).toBe(true);
const p2 = await qr.paginate(bobs(order), ts, { cursor: p1.nextCursor, pageSize: 1 });
expect(owners(p2.page)).toEqual(["bob"]);
expect(p2.page[0]!._id).not.toBe(p1.page[0]!._id);
expect(p2.hasMore).toBe(false);
}
});
});
37 changes: 35 additions & 2 deletions packages/storage/src/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,40 @@ import { isReclaimable } from "./context";
import type { StorageDoc } from "./modules";
import { STORAGE_TABLE_NUMBER } from "./system-table";

/**
* Media types safe to render inline from the app's own origin: none of them can run script in the
* page's origin. Everything else (notably `text/html` and `image/svg+xml`, but also unknown or
* absent types) is served as an attachment, since `contentType` is whatever the uploader claimed.
*/
const INLINE_SAFE_TYPES = new Set([
"image/png",
"image/jpeg",
"image/gif",
"image/webp",
"image/avif",
"text/plain",
"application/pdf",
]);

function isInlineSafe(contentType: string): boolean {
const essence = contentType.split(";")[0]!.trim().toLowerCase();
return INLINE_SAFE_TYPES.has(essence) || essence.startsWith("video/") || essence.startsWith("audio/");
}

/**
* Response headers for bytes streamed from the engine's origin. Uploaded files share that origin
* with the app (and, in dev, the dashboard), so an uploaded HTML/SVG file rendered inline would be
* stored XSS. `nosniff` stops the browser from upgrading a benign type to an active one;
* `attachment` makes a navigation to anything outside the allowlist a download, not a page. Neither
* affects `fetch()`, `<img>` or `<video>`, which ignore `Content-Disposition`.
*/
function servedFileHeaders(contentType: string | null): Headers {
const headers = new Headers({ "x-content-type-options": "nosniff" });
if (contentType !== null) headers.set("content-type", contentType);
if (contentType === null || !isInlineSafe(contentType)) headers.set("content-disposition", "attachment");
return headers;
}

export interface StorageRouteDeps {
runMutation(path: string, args: unknown): Promise<unknown>;
runQuery(path: string, args: unknown): Promise<unknown>;
Expand Down Expand Up @@ -244,8 +278,7 @@ export function storageRoutes(blobStore: BlobStore, deps: StorageRouteDeps): Sto
if (redirectUrl !== null) return new Response(null, { status: 302, headers: { location: redirectUrl } });

const size = doc.size ?? 0;
const headers = new Headers();
if (doc.contentType !== null) headers.set("content-type", doc.contentType);
const headers = servedFileHeaders(doc.contentType);

const range = parseRange(request.headers.get("range"));
if (range !== undefined) {
Expand Down
46 changes: 46 additions & 0 deletions packages/storage/test/http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -475,6 +475,52 @@ describe("GET /api/storage/:id — serve", () => {
expect(new Uint8Array(await response.arrayBuffer())).toEqual(bytes);
});

describe("uploaded bytes can't become active content on the app's origin", () => {
async function serveAs(contentType: string | undefined, range?: string): Promise<Response> {
const blobStore = new FakeBlobStore();
const runtime = await makeRuntime(blobStore);
const routes = storageRoutes(blobStore, routeDeps(runtime));
const id = await uploadReadyFile(runtime, routes, new TextEncoder().encode("<script>alert(1)</script>"), contentType, "public");
return findRoute(routes, "GET", `/api/storage/${id}`).handler(
new Request(`http://localhost/api/storage/${id}`, range !== undefined ? { headers: { range } } : {}),
);
}

it.each(["text/html", "text/html; charset=utf-8", "image/svg+xml", "application/xhtml+xml", "TEXT/HTML"])(
"%s is served nosniff as an attachment, never inline",
async (contentType) => {
const response = await serveAs(contentType);
expect(response.status).toBe(200);
expect(response.headers.get("x-content-type-options")).toBe("nosniff");
expect(response.headers.get("content-disposition")).toBe("attachment");
},
);

it("a file with no content-type is an attachment too (nothing for the browser to sniff into HTML)", async () => {
const response = await serveAs(undefined);
expect(response.headers.get("content-type")).toBeNull();
expect(response.headers.get("x-content-type-options")).toBe("nosniff");
expect(response.headers.get("content-disposition")).toBe("attachment");
});

it.each(["image/png", "image/jpeg", "video/mp4", "audio/mpeg", "text/plain; charset=utf-8", "application/pdf"])(
"%s stays inline (still nosniff)",
async (contentType) => {
const response = await serveAs(contentType);
expect(response.headers.get("content-type")).toBe(contentType);
expect(response.headers.get("x-content-type-options")).toBe("nosniff");
expect(response.headers.get("content-disposition")).toBeNull();
},
);

it("a Range response carries the same protections", async () => {
const response = await serveAs("text/html", "bytes=0-3");
expect(response.status).toBe(206);
expect(response.headers.get("x-content-type-options")).toBe("nosniff");
expect(response.headers.get("content-disposition")).toBe("attachment");
});
});

it("a Range request returns 206 with the correct partial bytes and Content-Range", async () => {
const blobStore = new FakeBlobStore();
const runtime = await makeRuntime(blobStore);
Expand Down
13 changes: 12 additions & 1 deletion packages/sync/src/handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import type { DiffableRange, DiffablePage } from "@concile/executor";
import {
encodeServerMessage,
parseClientMessage,
ProtocolError,
INITIAL_VERSION,
type ClientMessage,
type ServerMessage,
Expand Down Expand Up @@ -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);
Expand Down
1 change: 1 addition & 0 deletions packages/sync/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ export {
compareStateVersion,
isContiguous,
parseClientMessage,
ProtocolError,
encodeServerMessage,
} from "./protocol";

Expand Down
Loading
Loading