From 8d0a4260802c039df0fcce73d4bd3935cfdfc0d4 Mon Sep 17 00:00:00 2001 From: Josh de Leeuw Date: Sun, 13 Sep 2026 17:13:55 -0400 Subject: [PATCH 1/6] Harden the streaming-ingest tier's bounds and add a kill switch A review of the built staging tier (docs/streaming-ingest-design.md) found that one anonymous POST /api/session bought a capability to force up to 10,000 trials x 64 KiB = 640 MB into RTDB for a single session id, that the sweep's assembleSession read a whole session into memory before its MAX_ASSEMBLED_BYTES cap applied (and compared that cap against UTF-16 `.length` rather than real bytes), and that nothing bounded how many sessions one experiment id -- public by construction -- could have open at once. Rules caps: database.rules.json's $seq pattern is now 1-3 digits (1,000 trials) and the per-trial `.length` cap is 16384, both hand-copies of MAX_TRIALS_PER_SESSION / MAX_TRIAL_BYTES in the shared constant module (functions/src/staging-assembly.ts), which api-session-start.ts already read to answer the client. __tests__/rules-constants.test.js parses the rules file and asserts its literals still match those constants. The rules and code comments now say plainly that RTDB's `.length` counts UTF-16 code units, not bytes, and that MAX_ASSEMBLED_BYTES (measured with Buffer.byteLength) is the byte-accurate backstop. lastFlushAt: was `newData.isNumber()`, letting a client write any timestamp including one far in the future -- which made disconnectedSince() treat every later disconnect stamp as already answered, so such a session could never be swept as abandoned. Now requires `newData.val() == now`, like the other server-stamped timestamps in the tree. Bounded assembly read: assembleSession no longer does one `.get()` of a session's whole trials node. It now pages through orderByKey().startAfter(lastKey) via a new pure accumulator, assembleTrialsPaged (staging-assembly.ts), which stops asking for another page the moment the accumulated byte size would cross MAX_ASSEMBLED_BYTES -- so an over-cap session is truncated without ever being pulled into memory in full. Both assembleTrials and assembleTrialsPaged now cost each trial with Buffer.byteLength instead of `.length`. Per-experiment concurrency cap: openSessionCounts/{experimentId} (functions/src/staging.ts) is a transactional counter, incremented in openSession() and decremented in discardSession(), capping an experiment at MAX_OPEN_SESSIONS_PER_EXPERIMENT (500) concurrently open sessions. It self-heals negative drift in the transaction itself and the sweep runs reconcileOpenSessionCounts each pass, scoped to that run's candidate experiments, to correct a missed decrement. Deliberately not a per-IP counter. A refusal returns the existing 503 SESSION_START_ERROR shape, which the plugin already falls back from. Kill switch: STREAMING_ENABLED=false makes POST /api/session return that same 503 shape without touching Firestore or RTDB. Default (unset) is enabled, so no existing deployment's behavior changes. Updated pages/docs/api.js's example response and the design doc (new "Hardening pass" section) to match. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Q8xf16ov88M3KojPVfQwJT --- __tests__/database-rules.test.js | 70 ++++- database.rules.json | 76 ++++- docs/streaming-ingest-design.md | 56 ++++ functions/.env.datapipe-test | 7 + .../src/__tests__/rules-constants.test.js | 109 ++++++++ .../src/__tests__/staging-assembly.test.js | 153 ++++++++++- .../src/__tests__/staging-emulator.test.js | 117 +++++++- functions/src/api-session-start.ts | 15 + functions/src/scheduled-staging-sweep.ts | 45 +++ functions/src/staging-assembly.ts | 225 +++++++++++++-- functions/src/staging.ts | 260 ++++++++++++++++-- pages/docs/api.js | 10 +- 12 files changed, 1065 insertions(+), 78 deletions(-) create mode 100644 functions/src/__tests__/rules-constants.test.js diff --git a/__tests__/database-rules.test.js b/__tests__/database-rules.test.js index 080fbb2..c655f11 100644 --- a/__tests__/database-rules.test.js +++ b/__tests__/database-rules.test.js @@ -108,6 +108,26 @@ describe('reads', () => { it('denies reading a session meta node', async () => { await assertFails(client().ref(`staging/${OPEN}/meta`).once('value')); }); + + it('denies reading the per-experiment concurrency counter', async () => { + await assertFails(client().ref('openSessionCounts/exp123').once('value')); + }); +}); + +describe('openSessionCounts', () => { + // The per-experiment concurrency cap's ground truth + // (functions/src/staging.ts). Admin-SDK-only, on the same terms as + // openSessions: no client write path exists, or should ever be added. + it('denies a client writing its own counter', async () => { + await assertFails(client().ref('openSessionCounts/exp123').set(0)); + }); + + it('denies a client incrementing another experiment\'s counter', async () => { + await testEnv.withSecurityRulesDisabled(async (context) => { + await context.database().ref('openSessionCounts/exp123').set(499); + }); + await assertFails(client().ref('openSessionCounts/exp123').set(500)); + }); }); describe('the openSessions gate', () => { @@ -185,15 +205,18 @@ describe('append-only trials', () => { describe('size and shape caps', () => { // Property 3. Rules cannot count children, so these two caps are the only - // bound expressible here. - it('accepts a trial at exactly the 64 KiB cap', async () => { - const atCap = JSON.stringify({ v: 'x'.repeat(65536 - 12) }); - expect(atCap.length).toBeLessThanOrEqual(65536); + // bound expressible here. 16384 and 1,000 are MAX_TRIAL_BYTES and + // MAX_TRIALS_PER_SESSION (functions/src/staging-assembly.ts); see + // functions/src/__tests__/rules-constants.test.js for the test that keeps + // this file's literals from drifting away from those constants. + it('accepts a trial at exactly the 16 KiB cap', async () => { + const atCap = JSON.stringify({ v: 'x'.repeat(16384 - 12) }); + expect(atCap.length).toBeLessThanOrEqual(16384); await assertSucceeds(client().ref(`staging/${OPEN}/trials/0`).set(atCap)); }); it('denies a trial one byte over the cap', async () => { - await assertFails(client().ref(`staging/${OPEN}/trials/0`).set('x'.repeat(65537))); + await assertFails(client().ref(`staging/${OPEN}/trials/0`).set('x'.repeat(16385))); }); it('denies a non-string trial value', async () => { @@ -204,11 +227,11 @@ describe('size and shape caps', () => { }); it('accepts the highest permitted sequence number', async () => { - await assertSucceeds(client().ref(`staging/${OPEN}/trials/9999`).set('{"a":1}')); + await assertSucceeds(client().ref(`staging/${OPEN}/trials/999`).set('{"a":1}')); }); - it('denies a sequence number past the 10,000-trial ceiling', async () => { - await assertFails(client().ref(`staging/${OPEN}/trials/10000`).set('{"a":1}')); + it('denies a sequence number past the 1,000-trial ceiling', async () => { + await assertFails(client().ref(`staging/${OPEN}/trials/1000`).set('{"a":1}')); }); it('denies a non-numeric sequence key', async () => { @@ -223,9 +246,30 @@ describe('meta', () => { await assertFails(client().ref(`staging/${OPEN}/meta/startedAt`).set(0)); }); - it('allows lastFlushAt to be refreshed on every flush', async () => { - await assertSucceeds(client().ref(`staging/${OPEN}/meta/lastFlushAt`).set(Date.now())); - await assertSucceeds(client().ref(`staging/${OPEN}/meta/lastFlushAt`).set(Date.now() + 1)); + it('allows lastFlushAt to be refreshed on every flush, via the server timestamp placeholder', async () => { + await assertSucceeds(client().ref(`staging/${OPEN}/meta/lastFlushAt`).set(SERVER_TIME)); + await assertSucceeds(client().ref(`staging/${OPEN}/meta/lastFlushAt`).set(SERVER_TIME)); + }); + + it('denies a client-chosen lastFlushAt in the past', async () => { + // `newData.isNumber()` alone would accept any client clock. A stale value + // is not merely wrong, it is actively dangerous in the other direction + // from the future-dated case below: it could make a live participant's + // session look stalled if the plugin's own logic ever compared lastFlushAt + // against wall-clock time client-side, though today's exposure is smaller + // since disconnectedSince() only compares it against disconnect stamps. + await assertFails(client().ref(`staging/${OPEN}/meta/lastFlushAt`).set(Date.now() - 3600000)); + }); + + it('denies a client-chosen lastFlushAt in the future', async () => { + // THE BUG THIS CLOSES. disconnectedSince() (staging-assembly.ts) treats a + // flush timestamped after a disconnect stamp as proof the participant + // reconnected. `newData.isNumber()` alone let a client write ANY number, + // including one far in the future -- which would make every later + // disconnect stamp look "already answered" forever, so the session could + // never be swept as abandoned no matter how long the socket stayed + // dropped. + await assertFails(client().ref(`staging/${OPEN}/meta/lastFlushAt`).set(Date.now() + 86400000)); }); it('allows a disconnect stamp to be answered by a reconnect mark', async () => { @@ -259,7 +303,7 @@ describe('the real flush shape', () => { 'trials/0': '{"trial":0}', 'trials/1': '{"trial":1}', 'trials/2': '{"trial":2}', - 'meta/lastFlushAt': Date.now(), + 'meta/lastFlushAt': SERVER_TIME, }) ); }); @@ -272,7 +316,7 @@ describe('the real flush shape', () => { client().ref(`staging/${OPEN}`).update({ 'trials/1': '{"trial":"rewritten"}', 'trials/2': '{"trial":2}', - 'meta/lastFlushAt': Date.now(), + 'meta/lastFlushAt': SERVER_TIME, }) ); }); diff --git a/database.rules.json b/database.rules.json index b6c331f..0bc4c0e 100644 --- a/database.rules.json +++ b/database.rules.json @@ -29,10 +29,34 @@ // retrying client idempotent: re-flushing the same seq is a no-op refusal, // not a corruption. See "ORDERING AND DUPLICATES" below. // -// 3. SIZE-CAPPED. 64 KiB per trial, and a sequence number of at most four -// digits (10,000 trials). Rules cannot count children, so this is the only +// 3. SIZE-CAPPED. 16 KiB per trial, and a sequence number of at most three +// digits (1,000 trials). Rules cannot count children, so this is the only // bound expressible here; the sweep enforces a total-bytes ceiling of its -// own when it assembles a session (functions/src/staging.ts). +// own when it assembles a session, reading it in PAGES rather than in one +// piece so an over-cap session is never pulled into memory whole +// (functions/src/staging.ts, functions/src/staging-assembly.ts). +// +// THE SHARED CONSTANT MODULE. 16384 and the 3-digit pattern below are +// hand-copies of functions/src/staging-assembly.ts's MAX_TRIAL_BYTES and +// MAX_TRIALS_PER_SESSION -- rules cannot import a TypeScript module, so +// there is no way to make this file read them directly. +// __tests__/rules-constants.test.js parses this file and asserts the two +// literals still match those constants; api-session-start.ts sends the +// same constants to the client, so the rules, the endpoint's response and +// this comment cannot silently drift apart. Change a cap in ONE place +// (staging-assembly.ts) and update this file's regex / `.length` literal +// in the same commit, or that test fails. +// +// UNITS: `newData.val().length` below counts UTF-16 CODE UNITS, not +// bytes -- a JS/RTDB string-length primitive, not a byte count. A trial +// that is plain ASCII costs one byte per unit; a trial full of multibyte +// content (CJK, emoji, most non-Latin scripts) can cost up to ~3 bytes per +// unit for the same `.length`. So 16384 is a UTF-16 ceiling, and the real +// byte ceiling for an adversarial payload is up to ~3x that. This rule is +// still worth having -- it is the only bound RTDB can enforce before a +// byte ever reaches a function -- but it is NOT the byte-accurate cap. +// staging-assembly.ts's MAX_ASSEMBLED_BYTES, measured with +// Buffer.byteLength against the real read, is the byte-accurate backstop. // // 4. GATED ON AN OPEN SESSION. Every write requires // openSessions/{sessionId} to exist, and only POST /api/session @@ -102,6 +126,22 @@ ".indexOn": ["expiresAt"] }, + // openSessionCounts/{experimentId} = number of that experiment's + // currently-open sessions. + // + // The per-experiment concurrency cap's ground truth (tryAdmitSession, + // releaseOpenSessionSlot, reconcileOpenSessionCounts in + // functions/src/staging.ts), bounding what one experiment id -- public, + // by construction, in the experiment's own JavaScript -- can cost in + // staged RTDB storage no matter how many times POST /api/session is + // called for it. Written and read only by the Admin SDK, via + // transactions, exactly like openSessions above: no block here carries a + // `.read` or `.write`, so it inherits the root deny and stays invisible + // and unwritable to every client. Written down for the same reason as + // openSessions -- so nobody later "fixes" this absence by adding a rule, + // which could only weaken it. + "openSessionCounts": {}, + "staging": { "$sessionId": { "meta": { @@ -114,10 +154,22 @@ // Refreshed on every flush, in the same multi-path update that // carries the trials, so it costs nothing extra. It is the sweep's // backstop liveness signal for a client that died before - // onDisconnect could fire. + // onDisconnect could fire -- disconnectedSince() (staging-assembly.ts) + // treats a flush timestamped after a disconnect stamp as proof the + // participant is back. + // + // `== now`, not `isNumber()`: a client that could write ANY number + // could write one far in the future, and disconnectedSince() would + // then treat every later disconnect stamp as already answered -- + // the session would NEVER be markable abandoned, no matter how long + // the socket stayed dropped, because lastFlushAt would always sort + // after it. Requiring the server-resolved value (the same + // ServerValue.TIMESTAMP placeholder disconnects/reconnects already + // require) closes that off the same way it closes off a backdated + // disconnect stamp below. "lastFlushAt": { ".write": "root.child('openSessions').child($sessionId).exists()", - ".validate": "newData.isNumber()" + ".validate": "newData.val() == now" }, // ABANDONMENT: one write-once slot per connection, the standard // Firebase presence pattern. @@ -200,11 +252,15 @@ // no way to bound nesting, and the size cap below then means // exactly what it says. // - // 64 KiB is roughly 13x the top of the design doc's 0.5-5 KB - // per-trial range -- generous enough for a trial carrying an - // embedded response blob, tight enough that the 10,000-key - // ceiling is a real bound rather than a formality. - ".validate": "$seq.matches(/^[0-9]{1,4}$/) && newData.isString() && newData.val().length <= 65536" + // 16 KiB is roughly 3x the top of the design doc's 0.5-5 KB + // per-trial range -- generous enough for a trial carrying a small + // embedded response blob, tight enough that the 1,000-key ceiling + // below is a real structural bound: MAX_TRIALS_PER_SESSION x + // MAX_TRIAL_BYTES x the up-to-3x UTF-16-to-byte multiplier + // documented above stays well inside the sweep's memory + // allocation, which the previous 10,000 x 64 KiB combination did + // not (functions/src/staging-assembly.ts has the arithmetic). + ".validate": "$seq.matches(/^[0-9]{1,3}$/) && newData.isString() && newData.val().length <= 16384" } }, // Nothing but meta/ and trials/ may exist under a session. Defence in diff --git a/docs/streaming-ingest-design.md b/docs/streaming-ingest-design.md index 59ad5d9..f7fa316 100644 --- a/docs/streaming-ingest-design.md +++ b/docs/streaming-ingest-design.md @@ -419,6 +419,62 @@ Two prerequisites that are NOT code and are easy to miss: `firebase deploy --only database` will work. It is not auto-provisioned. 2. The deploy line in both workflows now includes `database`. +## Hardening pass (2026-09-13) + +A review of the built system found two gaps between what the rules and the +sweep were assumed to bound and what they actually bounded, plus two missing +operational controls. All four are answered here rather than in a new design: + +- **The rules caps were structurally too large.** 10,000 trials x 64 KiB was + the worst case one admitted session id could force the sweep to read -- + 640 MB before any UTF-16-to-byte multiplier, well past the sweep's 256MiB + allocation. Lowered to **1,000 trials x 16 KiB** (`$seq` is now 1-3 digits; + the per-trial `.length` cap is now 16384) in `database.rules.json`, and both + numbers are read from ONE constant module, + `functions/src/staging-assembly.ts` (`MAX_TRIALS_PER_SESSION`, + `MAX_TRIAL_BYTES`), which api-session-start.ts already used to answer the + client and which `__tests__/rules-constants.test.js` now asserts the rules + file's literals still match. RTDB's `.length` counts UTF-16 code units, not + bytes -- multibyte content can cost up to ~3x as many real bytes for the + same `.length` -- so the rules cap is a UTF-16 ceiling, not a byte ceiling; + `MAX_ASSEMBLED_BYTES`, measured with `Buffer.byteLength` against the real + read, is the byte-accurate backstop. + +- **The sweep read a whole session into memory before its size cap applied.** + `assembleSession` used to `.get()` the entire `staging/{id}/trials` node, + then apply `MAX_ASSEMBLED_BYTES` while walking the in-memory result -- so an + over-cap session was pulled into memory in full before being discarded. + `assembleSession` now pages the read (`orderByKey().startAfter(lastKey)`, + 200 trials per page) through a new pure accumulator, + `assembleTrialsPaged` (`staging-assembly.ts`), which stops asking for + another page the moment the accumulated byte size would cross the cap. The + same review also found the existing byte-cap comparison used `raw.length` + (UTF-16 units) against a byte constant; both `assembleTrials` and + `assembleTrialsPaged` now cost each trial with `Buffer.byteLength`. + +- **One anonymous `POST /api/session` bought an unbounded write capability.** + Nothing stopped that call being looped for one experiment id, which is + public by construction. `openSessionCounts/{experimentId}` + (`functions/src/staging.ts`) is a transactional counter, incremented in + `openSession()` and decremented in `discardSession()`, capping an + experiment at `MAX_OPEN_SESSIONS_PER_EXPERIMENT` (500) concurrently open + sessions -- comfortably above any lecture-hall study, per the instance- + ceiling reasoning above. It self-heals negative drift (a stored value below + zero is treated as zero) and the sweep runs `reconcileOpenSessionCounts` + each pass, scoped to that run's candidate experiments, to correct positive + drift (a missed decrement) within a few runs. Deliberately **not** a per-IP + counter -- DataPipe's own code must never read or store a participant's IP + address (`pages/docs/privacy.js`) -- it counts sessions per experiment, the + same unit `maxSessions` already limits. A refused start returns the + existing `503 SESSION_START_ERROR` shape, so the plugin's documented + fallback (submit once at the end) applies unchanged. + +- **A kill switch.** `STREAMING_ENABLED=false` (checked first, before any + Firestore or RTDB access) makes `POST /api/session` return the same `503 + SESSION_START_ERROR` shape without touching either store. Default (unset) + is enabled, so no existing deployment's behaviour changes. Documented as a + commented example in `functions/.env.datapipe-test`. + ## Related finding (2026-09-02) While measuring production volume: **`createdAt` and `lastRequestAt` are absent diff --git a/functions/.env.datapipe-test b/functions/.env.datapipe-test index 6987681..ca26764 100644 --- a/functions/.env.datapipe-test +++ b/functions/.env.datapipe-test @@ -41,3 +41,10 @@ ZENODO_REDIRECT_URI=https://datapipe-test.web.app/oauth2/connect # personal account. The id and the TEST_ZENODO_CLIENT_SECRET GitHub secret # come from that one registration and have to be rotated together. ZENODO_CLIENT_ID=i8LEvZPTBwO8ipksP6MmujspHM8QOh7FTrgSa1v9 + +# --- Streaming-ingest kill switch (docs/streaming-ingest-design.md) --- +# Uncomment to stop POST /api/session from touching Firestore or RTDB at all; +# the plugin's documented fallback is to submit once at the end, as it did +# before this feature existed. Any value other than the exact string "false", +# including unset (the default here), leaves the endpoint enabled. +# STREAMING_ENABLED=false diff --git a/functions/src/__tests__/rules-constants.test.js b/functions/src/__tests__/rules-constants.test.js new file mode 100644 index 0000000..2aceb1a --- /dev/null +++ b/functions/src/__tests__/rules-constants.test.js @@ -0,0 +1,109 @@ +/** + * @jest-environment node + * + * database.rules.json cannot import functions/src/staging-assembly.ts -- RTDB + * rules are not TypeScript -- so its `$seq` regex and per-trial `.length` cap + * are HAND-COPIES of MAX_TRIALS_PER_SESSION and MAX_TRIAL_BYTES. A hand-copy + * can drift from the constant it was copied from without either file's own + * tests noticing: database-rules.test.js only ever checks the rules' actual + * behaviour against literals it writes itself, and staging-assembly.test.js + * never reads database.rules.json at all. + * + * This suite is the thing that would catch that drift: it parses + * database.rules.json as data and asserts its literals equal the constants + * directly, so changing one without the other fails here rather than at + * runtime, when the endpoint tells a client one cap and the rules enforce a + * different one. + */ + +const fs = require("fs"); +const path = require("path"); +const { + MAX_TRIAL_BYTES, + MAX_TRIALS_PER_SESSION, +} = require("../../lib/staging-assembly.js"); + +const RULES_PATH = path.join(__dirname, "..", "..", "..", "database.rules.json"); + +/** + * database.rules.json is JSON with `//` line comments, which RTDB's rules + * loader accepts but JSON.parse does not. Strips them the same way a real + * JSONC-aware loader would: character by character, tracking whether the + * cursor is inside a double-quoted string (respecting `\"` so a quote inside + * a regex literal like `/^[0-9]{1,3}$/` -- itself inside a JSON string -- + * cannot be mistaken for the string's closing quote), so that `//` inside a + * string (there is none in this file today, but a rule expression could + * legally contain one) is never treated as a comment. + */ +function stripJsonComments(text) { + let out = ""; + let inString = false; + for (let i = 0; i < text.length; i++) { + const ch = text[i]; + if (inString) { + out += ch; + if (ch === "\\") { + // Copy the escaped character verbatim without re-examining it -- + // otherwise `\"` would end the string one character early. + i++; + if (i < text.length) out += text[i]; + continue; + } + if (ch === '"') inString = false; + continue; + } + if (ch === '"') { + inString = true; + out += ch; + continue; + } + if (ch === "/" && text[i + 1] === "/") { + while (i < text.length && text[i] !== "\n") i++; + out += "\n"; + continue; + } + out += ch; + } + return out; +} + +function loadRules() { + const raw = fs.readFileSync(RULES_PATH, "utf8"); + return JSON.parse(stripJsonComments(raw)); +} + +describe("database.rules.json agrees with functions/src/staging-assembly.ts", () => { + let rules; + let trialRule; + + beforeAll(() => { + rules = loadRules(); + trialRule = rules.rules.staging.$sessionId.trials.$seq[".validate"]; + }); + + it("caps the per-trial length at exactly MAX_TRIAL_BYTES", () => { + const match = trialRule.match(/newData\.val\(\)\.length\s*<=\s*(\d+)/); + expect(match).not.toBeNull(); + expect(Number(match[1])).toBe(MAX_TRIAL_BYTES); + }); + + it("bounds the $seq key to exactly the digit width MAX_TRIALS_PER_SESSION implies", () => { + // The number of digits the rule's regex allows is what actually bounds + // trial count (10 ** digits, since $seq is unpadded), not a literal copy + // of MAX_TRIALS_PER_SESSION itself -- so this derives the expected width + // from the constant instead of hand-copying "3". + const match = trialRule.match(/\$seq\.matches\(\/\^\[0-9\]\{1,(\d+)\}\$\/\)/); + expect(match).not.toBeNull(); + const digits = Number(match[1]); + expect(10 ** digits).toBe(MAX_TRIALS_PER_SESSION); + }); + + it("requires the disconnect/reconnect slot cap comment's own bound (sanity: rules file still parses as one object)", () => { + // Not a constants check -- just confirms stripJsonComments above did not + // silently produce a truncated or malformed object that the two tests + // above would pass against vacuously (e.g. `rules.rules` being undefined + // would make the `.length` and `$seq` lookups above throw, not silently + // pass, but this pins the shape explicitly all the same). + expect(rules.rules.staging.$sessionId.trials.$seq).toHaveProperty([".validate"]); + }); +}); diff --git a/functions/src/__tests__/staging-assembly.test.js b/functions/src/__tests__/staging-assembly.test.js index 8be5e1c..95365cb 100644 --- a/functions/src/__tests__/staging-assembly.test.js +++ b/functions/src/__tests__/staging-assembly.test.js @@ -9,8 +9,10 @@ import { assembleTrials, + assembleTrialsPaged, partialFilenameFor, disconnectedSince, + streamingEnabled, MAX_ASSEMBLED_BYTES, } from '../../lib/staging-assembly.js'; @@ -19,6 +21,23 @@ function staged(...jsonStrings) { return Object.fromEntries(jsonStrings.map((s, i) => [String(i), s])); } +/** + * An in-memory fetchPage for assembleTrialsPaged: pages through `entries` + * (already-sorted [key, value] pairs) `pageSize` at a time, and counts how + * many times it was called. + */ +function fakePager(entries, pageSize) { + const calls = { count: 0 }; + const fetchPage = async (afterKey) => { + calls.count++; + const startIndex = + afterKey === null ? 0 : entries.findIndex(([k]) => k === afterKey) + 1; + const page = entries.slice(startIndex, startIndex + pageSize); + return { entries: page, done: startIndex + page.length >= entries.length }; + }; + return { fetchPage, calls }; +} + describe('assembleTrials', () => { it('produces a JSON array that parses back to the original trials', () => { const result = assembleTrials(staged('{"trial":0}', '{"trial":1}', '{"trial":2}')); @@ -100,9 +119,9 @@ describe('assembleTrials', () => { }); it('truncates at the assembly ceiling rather than exhausting memory', () => { - // The rules permit 10,000 trials x 64 KiB. A crashed sweep is an outage - // that also lets every other abandoned session pile up behind it, so - // assembly stops and flags instead. + // The rules permit up to 1,000 trials x 16 KiB (functions/src/staging.ts). + // A crashed sweep is an outage that also lets every other abandoned + // session pile up behind it, so assembly stops and flags instead. const big = JSON.stringify({ v: 'x'.repeat(60000) }); const trials = {}; for (let i = 0; i < 500; i++) trials[String(i)] = big; @@ -116,6 +135,134 @@ describe('assembleTrials', () => { // half-written array. expect(() => JSON.parse(result.data)).not.toThrow(); }); + + it('measures the cap in bytes, not UTF-16 length', () => { + // The review finding this closes: the old comparison used `raw.length`, + // which for multibyte content undercounts real bytes by up to ~3x (see + // MAX_TRIAL_BYTES's doc in staging-assembly.ts). A trial made entirely of + // a 3-byte-per-unit character has to be judged on its real byte size, not + // its (much smaller-looking) UTF-16 length. + const wide = 'あ'; // U+3042, 1 UTF-16 unit, 3 UTF-8 bytes + const trial = JSON.stringify({ v: wide.repeat(10000) }); // ~10,000 units, ~30,000 bytes + expect(Buffer.byteLength(trial, 'utf8')).toBeGreaterThan(trial.length * 2); + + const trials = {}; + for (let i = 0; i < 2000; i++) trials[String(i)] = trial; + + const result = assembleTrials(trials); + + // If the comparison still used `.length`, this would fit ~800 copies + // before crossing MAX_ASSEMBLED_BYTES (24 MiB / ~30,000). Measured in + // real bytes it fits far fewer. + expect(result.truncated).toBe(true); + expect(Buffer.byteLength(result.data, 'utf8')).toBeLessThanOrEqual(MAX_ASSEMBLED_BYTES); + }); +}); + +describe('assembleTrialsPaged', () => { + it('produces the same result as assembleTrials for data that fits in one page', async () => { + const trials = staged('{"trial":0}', '{"trial":1}', '{"trial":2}'); + const entries = Object.entries(trials); + + const { fetchPage } = fakePager(entries, 10); + const result = await assembleTrialsPaged(fetchPage); + + expect(result.data).toBe(assembleTrials(trials).data); + expect(JSON.parse(result.data)).toEqual([ + { trial: 0 }, + { trial: 1 }, + { trial: 2 }, + ]); + expect(result.trialCount).toBe(3); + expect(result.gaps).toBe(0); + expect(result.truncated).toBe(false); + expect(result.pagesFetched).toBe(1); + }); + + it('preserves ascending order across a page boundary', async () => { + const trials = {}; + for (let i = 0; i < 25; i++) trials[String(i)] = `{"trial":${i}}`; + const entries = Object.keys(trials) + .sort((a, b) => Number(a) - Number(b)) + .map((k) => [k, trials[k]]); + + const { fetchPage, calls } = fakePager(entries, 10); + const result = await assembleTrialsPaged(fetchPage); + + expect(JSON.parse(result.data).map((t) => t.trial)).toEqual( + Array.from({ length: 25 }, (_, i) => i) + ); + expect(result.trialCount).toBe(25); + expect(calls.count).toBe(3); // 10 + 10 + 5 + expect(result.pagesFetched).toBe(3); + }); + + it('tolerates gaps and counts skipped trials across pages', async () => { + const entries = [ + ['0', '{"trial":0}'], + ['1', '{not json'], + ['4', '{"trial":4}'], + ]; + + const { fetchPage } = fakePager(entries, 2); + const result = await assembleTrialsPaged(fetchPage); + + expect(JSON.parse(result.data)).toEqual([{ trial: 0 }, { trial: 4 }]); + expect(result.skipped).toBe(1); + expect(result.gaps).toBe(2); // sequence numbers 2 and 3 never arrived + }); + + it('returns an empty assembly when there is nothing staged, in one call', async () => { + const { fetchPage, calls } = fakePager([], 200); + + const result = await assembleTrialsPaged(fetchPage); + + expect(result.data).toBe('[]'); + expect(result.trialCount).toBe(0); + expect(calls.count).toBe(1); + expect(result.pagesFetched).toBe(1); + }); + + it('stops fetching pages once the byte cap is crossed, without reading the rest', async () => { + // THE FIX THIS TEST PINS: the review found the sweep read a whole session + // into memory before its size cap applied. Here the fake backing store + // holds far more than fits under MAX_ASSEMBLED_BYTES, and the assertion + // is on `calls.count` -- proof the function stopped asking for more + // pages, not just that it stopped keeping what it already had. + const big = JSON.stringify({ v: 'x'.repeat(60000) }); // ~60KB/trial + const totalTrials = 1000; // ~60MB backing store if it were all read + const entries = Array.from({ length: totalTrials }, (_, i) => [String(i), big]); + const pageSize = 50; // ~3MB/page + + const { fetchPage, calls } = fakePager(entries, pageSize); + const result = await assembleTrialsPaged(fetchPage); + + expect(result.truncated).toBe(true); + expect(Buffer.byteLength(result.data, 'utf8')).toBeLessThanOrEqual(MAX_ASSEMBLED_BYTES); + // MAX_ASSEMBLED_BYTES (24MiB) / ~60KB per trial is ~400 trials, or 8 + // pages of 50 -- nowhere near the 20 pages totalTrials/pageSize would take + // to read everything. + const pagesToReadEverything = totalTrials / pageSize; + expect(calls.count).toBeLessThan(pagesToReadEverything); + expect(result.pagesFetched).toBe(calls.count); + }); +}); + +describe('streamingEnabled', () => { + it('defaults to enabled when unset', () => { + expect(streamingEnabled({})).toBe(true); + }); + + it('stays enabled for any value other than the literal string "false"', () => { + expect(streamingEnabled({ STREAMING_ENABLED: 'true' })).toBe(true); + expect(streamingEnabled({ STREAMING_ENABLED: '' })).toBe(true); + expect(streamingEnabled({ STREAMING_ENABLED: 'FALSE' })).toBe(true); + expect(streamingEnabled({ STREAMING_ENABLED: '0' })).toBe(true); + }); + + it('disables only on the exact string "false"', () => { + expect(streamingEnabled({ STREAMING_ENABLED: 'false' })).toBe(false); + }); }); describe('partialFilenameFor', () => { diff --git a/functions/src/__tests__/staging-emulator.test.js b/functions/src/__tests__/staging-emulator.test.js index 948abc7..1c6f99e 100644 --- a/functions/src/__tests__/staging-emulator.test.js +++ b/functions/src/__tests__/staging-emulator.test.js @@ -187,7 +187,7 @@ describe("POST /api/session", () => { // The plugin gets its configuration from the server, so one published // build can talk to both datapipe-test and production. expect(body.databaseURL).toEqual(expect.any(String)); - expect(body.maxTrialBytes).toBe(65536); + expect(body.maxTrialBytes).toBe(16384); // The disconnect-slot cap the rules enforce, so the plugin stops arming at // it instead of having stamps refused. expect(body.maxDisconnects).toBe(20); @@ -267,6 +267,100 @@ describe("POST /api/session", () => { }); }); +describe("the per-experiment concurrency cap", () => { + // MAX_OPEN_SESSIONS_PER_EXPERIMENT (functions/src/staging-assembly.ts). + // Seeded directly rather than opened by minting 500 real sessions: the + // mechanism under test is tryAdmitSession's transaction reading and + // bounding openSessionCounts/{experimentId}, not the accumulation of 500 + // HTTP round trips to get there. + const CAP = 500; + + it("refuses the (cap+1)th session and leaves the counter untouched", async () => { + const experimentID = await makeExperiment(); + await rtdb.ref(`openSessionCounts/${experimentID}`).set(CAP); + + const { status, body } = await startSession({ experimentID }); + + expect(status).toBe(503); + expect(body.error).toBe("SESSION_START_ERROR"); + // The transaction must abort without writing -- a refused admission must + // not itself be what pushes a counter around. + expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP); + // And no capability record was minted for the refused attempt. + const record = (await rtdb.ref("openSessions").get()).val() || {}; + expect(Object.values(record).some((r) => r.experimentId === experimentID)).toBe(false); + }); + + it("admits the cap-th session and then refuses the next one", async () => { + const experimentID = await makeExperiment(); + await rtdb.ref(`openSessionCounts/${experimentID}`).set(CAP - 1); + + const { status: admitted } = await startSession({ experimentID }); + expect(admitted).toBe(200); + expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP); + + const { status: refused, body } = await startSession({ experimentID }); + expect(refused).toBe(503); + expect(body.error).toBe("SESSION_START_ERROR"); + }); + + it("lets a new session in once a prior one completes and frees its slot", async () => { + const experimentID = await makeExperiment(); + await rtdb.ref(`openSessionCounts/${experimentID}`).set(CAP - 1); + + const { status: firstStatus, body: first } = await startSession({ experimentID }); + expect(firstStatus).toBe(200); + + const { status: refused } = await startSession({ experimentID }); + expect(refused).toBe(503); + + // Completion runs discardSession(), which must release the slot the first + // session reserved at admission. + await saveData({ + experimentID, + filename: "p01.csv", + data: "trial_type\nhtml-keyboard-response\n", + sessionId: first.sessionId, + }); + expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP - 1); + + const { status: secondStatus } = await startSession({ experimentID }); + expect(secondStatus).toBe(200); + }); + + it("lets a new session in once an abandoned one is swept, freeing its slot", async () => { + const experimentID = await makeExperiment(); + await rtdb.ref(`openSessionCounts/${experimentID}`).set(CAP - 1); + + const { body: first } = await startSession({ experimentID }); + await markAbandoned(first.sessionId); + + const stats = await sweepAbandonedSessions(new Set([first.sessionId])); + expect(stats.discarded).toBe(1); // nothing was staged, so it is discarded not recovered + expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP - 1); + + const { status } = await startSession({ experimentID }); + expect(status).toBe(200); + }); +}); + +describe("the streaming kill switch", () => { + // STREAMING_ENABLED=false is a pure predicate (streamingEnabled(), unit + // tested in staging-assembly.test.js) precisely because the functions + // emulator cannot be re-configured mid-suite to exercise the disabled + // branch here. What this suite CAN pin is the other half: every deployment + // today has STREAMING_ENABLED unset, and that must keep minting sessions + // exactly as it always has. + it("is enabled by default", async () => { + const experimentID = await makeExperiment(); + + const { status, body } = await startSession({ experimentID }); + + expect(status).toBe(200); + expect(body.sessionId).toEqual(expect.any(String)); + }); +}); + describe("completion", () => { it("drops the staged copy once a gate refuses the submission", async () => { // Keeping it would let the sweep re-offer the session as a .partial.json, @@ -474,6 +568,27 @@ describe("the abandonment sweep", () => { expect(entry.failureReason).not.toContain("missing"); }); + it("recovers a session with more trials than fit in one assembly page, in full", async () => { + // assembleSession (staging.ts) reads the staging tree in pages of 200 -- + // this is the one thing the pure paging tests (staging-assembly.test.js) + // cannot cover, because they fake the page fetcher rather than exercising + // RTDB's own orderByKey().startAfter() pagination. 250 trials forces a + // second page. + const experimentID = await makeExperiment(); + const { body } = await startSession({ experimentID, filename: "p20.csv" }); + await stageTrials(body.sessionId, 250); + await markAbandoned(body.sessionId); + + const stats = await sweepAbandonedSessions(new Set([body.sessionId])); + + expect(stats.recovered).toBe(1); + const entries = await queueEntriesFor(experimentID); + expect(entries).toHaveLength(1); + expect(entries[0].failureReason).toContain("250 trials"); + expect(entries[0].failureReason).not.toContain("missing"); + expect(entries[0].failureReason).not.toContain("truncated"); + }); + it("records sweep health on every run, including one that did nothing", async () => { // A broken sweep is the expensive failure here: orphaned staging data // accumulates at ~190x Cloud Storage's per-GB rate while the scheduled diff --git a/functions/src/api-session-start.ts b/functions/src/api-session-start.ts index 2b01507..a5aa976 100644 --- a/functions/src/api-session-start.ts +++ b/functions/src/api-session-start.ts @@ -57,6 +57,7 @@ import { FLUSH_INTERVAL_MS, FLUSH_EVERY_N_TRIALS, MAX_DISCONNECTS, + streamingEnabled, } from "./staging.js"; export const apiSessionStart = onRequest({ cors: true }, async (req, res) => { @@ -65,6 +66,20 @@ export const apiSessionStart = onRequest({ cors: true }, async (req, res) => { return; } + // THE KILL SWITCH. Checked before anything else touches Firestore or RTDB: + // an incident where the staging tier itself is the problem (a runaway + // write pattern, an RTDB-side outage, a bug in this endpoint) needs a lever + // that does not depend on the thing that might be broken. Same response + // shape as an unprovisioned RTDB instance below, because the plugin's + // documented behaviour for it is already exactly right: fall back to + // submitting once at the end. Default (unset) is enabled, so no existing + // deployment changes behaviour; set STREAMING_ENABLED=false in + // functions/.env. to flip it off. + if (!streamingEnabled(process.env)) { + res.status(503).json(MESSAGES.SESSION_START_ERROR); + return; + } + // `filename` is optional and advisory: it is used ONLY to name a recovered // partial session, so that an abandoned run lands in the researcher's // storage under a name they recognise instead of an opaque session id. A diff --git a/functions/src/scheduled-staging-sweep.ts b/functions/src/scheduled-staging-sweep.ts index b56573c..3c24be7 100644 --- a/functions/src/scheduled-staging-sweep.ts +++ b/functions/src/scheduled-staging-sweep.ts @@ -54,6 +54,7 @@ import { listOldestOpenSessions, listOpenSessions, partialFilenameFor, + reconcileOpenSessionCounts, OpenSession, ABANDON_GRACE_MS, disconnectedSince, @@ -89,6 +90,14 @@ export interface SweepStats { * was lost. A value that stays non-zero run after run is a broken write path. */ mirrorFixed: number; + /** + * openSessionCounts entries this run had to correct against the RTDB ground + * truth -- the per-experiment concurrency cap's drift-tolerance backstop + * (see reconcileOpenSessionCounts in staging.ts). Should be zero for the + * same reason mirrorFixed should: each fix is evidence a decrement was + * missed somewhere upstream. + */ + countersFixed: number; } export const scheduledStagingSweep = onSchedule( @@ -125,16 +134,25 @@ export async function sweepAbandonedSessions( skippedLive: 0, errors: 0, mirrorFixed: 0, + countersFixed: 0, }; let lastError: string | null = null; let openSessionCount: number | null = null; let mirror = { created: 0, updated: 0, deleted: 0 }; + // The experiments this run's candidates belong to -- the scope for + // reconcileOpenSessionCounts below. Populated even for candidates that turn + // out to be live and get skipped: a live session still proves its + // experiment's counter is worth checking this run, on the same + // rotating-window logic CANDIDATES_PER_RUN already applies to the sessions + // themselves. + const candidateExperimentIds = new Set(); try { const candidates = (await listOldestOpenSessions(CANDIDATES_PER_RUN)).filter( (s) => !only || only.has(s.sessionId) ); stats.candidates = candidates.length; + for (const s of candidates) candidateExperimentIds.add(s.experimentId); const now = Date.now(); @@ -182,10 +200,12 @@ export async function sweepAbandonedSessions( // the sessions just recovered or discarded are already gone from both sides // and are not counted as fixes. A separate try: a mirror failure must not // be what stops abandoned sessions being recovered, and vice versa. + let openSessionsForCounters: OpenSession[] = []; try { const readAt = Date.now(); const open = await listOpenSessions(); openSessionCount = open.length; + openSessionsForCounters = open; const inScope = (only ? open.filter((s) => only.has(s.sessionId)) : open).slice(0, MAX_RECONCILE); const entries = await Promise.all( inScope.map(async (session) => ({ session, meta: await getSessionMeta(session.sessionId) })) @@ -217,6 +237,30 @@ export async function sweepAbandonedSessions( stats.errors++; } + // Correct the per-experiment concurrency counters (staging.ts's + // openSessionCounts) for the experiments touched this run. A separate try, + // for the same reason as the mirror above: a counter-reconciliation failure + // must not be what stops the mirror or the recovery pass, or vice versa. + // Scoped to candidateExperimentIds rather than every experiment with an open + // session, both to keep this bounded (the same rotating-window reasoning as + // CANDIDATES_PER_RUN) and to keep a test's scoped sweep run from correcting + // a counter belonging to a different, concurrently-running test. + try { + stats.countersFixed = await reconcileOpenSessionCounts( + candidateExperimentIds, + openSessionsForCounters + ); + if (stats.countersFixed > 0) { + console.warn( + `Open-session concurrency counters needed ${stats.countersFixed} fix(es).` + ); + } + } catch (e) { + lastError = e instanceof Error ? e.message : "Unknown error"; + console.error(`Open-session counter reconciliation failed: ${lastError}`); + stats.errors++; + } + await recordSweepHealth(stats, lastError, openSessionCount, mirror); if (stats.recovered > 0 || stats.discarded > 0) { @@ -382,6 +426,7 @@ async function recordSweepHealth( skippedLive: stats.skippedLive, errors: stats.errors, mirrorFixed: stats.mirrorFixed, + countersFixed: stats.countersFixed, mirror, lastError, }, diff --git a/functions/src/staging-assembly.ts b/functions/src/staging-assembly.ts index d95eccd..769203a 100644 --- a/functions/src/staging-assembly.ts +++ b/functions/src/staging-assembly.ts @@ -16,14 +16,43 @@ // See functions/src/staging.ts for the Realtime Database half, and // docs/streaming-ingest-design.md for the design. -// Mirrors the `newData.val().length <= 65536` cap in database.rules.json. +// THE SHARED CONSTANT MODULE. database.rules.json's `$seq` pattern and +// per-trial `.length` cap are HAND-COPIES of MAX_TRIALS_PER_SESSION and +// MAX_TRIAL_BYTES below -- RTDB rules cannot import a TypeScript module -- so +// __tests__/rules-constants.test.js parses database.rules.json and asserts +// the two literals still match these constants. If you change either constant +// here, update database.rules.json's `$seq` regex / `.length` cap in the same +// change, or that test fails. +// +// Mirrors the `newData.val().length <= 16384` cap in database.rules.json. // Exported so api-session-start.ts can tell the client the number rather than // letting the plugin carry its own copy that could drift out of sync with the // rule that actually enforces it. -export const MAX_TRIAL_BYTES = 65536; +// +// NOTE ON UNITS: RTDB's `.length` counts UTF-16 CODE UNITS, not bytes. A +// trial that is entirely BMP text (Latin, most punctuation, plain ASCII) has +// length == byte count in UTF-8; a trial full of multibyte content (CJK, +// emoji, non-Latin scripts) can be up to ~3x that many bytes for the same +// `.length`. The rules cap is therefore a UTF-16 ceiling, not a byte ceiling +// -- database.rules.json's own comments repeat this where the cap is +// enforced. MAX_ASSEMBLED_BYTES below, measured with Buffer.byteLength (real +// UTF-8 bytes) rather than `.length`, is the byte-accurate backstop: it is +// enforced server-side, in a paged read that never has to trust the rules' +// worst case to stay small. +export const MAX_TRIAL_BYTES = 16384; -// Mirrors the 4-digit `$seq` cap in database.rules.json. -export const MAX_TRIALS_PER_SESSION = 10000; +// Mirrors the 3-digit `$seq` cap in database.rules.json (max 1,000 trials). +// Lowered from a 4-digit / 10,000-trial cap: the structural worst case of +// MAX_TRIALS_PER_SESSION x MAX_TRIAL_BYTES is what a buggy or malicious +// client can force the sweep to read for one session id, and it has to stay +// well inside the sweep's Cloud Function memory allocation (256MiB, see +// scheduled-staging-sweep.ts) even accounting for the up-to-3x UTF-16-to-UTF8 +// multiplier above. 1,000 x 16 KiB is ~16 MiB of UTF-16 length, ~48 MiB at the +// worst-case byte multiplier -- comfortably inside 256MiB, where 10,000 x +// 64 KiB (640 MB structural, before any multiplier) was not. The paged read +// in assembleTrialsPaged (below) is the other half of this fix: it bounds +// peak memory to one page regardless of what the rules would otherwise allow. +export const MAX_TRIALS_PER_SESSION = 1000; // Client flush cadence, sent to the plugin by api-session-start.ts for the // same reason as MAX_TRIAL_BYTES: one source of truth, tunable without a @@ -62,12 +91,31 @@ export const MAX_DISCONNECTS = 20; // Firebase never got to run. export const SESSION_TTL_MS = 24 * 60 * 60 * 1000; -// Assembly ceiling. The rules permit 10,000 trials x 64 KiB = 640 MB in the -// worst case, which no honest session approaches and which would OOM the -// 256MiB sweep instantly. Assembly stops here and flags the result as +// The per-experiment concurrency cap on admitted (open) staging sessions. +// api-session-start.ts's own description of the instance ceiling puts it +// plainly: "[twenty instances] ... A single lecture-hall study would saturate +// it." Five hundred concurrent open sessions is generously above any lecture +// hall or online-panel study DataPipe hosts today, and it bounds what one +// experiment id -- public, by construction, in the experiment's own +// JavaScript -- can cost in staged RTDB storage by looping POST /api/session +// without ever completing a session. This is NOT a per-participant or +// per-IP limit (DataPipe's own code must never read or store participant IP +// addresses -- see pages/docs/privacy.js); it counts open sessions for the +// experiment as a whole, the same unit maxSessions already limits. +export const MAX_OPEN_SESSIONS_PER_EXPERIMENT = 500; + +// Assembly ceiling, in real bytes (Buffer.byteLength), not UTF-16 `.length`. +// The rules permit up to ~48 MB in the worst case (MAX_TRIALS_PER_SESSION x +// MAX_TRIAL_BYTES x the UTF-16-to-UTF-8 multiplier documented above), which no +// honest session approaches. Assembly stops here and flags the result as // truncated rather than dying: a truncated recovery of an abusive or runaway // session is a diagnosis, a crashed sweep is an outage that also lets every // other abandoned session pile up behind it. +// +// This is a BACKSTOP, not the primary defence. The primary defence is that +// assembleSession (staging.ts) reads the staging tree in PAGES and stops +// fetching as soon as this cap is crossed -- so an over-cap session is never +// pulled into memory in full, whatever the rules would otherwise allow. export const MAX_ASSEMBLED_BYTES = 24 * 1024 * 1024; // Cap on the client-supplied filename carried through a session. Generous for @@ -171,6 +219,28 @@ export function partialFilenameFor(session: OpenSession): string { return `${base || fallback}.partial.json`; } +/** + * Whether a staged value is a real trial, and its BYTE cost if so. + * + * Shared by assembleTrials and assembleTrialsPaged so the two never disagree + * on what counts as a trial or what it costs against MAX_ASSEMBLED_BYTES. + * Byte cost is measured with Buffer.byteLength, not `.length` -- `.length` on + * a JS string counts UTF-16 code units, and a trial full of multibyte content + * can be up to ~3x that many real bytes for the same `.length` (see + * MAX_TRIAL_BYTES above). Comparing `.length` against a byte ceiling is the + * bug the review flagged: it under-counts exactly the content it exists to + * bound. + */ +function evaluateTrial(raw: unknown): { ok: true; raw: string; byteLength: number } | { ok: false } { + if (typeof raw !== "string") return { ok: false }; + try { + JSON.parse(raw); + } catch { + return { ok: false }; + } + return { ok: true, raw, byteLength: Buffer.byteLength(raw, "utf8") }; +} + /** * The pure half of assembleSession: everything above, minus the read. * @@ -178,6 +248,13 @@ export function partialFilenameFor(session: OpenSession): string { * without an emulator -- the same separation compaction-gate.ts makes for * compaction, and for the same reason: the interesting logic should not need * the infrastructure to exercise it. + * + * Takes the WHOLE trials object already in memory -- unlike assembleSession, + * which pages the RTDB read so an over-cap session is never pulled into + * memory whole. This is still exactly right for a pure unit test (there is no + * RTDB to page against) and for anything that already has the full object + * (a fixture, a one-off script). See assembleTrialsPaged for the bounded-read + * version staging.ts actually calls. */ export function assembleTrials( trials: Record @@ -192,24 +269,18 @@ export function assembleTrials( let truncated = false; for (const key of keys) { - const raw = trials[key]; - if (typeof raw !== "string") { - skipped++; - continue; - } - try { - JSON.parse(raw); - } catch { + const outcome = evaluateTrial(trials[key]); + if (!outcome.ok) { skipped++; continue; } - const cost = raw.length + (parts.length > 0 ? 1 : 0); // + the comma + const cost = outcome.byteLength + (parts.length > 0 ? 1 : 0); // + the comma if (bytes + cost > MAX_ASSEMBLED_BYTES) { truncated = true; break; } bytes += cost; - parts.push(raw); + parts.push(outcome.raw); } // Gaps are counted against the highest sequence number actually present, not @@ -230,6 +301,126 @@ export function assembleTrials( }; } +/** One page of staged trials, in ascending key order. */ +export interface TrialPage { + /** [sequence key, raw value] pairs, in ascending numeric key order. */ + entries: Array<[string, unknown]>; + /** True when this is the last page -- fetchPage will not be called again. */ + done: boolean; +} + +export interface PagedAssemblyResult extends AssembledSession { + /** + * How many pages fetchPage() was called for. The point of paging: this + * stays small (bounded by MAX_ASSEMBLED_BYTES / page size) even for a + * session with MAX_TRIALS_PER_SESSION trials staged, because assembly stops + * calling fetchPage the moment the byte cap is crossed. + */ + pagesFetched: number; +} + +/** + * assembleTrials, but fed PAGES instead of the whole node. + * + * This is the fix for the review finding that the sweep read a whole session + * into memory before its size cap applied: `fetchPage` is called page by + * page, in ascending key order, and assembly STOPS CALLING IT the moment the + * accumulated byte size would cross MAX_ASSEMBLED_BYTES -- so an over-cap + * session is truncated without ever being pulled into memory in full. A + * within-page trial that would push the total over the cap is itself dropped + * (truncated, not included), exactly as assembleTrials drops the same trial + * when it hits the cap mid-object. + * + * ORDERING: identical guarantee to assembleTrials -- ascending numeric + * sequence order -- provided fetchPage itself returns each page in ascending + * key order and pages themselves arrive in ascending order (staging.ts pages + * via `orderByKey().startAfter(lastKey)`, which RTDB guarantees for + * integer-valued keys). This function does no re-sorting of its own: sorting + * would require every key in memory at once, which is the exact thing paging + * exists to avoid. + * + * GAPS, TRUNCATED: gaps are counted against the highest key SEEN SO FAR, which + * for a truncated session is necessarily a prefix of the truth -- a session + * that stops early because of the byte cap does not get its gap count + * re-derived from a full read it deliberately never performed. That is the + * same trade the byte cap itself makes: an exact diagnosis of an abusive + * session is not worth reading the whole thing into memory to produce. + */ +export async function assembleTrialsPaged( + fetchPage: (afterKey: string | null) => Promise +): Promise { + const parts: string[] = []; + let bytes = 2; // the enclosing brackets + let skipped = 0; + let truncated = false; + let presentCount = 0; + let highest = -1; + let pagesFetched = 0; + let afterKey: string | null = null; + + outer: while (true) { + const page = await fetchPage(afterKey); + pagesFetched++; + + for (const [key, raw] of page.entries) { + presentCount++; + const n = Number(key); + if (Number.isFinite(n) && n > highest) highest = n; + + const outcome = evaluateTrial(raw); + if (!outcome.ok) { + skipped++; + continue; + } + const cost = outcome.byteLength + (parts.length > 0 ? 1 : 0); // + the comma + if (bytes + cost > MAX_ASSEMBLED_BYTES) { + truncated = true; + break outer; + } + bytes += cost; + parts.push(outcome.raw); + } + + if (page.done || page.entries.length === 0) break; + afterKey = page.entries[page.entries.length - 1][0]; + } + + const gaps = + highest >= 0 ? Math.max(0, highest + 1 - presentCount) : 0; + + return { + data: `[${parts.join(",")}]`, + trialCount: parts.length, + skipped, + truncated, + gaps, + pagesFetched, + }; +} + +/** + * Whether the streaming-ingest endpoint should mint new sessions. + * + * A PURE function of an environment-shaped object, not a live read of + * process.env -- so it can be unit-tested without the functions emulator, + * which cannot be re-configured mid-suite (each test worker holds one + * process for the whole run). api-session-start.ts calls this with + * process.env; tests call it with a plain object. + * + * DEFAULT ENABLED. Unset, empty, or any value other than the exact string + * "false" leaves every existing deployment unaffected -- this is a kill + * switch to reach for during an incident, not a flag a deployment has to set + * to keep working. Setting STREAMING_ENABLED=false in + * functions/.env. (or .env.local for the emulator) stops + * POST /api/session from touching Firestore or RTDB at all; the plugin's + * documented fallback (pages/docs/api.js) means every experiment keeps + * working by submitting once at the end, exactly as it did before this + * feature existed. + */ +export function streamingEnabled(env: Record): boolean { + return env.STREAMING_ENABLED !== "false"; +} + /** [slot number, timestamp] pairs from either shape RTDB may return. */ function slotEntries(slots: SlotMap | undefined): Array<[number, number]> { if (!slots) return []; diff --git a/functions/src/staging.ts b/functions/src/staging.ts index ef06ed6..ce32eb4 100644 --- a/functions/src/staging.ts +++ b/functions/src/staging.ts @@ -51,10 +51,12 @@ import { mirrorStart, removeLiveSession, connectionState } from "./live-sessions import { AssembledSession, MAX_FILENAME_LENGTH, + MAX_OPEN_SESSIONS_PER_EXPERIMENT, OpenSession, SessionMeta, SESSION_TTL_MS, - assembleTrials, + TrialPage, + assembleTrialsPaged, } from "./staging-assembly.js"; // Re-exported so callers have ONE import for the staging tier and do not have @@ -209,6 +211,125 @@ export function resetStagingHandleForTests(): void { cachedDb = null; } +// --------------------------------------------------------------------------- +// Per-experiment concurrency cap +// --------------------------------------------------------------------------- +// +// openSessionCounts/{experimentId} : number of currently-open sessions. +// +// A REVIEW FINDING THIS ANSWERS: one anonymous POST /api/session mints one +// session id, and nothing before this counter existed stopped that call being +// looped -- an experiment id is public (it ships in the experiment's own +// JavaScript), so the cost of an unbounded loop was 640MB-per-session-id of +// RTDB storage, times however many times it was called. The counter below +// bounds it to MAX_OPEN_SESSIONS_PER_EXPERIMENT sessions per experiment, +// however many times the endpoint is called. +// +// Deliberately NOT a per-IP counter: DataPipe's own code must never read or +// store a participant's IP address (pages/docs/privacy.js). This counts +// SESSIONS for an EXPERIMENT, the same unit maxSessions already limits, and +// carries no information about who is opening them. + +const OPEN_SESSION_COUNTS_PATH = "openSessionCounts"; + +/** + * Reserve one of MAX_OPEN_SESSIONS_PER_EXPERIMENT concurrent staging slots. + * Returns false, without writing anything else, once the experiment is at + * its cap. + * + * A TRANSACTION, not a read-then-write: two POST /api/session calls for the + * same experiment arriving together must not both read "499" and both + * proceed. RTDB retries a transaction against the server's current value on + * a conflicting write, which a plain get()-then-set() cannot do. + * + * SELF-HEALING ON NEGATIVE OR MALFORMED DRIFT: a stored value that is + * missing, not a number, or negative -- which releaseOpenSessionSlot's own + * floor should make impossible, but a hand edit or a bug predating this code + * could still produce -- is treated as zero rather than compounding the error + * into a cap that can never again be satisfied. This is the "recomputing when + * the counter goes negative" half of tolerating drift; reconcileOpenSessionCounts + * below is the other half, for drift that is positive (a missed decrement) + * rather than negative. + */ +export async function tryAdmitSession(experimentId: string): Promise { + const ref = rtdb().ref(`${OPEN_SESSION_COUNTS_PATH}/${experimentId}`); + const result = await ref.transaction((current: unknown) => { + const count = typeof current === "number" && current > 0 ? current : 0; + if (count >= MAX_OPEN_SESSIONS_PER_EXPERIMENT) return; // undefined aborts the transaction, writing nothing + return count + 1; + }); + return result.committed; +} + +/** + * Release a concurrency slot for an experiment. Floored at zero: a decrement + * that ever ran without (or twice for) a matching increment must not push the + * counter negative, which would then let in extra sessions until the count + * climbed back to zero on its own. + */ +export async function releaseOpenSessionSlot(experimentId: string): Promise { + const ref = rtdb().ref(`${OPEN_SESSION_COUNTS_PATH}/${experimentId}`); + await ref.transaction((current: unknown) => { + const count = typeof current === "number" ? current : 0; + return Math.max(0, count - 1); + }); +} + +/** + * Correct openSessionCounts for a scoped set of experiments against the RTDB + * ground truth (POSITIVE drift: a missed decrement -- see + * tryAdmitSession's doc for the negative-drift half of this). + * + * Scoped to `experimentIds`, not every counter that has ever existed: the + * sweep calls this once per run with the experiments its own candidates + * belong to, which is the same rotating-window approach + * CANDIDATES_PER_RUN already takes for the sessions themselves -- a + * stuck-too-high counter is corrected within a few runs rather than this + * needing an unbounded read of every experiment that has ever streamed. It + * also keeps a test that scopes a sweep run to its own experiment id (the + * `only` seam) from correcting -- or racing against -- a counter that + * belongs to a different, concurrently-running test. + * + * `openSessions` is the caller's already-fetched list of every currently-open + * session (listOpenSessions()), so this performs no additional read of the + * staging tier itself -- only of the small openSessionCounts table, and only + * for the experiment ids in scope. + */ +export async function reconcileOpenSessionCounts( + experimentIds: Iterable, + openSessions: OpenSession[] +): Promise { + const ids = [...new Set(experimentIds)]; + if (ids.length === 0) return 0; + + const trueCounts = new Map(ids.map((id) => [id, 0])); + for (const session of openSessions) { + if (trueCounts.has(session.experimentId)) { + trueCounts.set(session.experimentId, (trueCounts.get(session.experimentId) as number) + 1); + } + } + + const stored = await Promise.all( + ids.map((id) => rtdb().ref(`${OPEN_SESSION_COUNTS_PATH}/${id}`).get()) + ); + + const updates: Record = {}; + ids.forEach((id, i) => { + const storedValue = stored[i].exists() ? stored[i].val() : 0; + const storedCount = typeof storedValue === "number" ? storedValue : 0; + const truth = trueCounts.get(id) as number; + if (storedCount !== truth) { + updates[id] = truth === 0 ? null : truth; + } + }); + + const fixed = Object.keys(updates).length; + if (fixed > 0) { + await rtdb().ref(OPEN_SESSION_COUNTS_PATH).update(updates); + } + return fixed; +} + // --------------------------------------------------------------------------- // Session lifecycle // --------------------------------------------------------------------------- @@ -220,36 +341,59 @@ export function resetStagingHandleForTests(): void { * experiment is open. This function does not re-check, because it cannot: RTDB * has no view of Firestore, and that asymmetry is the whole reason session ids * are minted server-side. + * + * Throws (rather than returning a sentinel) when the experiment is at its + * concurrency cap. api-session-start.ts's existing try/catch around this call + * turns that into the same 503 SESSION_START_ERROR shape it already returns + * for an unprovisioned RTDB instance -- the plugin's documented fallback is + * to submit once at the end, which is the right behaviour here too. */ export async function openSession( experimentId: string, filename: string | undefined, owner: string ): Promise { - const sessionId = generateSessionId(); - const now = Date.now(); - const record: Record = { - experimentId, - owner, - startedAt: now, - expiresAt: now + SESSION_TTL_MS, - }; - // Omitted rather than written as undefined: RTDB rejects undefined values - // the same way Firestore does, and an absent filename is a normal state. - if (filename) record.filename = filename.slice(0, MAX_FILENAME_LENGTH); - await rtdb().ref(`openSessions/${sessionId}`).set(record); - - // The researcher's live dashboard copy (live-sessions.ts). Awaited, so it is - // written before the participant's page gets its response, but best-effort: - // mirrorStart swallows its own failure, and the sweep backfills a miss. - await mirrorStart(sessionId, { - experimentID: experimentId, - owner, - startedAt: now, - expiresAt: now + SESSION_TTL_MS, - ...connectionState({}), - }); - return sessionId; + const admitted = await tryAdmitSession(experimentId); + if (!admitted) { + throw new Error( + `Experiment ${experimentId} already has ${MAX_OPEN_SESSIONS_PER_EXPERIMENT} sessions open; ` + + "refusing another until one completes or is recovered." + ); + } + + try { + const sessionId = generateSessionId(); + const now = Date.now(); + const record: Record = { + experimentId, + owner, + startedAt: now, + expiresAt: now + SESSION_TTL_MS, + }; + // Omitted rather than written as undefined: RTDB rejects undefined values + // the same way Firestore does, and an absent filename is a normal state. + if (filename) record.filename = filename.slice(0, MAX_FILENAME_LENGTH); + await rtdb().ref(`openSessions/${sessionId}`).set(record); + + // The researcher's live dashboard copy (live-sessions.ts). Awaited, so it is + // written before the participant's page gets its response, but best-effort: + // mirrorStart swallows its own failure, and the sweep backfills a miss. + await mirrorStart(sessionId, { + experimentID: experimentId, + owner, + startedAt: now, + expiresAt: now + SESSION_TTL_MS, + ...connectionState({}), + }); + return sessionId; + } catch (e) { + // The slot was reserved above but nothing that would need its own + // teardown got written (or mirrorStart already swallowed its failure), so + // releasing it here is the only cleanup needed to avoid leaking a + // permanently-reserved slot on a failed admission. + await releaseOpenSessionSlot(experimentId); + throw e; + } } /** The capability record for a session id, or null if it is not open. */ @@ -323,9 +467,27 @@ export async function listOpenSessions(): Promise { return rows; } +// Trials fetched per RTDB round trip during assembly. Small enough that one +// page (ASSEMBLY_PAGE_SIZE x MAX_TRIAL_BYTES, worst case) is a fraction of +// MAX_ASSEMBLED_BYTES, so assembly can stop mid-page without having pulled +// anything close to a full session into memory first -- which is the whole +// point of paging the read at all (see the review finding at the top of this +// file: assembleSession used to `.get()` the entire trials node before its +// size cap applied). +const ASSEMBLY_PAGE_SIZE = 200; + /** * Read and assemble a session's staged trials. * + * READ IN PAGES, not one `.get()` of the whole node. `orderByKey()` gives + * RTDB's native ordering for integer-valued keys, which is numeric ascending + * -- the same order assembleTrials produces by sorting -- so + * `.startAfter(lastKey)` resumes exactly where the previous page left off + * with no re-sorting required. assembleTrialsPaged stops calling this fetcher + * the moment the accumulated byte size would cross MAX_ASSEMBLED_BYTES, so an + * over-cap session is truncated without this function ever holding more than + * one page in memory. + * * GAPS ARE TOLERATED, NOT REJECTED (design doc risk #5). A missing sequence * number means one flush never landed -- a dropped request, a tab closed * mid-write. The remaining trials are still the participant's real data, and @@ -341,11 +503,28 @@ export async function listOpenSessions(): Promise { export async function assembleSession( sessionId: string ): Promise { - const snap = await rtdb().ref(`staging/${sessionId}/trials`).get(); - if (!snap.exists()) { - return { data: "[]", trialCount: 0, skipped: 0, truncated: false, gaps: 0 }; - } - return assembleTrials(snap.val() as Record); + const trialsRef = rtdb().ref(`staging/${sessionId}/trials`); + + const fetchPage = async (afterKey: string | null): Promise => { + const query = + afterKey === null + ? trialsRef.orderByKey().limitToFirst(ASSEMBLY_PAGE_SIZE) + : trialsRef.orderByKey().startAfter(afterKey).limitToFirst(ASSEMBLY_PAGE_SIZE); + const snap = await query.get(); + if (!snap.exists()) return { entries: [], done: true }; + const entries: Array<[string, unknown]> = []; + snap.forEach((child) => { + entries.push([child.key as string, child.val()]); + return false; // keep iterating (forEach cancels on `true`) + }); + return { entries, done: entries.length < ASSEMBLY_PAGE_SIZE }; + }; + + // pagesFetched is a diagnostic for the caller's own tests, not part of the + // durable result -- discard it here so AssembledSession stays the one shape + // every caller (the sweep, its tests) already knows. + const { pagesFetched: _pagesFetched, ...assembled } = await assembleTrialsPaged(fetchPage); + return assembled; } /** @@ -377,6 +556,24 @@ export async function discardSession(sessionId: string): Promise { ); return; } + + // Read BEFORE the delete below removes it: releasing this session's + // concurrency slot needs to know which experiment it belonged to, and + // openSessions/{sessionId} is the only place that is recorded. Best-effort + // like everything else here -- a failed read just skips the release, and + // reconcileOpenSessionCounts corrects the resulting drift on the sweep's + // next run rather than this turning into a failed discard. + let experimentId: string | undefined; + try { + const snap = await rtdb().ref(`openSessions/${sessionId}/experimentId`).get(); + if (snap.exists()) experimentId = snap.val() as string; + } catch (e) { + const detail = e instanceof Error ? e.message : "Unknown error"; + console.error( + `Failed to read experimentId while discarding staging session ${sessionId}: ${detail}` + ); + } + try { await rtdb() .ref() @@ -388,6 +585,9 @@ export async function discardSession(sessionId: string): Promise { const detail = e instanceof Error ? e.message : "Unknown error"; console.error(`Failed to discard staging session ${sessionId}: ${detail}`); } + + if (experimentId) await releaseOpenSessionSlot(experimentId); + // And the researcher's dashboard row. Every way a session ends -- clean // completion, a gate refusing it, the sweep recovering or discarding it -- // comes through here, which is why this is the one place that removes it. diff --git a/pages/docs/api.js b/pages/docs/api.js index 155f596..a00734f 100644 --- a/pages/docs/api.js +++ b/pages/docs/api.js @@ -182,8 +182,10 @@ export default function ApiReferencePage() { not consume one from that limit — the count is still taken when a submission completes. A 503 with{" "} SESSION_START_ERROR means incremental upload is - unavailable and the experiment should simply submit at the end, as it - would otherwise. + unavailable — because the service is unreachable, because an + experiment already has an unusually large number of sessions open at + once, or because it has been switched off entirely — and the + experiment should simply submit at the end, as it would otherwise. @@ -193,8 +195,8 @@ export default function ApiReferencePage() { {`{ "sessionId": "8fKq2mXpR7vNwLzB4cTy1dHs", "databaseURL": "https://-default-rtdb.firebaseio.com", - "maxTrialBytes": 65536, - "maxTrials": 10000, + "maxTrialBytes": 16384, + "maxTrials": 1000, "flushIntervalMs": 10000, "flushEveryNTrials": 10 }`} From fbad44602dd46c468804b5b3e3c1733f66db04dd Mon Sep 17 00:00:00 2001 From: Josh de Leeuw Date: Sun, 13 Sep 2026 17:16:44 -0400 Subject: [PATCH 2/6] Harden the abandoned-session sweep against starvation, mis-reported discards, and filename collisions Six review findings against the streaming-ingest staging tier (docs/streaming-ingest-design.md), all addressed in this PR: 1. Candidate starvation. scheduled-staging-sweep.ts read one page of 30 oldest-open-session candidates per run and never advanced past it, so a page's worth of long-lived or zombie sessions at the head of the queue (ordered by expiresAt) could block recovery of everything abandoned behind them for up to the 24-hour TTL. listOldestOpenSessions (staging.ts) now takes an optional cursor and the sweep pages through candidates -- startAfter(expiresAt, sessionId) -- until MAX_SESSIONS_PER_RUN sessions are actually recovered or discarded, the table is exhausted, or a hard MAX_PAGES_PER_RUN (20) ceiling is hit. Pages fetched are recorded in SweepStats and systemStatus/staging as `pages`. 2. Discard outcome accuracy. discardSession (staging.ts) swallowed its own RTDB error and returned void, so the sweep reported "recovered" regardless of whether the staging node was actually removed. If the discard fails after a queue entry was written, the next run reassembles and re-queues the same session under the same (now deterministic) filename; if that first entry had already completed, a duplicate partial reaches the provider. discardSession now returns a boolean (still never throws); the sweep only counts/logs "recovered" or "discarded" when the removal actually succeeded, otherwise counting it as an error so a persistently failing discard path is visible. recoverSession also checks for an existing uploadQueue document with the same deduplication key (queue-upload.ts's new queueDocIdFor helper) that is already "completed" before queueing, and skips the re-queue if so. 3. Concurrency limiter. The live-sessions reconciliation pass ran an unbounded Promise.all over up to 500 getSessionMeta calls. Replaced with mapWithConcurrency (new functions/src/concurrency-limit.ts), a small local worker-pool helper with no new dependency, at a concurrency of 20. 4. Partial filename collisions. partialFilenameFor (staging-assembly.ts) used only the client-supplied filename, so two abandoned sessions named e.g. "data.csv" produced the identical experimentID:filename deduplication key and the second recovery silently overwrote the first's payload in Cloud Storage. A client-supplied name now always gets an 8-hex-char suffix (a sha256 of the session id); the id-only fallback for sessions with no filename is unchanged, since it is already unique. Updated the filename description in pages/docs/experiments/sending-data.js and docs/streaming-ingest-design.md. 5. Salvage instead of discard on data-content refusals. api-data.ts discarded the staged copy on all eight refusal branches, including INVALID_DATA and both duplicate-filename refusals (the collision cache's "duplicate" verdict and the provider's own NAME_CONFLICT dual-run backstop) -- refusals about THIS SUBMISSION, not about the experiment's willingness to accept data. Those three branches now leave staging alone; the sweep's second door (finalized/active re-check) still covers a since-closed experiment. The four experiment-state gates (finalized, inactive, session-cap, unknown experiment) are unchanged. 6. Confirmed (read, did not need to change) that a missing or unreachable RTDB is caught at every call site in the sweep and only ever surfaces as a recorded systemStatus/staging error, never an uncaught throw. Added a cheap regression test pointing STAGING_DATABASE_URL at a refused local port. Tests added, all in functions/src/__tests__/: - concurrency-limit.test.js (new, pure, no emulator): mapWithConcurrency ordering, concurrency ceiling, exhaustiveness, rejection propagation, edge cases. - staging-assembly.test.js: partialFilenameFor cases updated for the new hash suffix, plus a same-filename/different-session collision case. - staging-emulator.test.js: candidate-paging-past-a-wall-of-live-sessions, same-filename distinct queue entries, discard-failure not counted as recovered, no re-queue of an already-completed session after a discard failure, unreachable-database survival, and INVALID_DATA leaving staging in place (existing finalized-still-discards test serves as the control). Existing p07/p09 filename assertions updated for the new hash suffix. `cd functions && npm run build` and `npm run lint` (repo root) are both clean. Emulator suites were not run here per the parent's instructions; see the report for the exact files to run. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Q8xf16ov88M3KojPVfQwJT --- docs/streaming-ingest-design.md | 10 +- .../src/__tests__/concurrency-limit.test.js | 80 ++++++ .../src/__tests__/staging-assembly.test.js | 32 ++- .../src/__tests__/staging-emulator.test.js | 243 +++++++++++++++++- functions/src/api-data.ts | 62 ++++- functions/src/concurrency-limit.ts | 47 ++++ functions/src/queue-upload.ts | 17 +- functions/src/scheduled-staging-sweep.ts | 234 ++++++++++++----- functions/src/staging-assembly.ts | 31 ++- functions/src/staging.ts | 75 ++++-- pages/docs/experiments/sending-data.js | 8 +- 11 files changed, 718 insertions(+), 121 deletions(-) create mode 100644 functions/src/__tests__/concurrency-limit.test.js create mode 100644 functions/src/concurrency-limit.ts diff --git a/docs/streaming-ingest-design.md b/docs/streaming-ingest-design.md index 59ad5d9..19a86b7 100644 --- a/docs/streaming-ingest-design.md +++ b/docs/streaming-ingest-design.md @@ -17,9 +17,11 @@ sections they affect are annotated inline. In summary: dataset, because the browser still has it; `sessionId` only names the staged copy to discard. This removes the only line item in the cost table below and keeps validation, CSV support and the metadata pipeline byte-identical. -3. **Partial sessions** upload as `.partial.json`, do not increment - `sessions`, and do not arm the upload-failure notifier. This answers open - question 2. +3. **Partial sessions** upload as `-.partial.json` (the hash is a + short digest of the session id, added so two sessions that happen to share + a client-supplied filename cannot collide on the same recovered file), do + not increment `sessions`, and do not arm the upload-failure notifier. This + answers open question 2. 4. **Staged trials are not encrypted at the application layer**, contrary to this document's "safe default" in open question 4. The writer is the participant's browser: it has no key, and a key shipped in a plugin bundle @@ -315,7 +317,7 @@ option here ships as a coordinated pair of releases. badly with experiments served from arbitrary hosts. **Unresolved — this needs an answer before the spike, not after.** -2. **ANSWERED (decision 3): `.partial.json`, uncounted, no notification.** +2. **ANSWERED (decision 3): `-.partial.json`, uncounted, no notification.** The sweep re-checks `finalized` and `active` before queueing, because a researcher can seal an experiment between staging and recovery, and a file landing outside a merged archive is exactly what `docs/finalization-spec.md` diff --git a/functions/src/__tests__/concurrency-limit.test.js b/functions/src/__tests__/concurrency-limit.test.js new file mode 100644 index 0000000..1351893 --- /dev/null +++ b/functions/src/__tests__/concurrency-limit.test.js @@ -0,0 +1,80 @@ +/** + * @jest-environment node + * + * mapWithConcurrency (functions/src/concurrency-limit.ts): the bounded- + * concurrency helper that replaced the unbounded Promise.all in + * scheduled-staging-sweep.ts's live-sessions reconciliation pass. Pure and + * infrastructure-free, so no emulator is needed. + */ + +import { mapWithConcurrency } from '../../lib/concurrency-limit.js'; + +describe('mapWithConcurrency', () => { + it('returns results in input order regardless of completion order', async () => { + const delays = [30, 10, 20, 0, 15]; + + const results = await mapWithConcurrency(delays, 3, async (delay, index) => { + await new Promise((resolve) => setTimeout(resolve, delay)); + return index; + }); + + expect(results).toEqual([0, 1, 2, 3, 4]); + }); + + it('never runs more than `concurrency` calls at once', async () => { + let inFlight = 0; + let maxInFlight = 0; + const items = Array.from({ length: 20 }, (_, i) => i); + + await mapWithConcurrency(items, 4, async (item) => { + inFlight++; + maxInFlight = Math.max(maxInFlight, inFlight); + await new Promise((resolve) => setTimeout(resolve, 5)); + inFlight--; + return item * 2; + }); + + expect(maxInFlight).toBeLessThanOrEqual(4); + }); + + it('runs every item exactly once', async () => { + const seen = []; + const items = Array.from({ length: 37 }, (_, i) => i); + + await mapWithConcurrency(items, 5, async (item) => { + seen.push(item); + }); + + expect(seen.slice().sort((a, b) => a - b)).toEqual(items); + }); + + it('propagates a rejection rather than swallowing it', async () => { + await expect( + mapWithConcurrency([1, 2, 3], 2, async (item) => { + if (item === 2) throw new Error('boom'); + return item; + }) + ).rejects.toThrow('boom'); + }); + + it('handles concurrency higher than the item count', async () => { + const results = await mapWithConcurrency([1, 2], 20, async (item) => item + 1); + expect(results).toEqual([2, 3]); + }); + + it('handles an empty list', async () => { + const results = await mapWithConcurrency([], 5, async (item) => item); + expect(results).toEqual([]); + }); + + it('handles a concurrency of 1 by running strictly sequentially', async () => { + const order = []; + await mapWithConcurrency([1, 2, 3], 1, async (item) => { + order.push(`start-${item}`); + await new Promise((resolve) => setTimeout(resolve, 5)); + order.push(`end-${item}`); + }); + + expect(order).toEqual(['start-1', 'end-1', 'start-2', 'end-2', 'start-3', 'end-3']); + }); +}); diff --git a/functions/src/__tests__/staging-assembly.test.js b/functions/src/__tests__/staging-assembly.test.js index 8be5e1c..e1a7e51 100644 --- a/functions/src/__tests__/staging-assembly.test.js +++ b/functions/src/__tests__/staging-assembly.test.js @@ -7,6 +7,7 @@ * to exercise. */ +import { createHash } from 'crypto'; import { assembleTrials, partialFilenameFor, @@ -14,6 +15,11 @@ import { MAX_ASSEMBLED_BYTES, } from '../../lib/staging-assembly.js'; +// Computed here rather than imported, the same way staging-emulator.test.js +// recomputes the live-sessions mirror id: the suite should check the actual +// hash a session id produces, not the module's opinion of it. +const shortHash = (sessionId) => createHash('sha256').update(sessionId).digest('hex').slice(0, 8); + /** The shape RTDB hands back: a map of sequence key -> trial JSON string. */ function staged(...jsonStrings) { return Object.fromEntries(jsonStrings.map((s, i) => [String(i), s])); @@ -132,20 +138,36 @@ describe('partialFilenameFor', () => { // submitted, so a recovered fragment of a CSV study is a .json file and // has to say so. expect(partialFilenameFor(session({ filename: 'subject42.csv' }))).toBe( - 'subject42.partial.json' + `subject42-${shortHash('abc123')}.partial.json` ); }); it('replaces an existing json extension rather than doubling it', () => { expect(partialFilenameFor(session({ filename: 'subject42.json' }))).toBe( - 'subject42.partial.json' + `subject42-${shortHash('abc123')}.partial.json` ); }); it('falls back to the session id when no filename was captured', () => { + // No hash appended here: the fallback name already contains the full + // session id and is unique on its own. expect(partialFilenameFor(session())).toBe('session-abc123.partial.json'); }); + it('gives two sessions with the same client-supplied filename distinct names', () => { + // The bug this exists to prevent: without a per-session suffix, two + // abandoned sessions named "data.csv" would collide on the same + // `experimentID:filename` deduplication key in queue-upload.ts and the + // second recovery would silently overwrite the first's payload in Cloud + // Storage. + const a = partialFilenameFor(session({ sessionId: 'session-aaa', filename: 'data.csv' })); + const b = partialFilenameFor(session({ sessionId: 'session-bbb', filename: 'data.csv' })); + + expect(a).not.toBe(b); + expect(a).toBe(`data-${shortHash('session-aaa')}.partial.json`); + expect(b).toBe(`data-${shortHash('session-bbb')}.partial.json`); + }); + it('neutralises path separators in a client-supplied name', () => { // The name is CLIENT-SUPPLIED and becomes a path in a researcher's Drive, // OSF or Zenodo container. Nothing downstream re-checks it. @@ -158,7 +180,7 @@ describe('partialFilenameFor', () => { expect(result.endsWith('.partial.json')).toBe(true); expect(partialFilenameFor(session({ filename: 'a/b\\c.csv' }))).toBe( - 'a_b_c.partial.json' + `a_b_c-${shortHash('abc123')}.partial.json` ); }); @@ -166,11 +188,11 @@ describe('partialFilenameFor', () => { // `\.[^.]*$` looks correct and is not: on a name with embedded dots but no // extension it eats the last segment. expect(partialFilenameFor(session({ filename: 'etc.d/passwd' }))).toBe( - 'etc.d_passwd.partial.json' + `etc.d_passwd-${shortHash('abc123')}.partial.json` ); // A version-style name keeps the version. expect(partialFilenameFor(session({ filename: 'data.2026.csv' }))).toBe( - 'data.2026.partial.json' + `data.2026-${shortHash('abc123')}.partial.json` ); }); diff --git a/functions/src/__tests__/staging-emulator.test.js b/functions/src/__tests__/staging-emulator.test.js index 948abc7..980ddbb 100644 --- a/functions/src/__tests__/staging-emulator.test.js +++ b/functions/src/__tests__/staging-emulator.test.js @@ -39,6 +39,7 @@ const { sweepAbandonedSessions, ABANDON_GRACE_MS, } = require("../../lib/scheduled-staging-sweep.js"); +const { generateSessionId, resetStagingHandleForTests } = require("../../lib/staging.js"); const FUNCTIONS_HOST = process.env.FUNCTIONS_EMULATOR_HOST || "localhost:5001"; const PROJECT_ID = "datapipe-test"; @@ -324,6 +325,37 @@ describe("completion", () => { expect(status).toBe(400); expect(body.error).toBe("EXPERIMENT_FINALIZED"); }); + + it("leaves the staged copy in place when the submission itself is refused as invalid", async () => { + // INVALID_DATA is a refusal about THIS SUBMISSION -- this string failed + // validation -- not about whether the experiment accepts data at all. The + // participant's staged trials are unaffected by that verdict, and + // discarding them would destroy the one recoverable copy in exactly the + // case the staging tier exists for: a browser that is never coming back + // to retry with a better-formed payload. Contrast with the finalized-gate + // test above, which still discards. + const experimentID = await makeExperiment({ + useValidation: true, + allowJSON: true, + allowCSV: false, + requiredFields: [], + }); + const { body } = await startSession({ experimentID }); + await stageTrials(body.sessionId, 3); + + const { status, body: response } = await saveData({ + experimentID, + filename: "p01.json", + data: "this is not valid json", + sessionId: body.sessionId, + }); + + expect(status).toBe(400); + expect(response.error).toBe("INVALID_DATA"); + expect((await rtdb.ref(`staging/${body.sessionId}/trials`).get()).numChildren()).toBe(3); + expect((await rtdb.ref(`openSessions/${body.sessionId}`).get()).exists()).toBe(true); + }); + }); describe("the abandonment sweep", () => { @@ -340,8 +372,11 @@ describe("the abandonment sweep", () => { const entries = await queueEntriesFor(experimentID); expect(entries).toHaveLength(1); const [entry] = entries; - // Marked in the name, so a researcher can tell a fragment from a session. - expect(entry.filename).toBe("p07.partial.json"); + // Marked in the name, so a researcher can tell a fragment from a session + // -- and suffixed with a hash of the session id, so two sessions named + // "p07.csv" cannot collide on the same recovered file. + const suffix = createHash("sha256").update(body.sessionId).digest("hex").slice(0, 8); + expect(entry.filename).toBe(`p07-${suffix}.partial.json`); expect(entry.partial).toBe(true); expect(entry.status).toBe("pending"); expect(entry.failureReason).toContain("4 trials"); @@ -409,7 +444,8 @@ describe("the abandonment sweep", () => { const stats = await sweepAbandonedSessions(new Set([body.sessionId])); expect(stats.recovered).toBe(1); - expect((await queueEntriesFor(experimentID))[0].filename).toBe("p09.partial.json"); + const suffix = createHash("sha256").update(body.sessionId).digest("hex").slice(0, 8); + expect((await queueEntriesFor(experimentID))[0].filename).toBe(`p09-${suffix}.partial.json`); }); it("discards rather than uploads when the experiment was finalized meanwhile", async () => { @@ -484,8 +520,209 @@ describe("the abandonment sweep", () => { expect(status.exists).toBe(true); expect(status.data().lastRunAt.toMillis()).toBeGreaterThan(Date.now() - 60000); expect(status.data().openSessionCount).toEqual(expect.any(Number)); + expect(status.data().pages).toEqual(expect.any(Number)); expect(status.data().lastError).toBeNull(); }); + + it("pages past a wall of live sessions to recover an abandoned one behind them", async () => { + // Regression for candidate starvation: sessions skipped as live used to + // consume candidate slots with no cursor advancing, so a page's worth of + // long-lived or zombie sessions at the head of the queue could block + // recovery of everything behind them for up to 24 hours. + const experimentID = await makeExperiment(); + // Far enough out that these fixtures sort ahead of every other open + // session in the shared emulator (which default to ~24h out), but the + // exact value doesn't matter -- only the ordering between these fixtures + // does. + const base = Date.now() + 5 * 60000; + + // More than one page (CANDIDATES_PER_PAGE = MAX_SESSIONS_PER_RUN * 3 = 30) + // of live sessions, each with a smaller `expiresAt` than the target below + // -- so they sort first and the target lands on a later page. Written + // directly to RTDB (bypassing the session-start endpoint) so the fixture + // stays fast: only the fields the sweep actually reads matter here. + const wallSize = 35; + for (let i = 0; i < wallSize; i++) { + const sessionId = generateSessionId(); + await rtdb.ref(`openSessions/${sessionId}`).set({ + experimentId: experimentID, + owner: "staging-testuser", + startedAt: Date.now(), + expiresAt: base + i, + }); + } + + const { body } = await startSession({ experimentID, filename: "behind-the-wall.csv" }); + // Push this session's expiry behind the whole wall above. + await rtdb.ref(`openSessions/${body.sessionId}/expiresAt`).set(base + wallSize + 1000); + await stageTrials(body.sessionId, 3); + await markAbandoned(body.sessionId); + + const stats = await sweepAbandonedSessions(new Set([body.sessionId])); + + expect(stats.recovered).toBe(1); + // The whole point: more than one page had to be fetched to reach it. + expect(stats.pages).toBeGreaterThan(1); + + const suffix = createHash("sha256").update(body.sessionId).digest("hex").slice(0, 8); + const entries = await queueEntriesFor(experimentID); + expect(entries.map((e) => e.filename)).toContain(`behind-the-wall-${suffix}.partial.json`); + }); + + it("keeps two abandoned sessions that share a client-supplied filename as distinct queue entries", async () => { + // Without a session-id suffix on the recovered filename, both sessions + // would compute the SAME `experimentID:filename` deduplication key in + // queue-upload.ts, and the second recovery would silently overwrite the + // first session's payload in Cloud Storage before either one is + // delivered. + const experimentID = await makeExperiment(); + + const first = await startSession({ experimentID, filename: "data.csv" }); + await stageTrials(first.body.sessionId, 2); + await markAbandoned(first.body.sessionId); + + const second = await startSession({ experimentID, filename: "data.csv" }); + await stageTrials(second.body.sessionId, 5); + await markAbandoned(second.body.sessionId); + + const stats = await sweepAbandonedSessions( + new Set([first.body.sessionId, second.body.sessionId]) + ); + + expect(stats.recovered).toBe(2); + const entries = await queueEntriesFor(experimentID); + expect(entries).toHaveLength(2); + const filenames = entries.map((e) => e.filename); + // Distinct docs, distinct filenames -- not one overwriting the other. + expect(new Set(filenames).size).toBe(2); + filenames.forEach((f) => expect(f).toMatch(/^data-[0-9a-f]{8}\.partial\.json$/)); + }); + + it("does not report a session as recovered when its discard fails", async () => { + // Regression: discardSession used to swallow its own RTDB error and the + // sweep reported "recovered" regardless of whether the staging node was + // actually removed. If the failure happens AFTER the queue entry is + // written, the session is still sitting in openSessions afterwards and + // the next run reassembles and re-queues it -- a duplicate delivery if + // the first entry has already completed by then (covered separately + // below). + // + // Faked here via the seam discardSession already has: an id that fails + // isValidSessionId is refused before it ever touches RTDB, returning + // false without throwing -- exactly the "discard failed, non-throwing" + // contract being tested. openSession() never mints an id shaped like + // this; writing one directly is what stands in for "the RTDB write + // failed" here. + const experimentID = await makeExperiment(); + const badId = "not-a-real-session-id"; + await rtdb.ref(`openSessions/${badId}`).set({ + experimentId: experimentID, + owner: "staging-testuser", + startedAt: Date.now(), + expiresAt: Date.now() + 60000, + filename: "p01.csv", + }); + await stageTrials(badId, 3); + await markAbandoned(badId); + + try { + const stats = await sweepAbandonedSessions(new Set([badId])); + + expect(stats.recovered).toBe(0); + expect(stats.discarded).toBe(0); + expect(stats.errors).toBeGreaterThanOrEqual(1); + + // The data WAS queued -- this is about the REPORT, not about queueing + // having failed too. + const entries = await queueEntriesFor(experimentID); + expect(entries).toHaveLength(1); + expect(entries[0].status).toBe("pending"); + + // And the staging node genuinely still exists, exactly as a false + // discardOk promised -- next run will see it again. + expect((await rtdb.ref(`staging/${badId}`).get()).exists()).toBe(true); + expect((await rtdb.ref(`openSessions/${badId}`).get()).exists()).toBe(true); + } finally { + // discardSession refuses this id by design, so nothing else will clean + // it up. + await rtdb.ref(`staging/${badId}`).remove(); + await rtdb.ref(`openSessions/${badId}`).remove(); + } + }); + + it("does not re-queue a session whose recovery already completed before a discard failure", async () => { + // The other half of the fix above: partialFilenameFor is a pure function + // of the session, so a session that is STILL staged only because its + // discard failed computes the identical deduplication key on the next + // run. queueUpload's own dedup logic only special-cases "pending" and + // "processing" -- a "completed" doc falls through and gets freshly + // re-queued -- so without this check a second copy of an already- + // delivered partial would reach the provider. + const experimentID = await makeExperiment(); + const badId = "already-delivered-fixture"; // fails isValidSessionId, as above + await rtdb.ref(`openSessions/${badId}`).set({ + experimentId: experimentID, + owner: "staging-testuser", + startedAt: Date.now(), + expiresAt: Date.now() + 60000, + filename: "p02.csv", + }); + await stageTrials(badId, 3); + await markAbandoned(badId); + + const suffix = createHash("sha256").update(badId).digest("hex").slice(0, 8); + const filename = `p02-${suffix}.partial.json`; + const docId = `${experimentID}:${filename}`.replace(/[/\\]/g, "_"); + await db.collection("uploadQueue").doc(docId).set({ + experimentID, + owner: "staging-testuser", + filename, + storagePath: `upload-queue/${docId}`, + dataType: "data", + status: "completed", + errorCode: 0, + retryCount: 0, + maxRetries: 5, + createdAt: new Date(), + completedAt: new Date(), + deduplicationKey: `${experimentID}:${filename}`, + sessionIncremented: false, + partial: true, + }); + + try { + await sweepAbandonedSessions(new Set([badId])); + + // The pre-existing completed doc must be untouched, not overwritten + // back to "pending" by a fresh re-queue. + const doc = await db.collection("uploadQueue").doc(docId).get(); + expect(doc.data().status).toBe("completed"); + const entries = await queueEntriesFor(experimentID); + expect(entries).toHaveLength(1); + } finally { + await rtdb.ref(`staging/${badId}`).remove(); + await rtdb.ref(`openSessions/${badId}`).remove(); + await db.collection("uploadQueue").doc(docId).delete(); + } + }); + + it("records health and does not throw when the staging database is unreachable", async () => { + // The design doc's requirement: a broken sweep must show up as a recorded + // failure, never as an uncaught exception that takes the scheduled + // function down without a trace. Port 1 refuses the connection + // immediately, so this stays fast and needs no real network access. + const originalUrl = process.env.STAGING_DATABASE_URL; + process.env.STAGING_DATABASE_URL = "http://127.0.0.1:1/?ns=unreachable-staging-test"; + resetStagingHandleForTests(); + + try { + const stats = await sweepAbandonedSessions(new Set()); + expect(stats.errors).toBeGreaterThan(0); + } finally { + process.env.STAGING_DATABASE_URL = originalUrl; + resetStagingHandleForTests(); + } + }); }); // --------------------------------------------------------------------------- diff --git a/functions/src/api-data.ts b/functions/src/api-data.ts index eb4ad43..91acacf 100644 --- a/functions/src/api-data.ts +++ b/functions/src/api-data.ts @@ -35,10 +35,13 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 // // WHERE THIS IS CALLED, AND WHY THERE // - // Exactly where cleanupPending() is called, plus the four gates above that - // reject before a pending copy exists. Both sets are the same predicate: - // DataPipe either HAS the data somewhere durable (uploaded, or in the - // encrypted upload queue) or has DEFINITIVELY REFUSED this session. + // Everywhere cleanupPending() is called, plus the finalized / inactive / + // session-cap / unknown-experiment gates above that reject before a pending + // copy even exists. Both sets are (almost) the same predicate: DataPipe + // either HAS the data somewhere durable (uploaded, or in the encrypted + // upload queue) or has DEFINITIVELY REFUSED THE EXPERIMENT, not merely this + // submission -- see the "NOT CALLED ON EVERY REFUSAL" note below for the + // two exceptions. // // Both halves matter, and for opposite reasons. // @@ -53,6 +56,24 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 // sweep re-checks the gates itself, so this is the first of two doors, // not the only one.) // + // NOT CALLED ON EVERY REFUSAL, THOUGH. The four gates below this comment + // (finalized / inactive / session-cap / unknown experiment -- the last one + // is EXPERIMENT_NOT_FOUND, above this comment, which never staged anything + // in the first place) are refusals about the EXPERIMENT: nothing this + // submission does will ever be accepted, so the staged copy is genuinely + // worthless and discarding it is correct. INVALID_DATA and the + // duplicate-filename refusals (OSF_FILE_EXISTS, from either the collision + // cache's "duplicate" verdict or the provider's own NAME_CONFLICT) are + // refusals about THIS SUBMISSION -- this exact string failed validation, or + // this exact filename collided -- and say nothing about whether the + // experiment would accept the participant's trials under a different name + // or format. Discarding on those would destroy the one recoverable copy in + // exactly the case the staging tier exists for: a participant whose + // browser is never coming back to retry. Leaving it staged costs nothing + // extra -- the sweep's second door re-checks finalized/active before ever + // promoting it, so a since-finalized or since-deactivated experiment is + // still covered. + // // It is deliberately NOT called on DataPipe's own failures -- a persist // error, a token failure, a metadata error, an exception path that could not // even queue. Those are the cases the staging tier is FOR: the participant @@ -60,15 +81,17 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 // standing between that and lost data. Same reasoning as the "pending-data // copy is deliberately kept" note in the metadata branch below. // - // Never throws: discardSession swallows its own errors, because orphaned + // Never throws: discardSession reports its own errors through its boolean + // return value (ignored here) rather than throwing, because orphaned // staging data is a sweep's problem and must never turn a 201 into a 500. // // ORDERING: always awaited BEFORE res.json(), never after. On the branches // that reach cleanupPending() that is already true, because those clean up - // ahead of responding. On the four gates above it is a deliberate departure - // from the surrounding style -- writeLog() there runs AFTER the response -- - // and the difference is that a log write losing a race costs a log line, - // while this one leaves a participant's trials sitting in RTDB. Work queued + // ahead of responding. On the three experiment-state gates above it is a + // deliberate departure from the surrounding style -- writeLog() there runs + // AFTER the response -- and the difference is that a log write losing a + // race costs a log line, while this one leaves a participant's trials + // sitting in RTDB. Work queued // after a response is not guaranteed to run: the instance can be frozen or // scaled down the moment the response is flushed. Same reasoning as the // "logs are written BEFORE the response here" note on the NAME_CONFLICT @@ -81,8 +104,8 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 // away the empty path segments, resolves to "/staging" and "/openSessions" // themselves -- wiping every in-progress session for every experiment. This // gate runs on every request that reaches discardStaging, including the - // four unauthenticated ones above (finalized / inactive / session-cap / - // validation), so a closed experiment id was, before this check, enough to + // three unauthenticated experiment-state gates below (finalized / inactive / + // session-cap), so a closed experiment id was, before this check, enough to // reach it. discardSession guards the same thing again on its own input; // this is the first of the two doors, not the only one. const discardStaging = async () => { @@ -164,7 +187,10 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 } } if (!valid) { - await discardStaging(); + // Staging is deliberately LEFT ALONE here -- see the comment on + // discardStaging() above. This submission's string failed validation; + // the trials sitting in RTDB did not, and the sweep is what gives a + // participant who cannot retry a second chance at being recovered. res.status(400).json(MESSAGES.INVALID_DATA); await writeLog(experimentID, "logError", MESSAGES.INVALID_DATA, logContext); return; @@ -323,8 +349,13 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 if (!claimResult.claimed) { if (claimResult.reason === "duplicate") { + // cleanupPending() runs (the Cloud Storage pending copy really is + // superseded -- persist-pending.ts's own sweep would just rediscover an + // identical failure), but staging is deliberately LEFT ALONE -- see the + // comment on discardStaging() above. This filename collided; the + // participant's trials did not, and the sweep can still recover them + // under partialFilenameFor's own (hash-suffixed) name. await cleanupPending(pendingPath); - await discardStaging(); res.status(400).json({...MESSAGES.OSF_FILE_EXISTS, metadataMessage}); await writeLog(experimentID, "logError", MESSAGES.OSF_FILE_EXISTS, logContext); return; @@ -461,7 +492,10 @@ export const apiData = onRequest({ cors: true, memory: "512MiB", concurrency: 1 direction: "cache-free-provider-conflict", }, logContext); await cleanupPending(pendingPath); - await discardStaging(); + // Staging is deliberately LEFT ALONE here too -- see the comment on + // discardStaging() above. The provider refused this filename; the + // participant's trials did not, and the sweep can still recover them + // under partialFilenameFor's own (hash-suffixed) name. res.status(400).json({...MESSAGES.OSF_FILE_EXISTS, metadataMessage}); return; } diff --git a/functions/src/concurrency-limit.ts b/functions/src/concurrency-limit.ts new file mode 100644 index 0000000..df9f133 --- /dev/null +++ b/functions/src/concurrency-limit.ts @@ -0,0 +1,47 @@ +// A small bounded-concurrency map, used in place of an unbounded Promise.all. +// +// scheduled-staging-sweep.ts used to fan out one getSessionMeta call per open +// session -- up to MAX_RECONCILE (500) of them -- through a single +// `Promise.all(sessions.map(...))`. RTDB does not bill operations, so this was +// never a cost problem, but it is still 500 concurrent reads launched at once +// from a single 256MiB function instance with nothing bounding how many are +// in flight together. This runs the same work with at most `concurrency` +// promises outstanding at a time. +// +// Deliberately not a dependency (p-limit and friends): the whole thing is a +// worker-pool over a shared cursor, and pulling in a package for it would cost +// every future importer of this module a transitive dependency for a dozen +// lines of code. +// +// Pure and infrastructure-free on purpose -- same reasoning as +// staging-assembly.ts's split from staging.ts -- so it is testable without an +// emulator. + +/** + * Run `fn` over `items` with at most `concurrency` calls in flight at once. + * + * Results are returned in the same order as `items`, regardless of which + * call finishes first -- callers can treat this as a drop-in replacement for + * `Promise.all(items.map(fn))`. + */ +export async function mapWithConcurrency( + items: readonly T[], + concurrency: number, + fn: (item: T, index: number) => Promise +): Promise { + const results: R[] = new Array(items.length); + let nextIndex = 0; + + async function worker(): Promise { + for (;;) { + const index = nextIndex++; + if (index >= items.length) return; + results[index] = await fn(items[index], index); + } + } + + const workerCount = Math.max(1, Math.min(concurrency, items.length)); + await Promise.all(Array.from({ length: workerCount }, worker)); + + return results; +} diff --git a/functions/src/queue-upload.ts b/functions/src/queue-upload.ts index 0767949..a89b233 100644 --- a/functions/src/queue-upload.ts +++ b/functions/src/queue-upload.ts @@ -101,9 +101,24 @@ export function isProbeRetry(code?: string | null): boolean { return !!code && PROBE_RETRY_CODES.has(code); } +/** + * The uploadQueue document id for a given experiment/filename pair -- the + * same value stored as `deduplicationKey` below, sanitised into a legal + * Firestore document id. + * + * Exported so a caller that needs to know whether an entry ALREADY EXISTS for + * a filename -- scheduled-staging-sweep.ts, before re-queueing a recovered + * partial whose earlier discard may have failed -- computes the identical id + * this module uses, rather than keeping a second copy that could drift out of + * sync with it. + */ +export function queueDocIdFor(experimentID: string, filename: string): string { + return `${experimentID}:${filename}`.replace(/[/\\]/g, "_"); +} + export default async function queueUpload(params: QueueUploadParams): Promise { const deduplicationKey = `${params.experimentID}:${params.filename}`; - const docId = deduplicationKey.replace(/[/\\]/g, "_"); + const docId = queueDocIdFor(params.experimentID, params.filename); const docRef = db.collection("uploadQueue").doc(docId); diff --git a/functions/src/scheduled-staging-sweep.ts b/functions/src/scheduled-staging-sweep.ts index b56573c..1e6e424 100644 --- a/functions/src/scheduled-staging-sweep.ts +++ b/functions/src/scheduled-staging-sweep.ts @@ -42,9 +42,10 @@ import { onSchedule } from "firebase-functions/v2/scheduler"; import { Timestamp } from "firebase-admin/firestore"; import { db } from "./app.js"; -import queueUpload from "./queue-upload.js"; +import queueUpload, { queueDocIdFor } from "./queue-upload.js"; import { uploadPathFor } from "./metadata-derived-files.js"; import { ExperimentData } from "./interfaces.js"; +import { mapWithConcurrency } from "./concurrency-limit.js"; import { assembleSession, countOpenSessions, @@ -55,6 +56,7 @@ import { listOpenSessions, partialFilenameFor, OpenSession, + OpenSessionsCursor, ABANDON_GRACE_MS, disconnectedSince, } from "./staging.js"; @@ -71,11 +73,37 @@ export { ABANDON_GRACE_MS }; // scheduled-pending-recovery.ts's MAX_FILES_PER_RUN. const MAX_SESSIONS_PER_RUN = 10; -// Candidates read per run. More than can be processed, because most of the -// oldest open sessions on a busy deployment are LIVE rather than abandoned and -// are skipped without costing anything but a meta read. Mirrors the -// `maxResults: MAX_FILES_PER_RUN * 2` in the pending sweep. -const CANDIDATES_PER_RUN = MAX_SESSIONS_PER_RUN * 3; +// Candidates read PER PAGE. More than can be processed from a single page, +// because most of the oldest open sessions on a busy deployment are LIVE +// rather than abandoned and are skipped without costing anything but a meta +// read. Mirrors the `maxResults: MAX_FILES_PER_RUN * 2` in the pending sweep. +// +// THIS IS A PAGE SIZE, NOT A PER-RUN CAP. A single page used to be the whole +// candidate set: if the CANDIDATES_PER_PAGE oldest open sessions were all +// long-lived or zombie (their onDisconnect never fired, their tab is still +// technically open, whatever the reason), every one of them was skipped as +// live, the cursor never advanced, and everything ABANDONED behind them in +// the queue waited -- potentially for the full 24-hour TTL -- because nothing +// ever looked past position CANDIDATES_PER_PAGE. Paging past a skipped page is +// what fixes that; see the loop in sweepAbandonedSessions. +const CANDIDATES_PER_PAGE = MAX_SESSIONS_PER_RUN * 3; + +// Hard ceiling on pages fetched in one run, independent of how many sessions +// get skipped as live. Without this, a deployment with thousands of +// simultaneously live (not abandoned) sessions would have the sweep page +// through the entire table every five minutes looking for the few that are +// actually abandoned -- bounded work turning unbounded. 20 pages of +// CANDIDATES_PER_PAGE (30) is 600 sessions inspected per run at the most, which +// keeps a run's RTDB reads bounded the same way MAX_SESSIONS_PER_RUN bounds +// its writes. +const MAX_PAGES_PER_RUN = 20; + +// How many getSessionMeta calls the live-sessions reconciliation pass runs at +// once. Was an unbounded `Promise.all` over up to MAX_RECONCILE (500) open +// sessions; RTDB does not bill operations, so this was never a cost problem, +// but nothing bounded how many reads one 256MiB instance had in flight +// together. See concurrency-limit.ts. +const RECONCILE_CONCURRENCY = 20; export interface SweepStats { candidates: number; @@ -83,6 +111,8 @@ export interface SweepStats { discarded: number; skippedLive: number; errors: number; + /** Candidate pages fetched this run (see CANDIDATES_PER_PAGE / MAX_PAGES_PER_RUN). */ + pages: number; /** * Live-sessions mirror documents this run had to create, correct or delete. * Should be zero: each one is a write on the session lifecycle that failed or @@ -124,53 +154,96 @@ export async function sweepAbandonedSessions( discarded: 0, skippedLive: 0, errors: 0, + pages: 0, mirrorFixed: 0, }; let lastError: string | null = null; let openSessionCount: number | null = null; let mirror = { created: 0, updated: 0, deleted: 0 }; - try { - const candidates = (await listOldestOpenSessions(CANDIDATES_PER_RUN)).filter( - (s) => !only || only.has(s.sessionId) - ); - stats.candidates = candidates.length; + // `only`-scoped runs (tests) can stop as soon as every id they care about + // has been seen, rather than paging until the whole (possibly large, shared + // emulator) table is exhausted. Production runs (`only` undefined) always + // page until one of the other three stopping conditions below fires. + const pending = only ? new Set(only) : null; + try { const now = Date.now(); - - for (const session of candidates) { - if (stats.recovered + stats.discarded >= MAX_SESSIONS_PER_RUN) break; - - try { - const meta = await getSessionMeta(session.sessionId); - - const since = disconnectedSince(meta); - const abandoned = since !== null && now - since >= ABANDON_GRACE_MS; - // The backstop, for a client that died before it could register an - // onDisconnect at all, or one whose onDisconnect Firebase never ran. - // Without it such a session would sit in RTDB forever, being paid for. - const expired = - typeof session.expiresAt === "number" && now >= session.expiresAt; - - if (!abandoned && !expired) { - stats.skippedLive++; - continue; + let cursor: OpenSessionsCursor | undefined; + + pageLoop: for (let page = 0; page < MAX_PAGES_PER_RUN; page++) { + const rawPage = await listOldestOpenSessions(CANDIDATES_PER_PAGE, cursor); + if (rawPage.length === 0) break; + stats.pages++; + + // Advance the cursor off the RAW page (not the `only`-filtered one) + // regardless of whether anything on it matched -- this is what lets a + // wall of skipped-or-out-of-scope candidates be paged PAST instead of + // re-read forever. See CANDIDATES_PER_PAGE's comment. + const lastRow = rawPage[rawPage.length - 1]; + cursor = { expiresAt: lastRow.expiresAt, sessionId: lastRow.sessionId }; + + const pageCandidates = only ? rawPage.filter((s) => only.has(s.sessionId)) : rawPage; + stats.candidates += pageCandidates.length; + + for (const session of pageCandidates) { + if (stats.recovered + stats.discarded >= MAX_SESSIONS_PER_RUN) break pageLoop; + + try { + const meta = await getSessionMeta(session.sessionId); + + const since = disconnectedSince(meta); + const abandoned = since !== null && now - since >= ABANDON_GRACE_MS; + // The backstop, for a client that died before it could register an + // onDisconnect at all, or one whose onDisconnect Firebase never ran. + // Without it such a session would sit in RTDB forever, being paid for. + const expired = + typeof session.expiresAt === "number" && now >= session.expiresAt; + + if (!abandoned && !expired) { + stats.skippedLive++; + pending?.delete(session.sessionId); + continue; + } + + const result = await recoverSession(session); + if (result.discardOk) { + if (result.status === "recovered") stats.recovered++; + else stats.discarded++; + } else { + // The queue write (if any) happened, but the staging node is + // STILL THERE -- discardSession returned false rather than + // throwing. Reporting this as "recovered" or "discarded" would + // describe cleanup that has not actually happened: the session + // will surface again as a candidate next run (still open, still + // abandoned) and, if it was already queued, recoverSession's own + // dedup check is what stops that from becoming a duplicate + // delivery -- not this branch. Counted as an error so a discard + // path that is silently and persistently failing is visible in + // systemStatus/staging instead of being folded into "recovered". + console.error( + `Staging session ${session.sessionId} was ${result.status} but its RTDB node ` + + `could not be removed; it remains staged and will be retried next run.` + ); + stats.errors++; + } + pending?.delete(session.sessionId); + } catch (e) { + const detail = e instanceof Error ? e.message : "Unknown error"; + // One bad session must not stop the run: the others behind it are + // accruing storage cost, and a session that throws every time would + // otherwise block the queue permanently. + console.error( + `Failed to recover staging session ${session.sessionId}: ${detail}` + ); + lastError = detail; + stats.errors++; + pending?.delete(session.sessionId); } - - const outcome = await recoverSession(session); - if (outcome === "recovered") stats.recovered++; - else stats.discarded++; - } catch (e) { - const detail = e instanceof Error ? e.message : "Unknown error"; - // One bad session must not stop the run: the others behind it are - // accruing storage cost, and a session that throws every time would - // otherwise block the queue permanently. - console.error( - `Failed to recover staging session ${session.sessionId}: ${detail}` - ); - lastError = detail; - stats.errors++; } + + if (rawPage.length < CANDIDATES_PER_PAGE) break; // table exhausted + if (pending && pending.size === 0) break; // everything in scope was found } } catch (e) { lastError = e instanceof Error ? e.message : "Unknown error"; @@ -187,9 +260,13 @@ export async function sweepAbandonedSessions( const open = await listOpenSessions(); openSessionCount = open.length; const inScope = (only ? open.filter((s) => only.has(s.sessionId)) : open).slice(0, MAX_RECONCILE); - const entries = await Promise.all( - inScope.map(async (session) => ({ session, meta: await getSessionMeta(session.sessionId) })) - ); + // Bounded, not `Promise.all`: up to MAX_RECONCILE (500) sessions here, and + // nothing should put 500 concurrent RTDB reads in flight from one 256MiB + // instance at once. See concurrency-limit.ts. + const entries = await mapWithConcurrency(inScope, RECONCILE_CONCURRENCY, async (session) => ({ + session, + meta: await getSessionMeta(session.sessionId), + })); // Only for sessions opened before owner was recorded on openSessions; one // read per experiment per run. @@ -229,18 +306,33 @@ export async function sweepAbandonedSessions( return stats; } +/** What became of one candidate session, and whether its RTDB node is actually gone. */ +export interface RecoverResult { + /** "recovered" when a queue entry was written this run, "discarded" otherwise. */ + status: "recovered" | "discarded"; + /** + * Whether discardSession actually removed the staging node. When false, the + * session is still open and will reappear as a candidate next run -- see + * discardSession's own doc comment for why this matters more for + * "recovered" than for "discarded". + */ + discardOk: boolean; +} + /** * Turn one abandoned session into a queued partial upload, or discard it. * - * Returns "recovered" when a queue entry was written, "discarded" when the - * session was dropped without one. Either way the staging node is gone - * afterwards -- this function never leaves data behind for the next run to - * rediscover, because a session that cannot be recovered is a session that - * would otherwise be paid for forever. + * "Discarded" means the session was dropped without a queue entry. + * "Recovered" means one was written. Either way this function ATTEMPTS to + * remove the staging node before returning, but does not assume it succeeded + * -- see `discardOk` and discardSession's doc comment. A session whose discard + * failed is not lost: it stays in openSessions and is picked up again next + * run, and the completed-queue-entry check below is what makes that safe to + * retry rather than a source of duplicate deliveries. */ export async function recoverSession( session: OpenSession -): Promise<"recovered" | "discarded"> { +): Promise { const { sessionId, experimentId } = session; const expDoc = await db.collection("experiments").doc(experimentId).get(); @@ -248,8 +340,7 @@ export async function recoverSession( console.warn( `Staging session ${sessionId} belongs to missing experiment ${experimentId}; discarding.` ); - await discardSession(sessionId); - return "discarded"; + return { status: "discarded", discardOk: await discardSession(sessionId) }; } const expData = expDoc.data() as ExperimentData; @@ -258,8 +349,7 @@ export async function recoverSession( console.warn( `Experiment ${experimentId} has no owner; discarding staging session ${sessionId}.` ); - await discardSession(sessionId); - return "discarded"; + return { status: "discarded", discardOk: await discardSession(sessionId) }; } // THE SECOND DOOR. api-session-start.ts checked these gates when the session @@ -277,16 +367,14 @@ export async function recoverSession( console.log( `Experiment ${experimentId} is finalized; discarding staging session ${sessionId}.` ); - await discardSession(sessionId); - return "discarded"; + return { status: "discarded", discardOk: await discardSession(sessionId) }; } if (!expData.active) { console.log( `Experiment ${experimentId} is not collecting; discarding staging session ${sessionId}.` ); - await discardSession(sessionId); - return "discarded"; + return { status: "discarded", discardOk: await discardSession(sessionId) }; } const assembled = await assembleSession(sessionId); @@ -296,8 +384,7 @@ export async function recoverSession( // left. There is no data to recover and an empty file would be noise in the // researcher's dataset. if (assembled.trialCount === 0) { - await discardSession(sessionId); - return "discarded"; + return { status: "discarded", discardOk: await discardSession(sessionId) }; } const filename = partialFilenameFor(session); @@ -308,6 +395,25 @@ export async function recoverSession( // pending-recovery sweep carries. const uploadFilename = uploadPathFor(expData.metadataActive, filename); + // partialFilenameFor is a PURE function of the session -- same session id, + // same filename, every time. So if an earlier run already queued this exact + // session and that entry has since COMPLETED (the retry worker delivered + // it), the only reason this session is still here to be recovered again is + // that its discardSession call failed afterwards (see discardSession's doc + // comment). Re-queueing would land a second copy of the same partial with + // the provider: queueUpload's own dedup logic only special-cases "pending" + // and "processing" entries, and a "completed" one falls through and gets + // freshly re-queued. Checking here, before the call, is what stops that. + const docId = queueDocIdFor(experimentId, uploadFilename); + const existing = await db.collection("uploadQueue").doc(docId).get(); + if (existing.exists && existing.data()?.status === "completed") { + console.warn( + `Staging session ${sessionId} already delivered as uploadQueue/${docId}; ` + + `skipping re-queue and just clearing the leftover staging node.` + ); + return { status: "discarded", discardOk: await discardSession(sessionId) }; + } + const notes: string[] = [`${assembled.trialCount} trials`]; if (assembled.gaps > 0) notes.push(`${assembled.gaps} missing`); if (assembled.skipped > 0) notes.push(`${assembled.skipped} unreadable`); @@ -341,8 +447,7 @@ export async function recoverSession( failureReason: `Recovered from an abandoned session (${notes.join(", ")})`, }); - await discardSession(sessionId); - return "recovered"; + return { status: "recovered", discardOk: await discardSession(sessionId) }; } /** @@ -377,6 +482,7 @@ async function recordSweepHealth( lastRunAt: Timestamp.now(), openSessionCount, candidates: stats.candidates, + pages: stats.pages, recovered: stats.recovered, discarded: stats.discarded, skippedLive: stats.skippedLive, diff --git a/functions/src/staging-assembly.ts b/functions/src/staging-assembly.ts index d95eccd..315fd1a 100644 --- a/functions/src/staging-assembly.ts +++ b/functions/src/staging-assembly.ts @@ -16,6 +16,11 @@ // See functions/src/staging.ts for the Realtime Database half, and // docs/streaming-ingest-design.md for the design. +// crypto is a Node builtin, not an ESM-only package like nanoid -- importing +// it here does not reintroduce the problem staging.ts's module boundary +// exists to avoid (see the module comment above). +import { createHash } from "crypto"; + // Mirrors the `newData.val().length <= 65536` cap in database.rules.json. // Exported so api-session-start.ts can tell the client the number rather than // letting the plugin carry its own copy that could drift out of sync with the @@ -132,7 +137,7 @@ export interface AssembledSession { /** * The name a recovered partial session is stored under. * - * Three things happen here, and each one is load-bearing: + * Four things happen here, and each one is load-bearing: * * - `/` and `\` are replaced, exactly as persistPending does. The base name * is CLIENT-SUPPLIED, and it becomes a path in a researcher's Drive, OSF or @@ -147,13 +152,25 @@ export interface AssembledSession { * storage. A recovered fragment sitting in a dataset under an ordinary name * would quietly make the record non-Psych-DS -- the same class of problem * docs/finalization-spec.md addresses for archives. + * - A SHORT HASH OF THE SESSION ID is appended to a client-supplied name, + * unconditionally. Two participants in the same experiment can supply the + * same filename -- "data.csv" is a common default -- and without this both + * abandoned sessions produce the identical `experimentID:filename` + * deduplication key in queue-upload.ts. That is not a queueing conflict: + * the second recovery silently OVERWRITES the first session's payload in + * Cloud Storage before it is ever delivered, so one participant's data is + * destroyed with no error anywhere. The session id is unforgeable and + * unique by construction (staging.ts), so hashing it is enough to make + * every recovered filename unique regardless of what the browser sent -- + * without the hash needing to be reversible or human-legible. * - * Falls back to the session id when no filename was captured. Opaque, but - * unique and traceable back to a sweep log line. + * Falls back to the session id (with no hash appended -- it is already + * unique) when no filename was captured. Opaque, but unique and traceable + * back to a sweep log line. */ export function partialFilenameFor(session: OpenSession): string { const fallback = `session-${session.sessionId}`; - const base = (session.filename || fallback) + const sanitized = (session.filename || "") // Path separators first: the name becomes a path in a researcher's Drive, // OSF or Zenodo container, and nothing downstream re-checks it. .replace(/[/\\]/g, "_") @@ -168,7 +185,11 @@ export function partialFilenameFor(session: OpenSession): string { // A leading dot makes a hidden file on every POSIX system, and would hide // a participant's recovered data from the researcher looking for it. .replace(/^\.+/, ""); - return `${base || fallback}.partial.json`; + + if (!sanitized) return `${fallback}.partial.json`; + + const suffix = createHash("sha256").update(session.sessionId).digest("hex").slice(0, 8); + return `${sanitized}-${suffix}.partial.json`; } /** diff --git a/functions/src/staging.ts b/functions/src/staging.ts index ef06ed6..7fd8851 100644 --- a/functions/src/staging.ts +++ b/functions/src/staging.ts @@ -268,28 +268,39 @@ export async function getSessionMeta(sessionId: string): Promise { return snap.exists() ? (snap.val() as SessionMeta) : {}; } +/** Where a page of `listOldestOpenSessions` left off, for the next page. */ +export interface OpenSessionsCursor { + expiresAt: number; + sessionId: string; +} + /** - * Oldest-first candidates for the sweep. + * Oldest-first candidates for the sweep, one page at a time. * * Ordered by `expiresAt`, which for a fixed TTL is the same order as - * `startedAt` -- so the sessions most likely to be abandoned surface first, - * and the read is bounded regardless of how many sessions are in flight. The - * caller filters live ones out by reading each candidate's meta; this only has - * to make sure it never reads the whole table. + * `startedAt` -- so the sessions most likely to be abandoned surface first. + * The caller filters live ones out by reading each candidate's meta; this only + * has to make sure it never reads the whole table in one request. + * + * `after`, when given, resumes past the last row of a previous page rather + * than re-reading from the top -- this is what lets the sweep keep paging + * instead of being stuck re-fetching the same `limit` oldest sessions every + * run. RTDB's `startAfter(value, key)` breaks ties on the key, which is why + * the cursor carries the session id alongside `expiresAt`: two sessions can + * share an `expiresAt` (same TTL, same millisecond), and a cursor keyed on + * `expiresAt` alone could skip or repeat a row at that boundary. * * Requires the `.indexOn: ["expiresAt"]` directive in database.rules.json -- * without it RTDB still answers, but by downloading the entire node and - * sorting in the client, which is precisely what the bound above exists to - * avoid. + * sorting in the client, which is precisely what paging exists to avoid. */ export async function listOldestOpenSessions( - limit: number + limit: number, + after?: OpenSessionsCursor ): Promise { - const snap = await rtdb() - .ref("openSessions") - .orderByChild("expiresAt") - .limitToFirst(limit) - .get(); + const ordered = rtdb().ref("openSessions").orderByChild("expiresAt"); + const query = after ? ordered.startAfter(after.expiresAt, after.sessionId) : ordered; + const snap = await query.limitToFirst(limit).get(); if (!snap.exists()) return []; const rows: OpenSession[] = []; snap.forEach((child) => { @@ -357,26 +368,41 @@ export async function assembleSession( * gone -- so a completed session cannot later be marked abandoned by a socket * closing. * - * Best-effort by contract. Every caller has already got the participant's data - * somewhere durable by the time it calls this, so a failure here is orphaned - * staging data (which the sweep collects on its next pass) rather than lost - * data. It must never turn a successful submission into an error response. + * Best-effort by contract: NEVER THROWS. Every caller has already got the + * participant's data somewhere durable by the time it calls this, so a + * failure here is orphaned staging data rather than lost data, and it must + * never turn a successful submission into an error response. + * + * Returns whether the RTDB removal actually happened, so callers that matter + * -- scheduled-staging-sweep.ts, specifically -- can tell "gone" from + * "still there, try again next run" instead of assuming success. This is not + * a formality: the sweep hands a recovered session's data to queueUpload + * BEFORE calling this, so a discard that fails silently after that write + * leaves the staging node in place. The next run would otherwise reassemble + * and re-queue the same session under the same deterministic filename + * (partialFilenameFor is a pure function of the session), and if the first + * queue entry had already reached the provider by then, that re-queue lands a + * second copy of the same partial. api-data.ts's call sites ignore the return + * value -- a completion response was already sent by the time discardStaging + * runs there, so there is nothing left to condition on. * * Guards its own input, in addition to every caller checking first: this is * exported and called from several places (api-data.ts, the sweep), and it is * the function that actually splices sessionId into an RTDB path. A value * that fails isValidSessionId is refused here even if some future caller - * forgets to -- logged and returned, never thrown, matching the best-effort - * contract above. See isValidSessionId for why this matters: "/", "//", and - * "../x" all normalize to paths this function must never touch. + * forgets to -- logged and reported as a failed discard, never thrown, + * matching the best-effort contract above. See isValidSessionId for why this + * matters: "/", "//", and "../x" all normalize to paths this function must + * never touch. */ -export async function discardSession(sessionId: string): Promise { +export async function discardSession(sessionId: string): Promise { if (!isValidSessionId(sessionId)) { console.warn( `Refusing to discard session with an invalid id: ${JSON.stringify(sessionId)}` ); - return; + return false; } + let removed = true; try { await rtdb() .ref() @@ -387,11 +413,16 @@ export async function discardSession(sessionId: string): Promise { } catch (e) { const detail = e instanceof Error ? e.message : "Unknown error"; console.error(`Failed to discard staging session ${sessionId}: ${detail}`); + removed = false; } // And the researcher's dashboard row. Every way a session ends -- clean // completion, a gate refusing it, the sweep recovering or discarding it -- // comes through here, which is why this is the one place that removes it. // Separate from the RTDB update on purpose: either can fail without the // other being skipped, and the sweep's reconciliation collects what is left. + // Not reflected in the return value: it is the RTDB removal above that a + // caller needs to know about to avoid re-queueing, and a mirror-delete + // failure alone cannot cause that. await removeLiveSession(sessionId); + return removed; } diff --git a/pages/docs/experiments/sending-data.js b/pages/docs/experiments/sending-data.js index 2a25710..2eaafc5 100644 --- a/pages/docs/experiments/sending-data.js +++ b/pages/docs/experiments/sending-data.js @@ -115,9 +115,11 @@ export default function SendingDataPage() { An abandoned session becomes a second kind of file.{" "} DataPipe assembles the trials it received and stores them as{" "} - <your filename>.partial.json — JSON even if your - experiment submits CSV, because it is rebuilt from individual trials - rather than from the string your experiment would have sent. Plan + <your filename>-<id>.partial.json — JSON + even if your experiment submits CSV, because it is rebuilt from + individual trials rather than from the string your experiment would + have sent. The short id keeps two participants who happened to use + the same filename from colliding on the same recovered file. Plan for that in your analysis, and treat a partial file as a participant who did not finish. Partial sessions do not count toward your session limit. From 9d59c799680b82aecccc99c045753f7a6bec9968 Mon Sep 17 00:00:00 2001 From: Josh de Leeuw Date: Sun, 13 Sep 2026 17:31:23 -0400 Subject: [PATCH 3/6] Fix flaky concurrency-cap tests: direct discard, realistic reconcile fixture Two of the new "per-experiment concurrency cap" tests (staging-emulator.test.js) failed against a real emulator run: "lets a new session in once a prior one completes" routed the release through a full /api/data completion, which also exercises validation, the provider write and the metadata pipeline -- none of which this test is about -- so its pass/fail depended on that unrelated round trip instead of on discardSession's own slot release. Calls discardSession() directly now, the same way sweepAbandonedSessions is already exercised directly elsewhere in this file. "lets a new session in once an abandoned one is swept" seeded openSessionCounts directly with no matching openSessions entries behind it, then ran a real sweep. reconcileOpenSessionCounts recomputes an experiment's counter from the real entries under openSessions on every run -- correct behavior, since the counter has no legitimate reason to differ from what is actually open -- so it rightly overwrote the unsupported seeded value back to the true (lower) count. The fixture was wrong, not the reconcile: a new seedOpenSessions() helper seeds real (if synthetic) openSessions entries to match, parked a year out on expiresAt so they never crowd a same-run real session out of listOldestOpenSessions' oldest-30 window. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Q8xf16ov88M3KojPVfQwJT --- .../src/__tests__/staging-emulator.test.js | 83 ++++++++++++++++--- 1 file changed, 73 insertions(+), 10 deletions(-) diff --git a/functions/src/__tests__/staging-emulator.test.js b/functions/src/__tests__/staging-emulator.test.js index 1c6f99e..d68dae6 100644 --- a/functions/src/__tests__/staging-emulator.test.js +++ b/functions/src/__tests__/staging-emulator.test.js @@ -39,6 +39,7 @@ const { sweepAbandonedSessions, ABANDON_GRACE_MS, } = require("../../lib/scheduled-staging-sweep.js"); +const { discardSession } = require("../../lib/staging.js"); const FUNCTIONS_HOST = process.env.FUNCTIONS_EMULATOR_HOST || "localhost:5001"; const PROJECT_ID = "datapipe-test"; @@ -122,6 +123,40 @@ const queueEntriesFor = async (experimentID) => (await db.collection("uploadQueue").where("experimentID", "==", experimentID).get()) .docs.map((d) => d.data()); +/** + * Seed `count` synthetic openSessions entries for `experimentID` -- fixtures + * for reconcileOpenSessionCounts's ground truth, not sessions admitted + * through tryAdmitSession. + * + * reconcileOpenSessionCounts (staging.ts) recomputes an experiment's counter + * from the REAL entries under openSessions on every sweep run; that is + * correct behaviour, not drift-tolerance gone wrong, since the counter has no + * legitimate reason to differ from what is actually open. A test that sets + * openSessionCounts/{id} directly without any matching openSessions entries + * is therefore asking the sweep to preserve a value ground truth does not + * support, and the sweep is right to overwrite it. Tests that seed the + * counter directly and then run a sweep must seed matching entries here too. + * + * expiresAt is set far in the future (a year out, vs. a real session's ~24h) + * on purpose: listOldestOpenSessions(CANDIDATES_PER_RUN) -- unlike + * listOpenSessions(), which reconcileOpenSessionCounts uses and does not + * care about ordering -- orders by expiresAt ascending and takes the + * OLDEST 30. Hundreds of fixtures at a real session's ~24h expiry would + * crowd a same-run real session out of that window and make the sweep skip + * it entirely; parked a year out, they never compete for those slots. + */ +async function seedOpenSessions(experimentID, count) { + const updates = {}; + for (let i = 0; i < count; i++) { + updates[`openSessions/concurrency-fixture-${experimentID}-${i}`] = { + experimentId: experimentID, + startedAt: Date.now(), + expiresAt: Date.now() + 365 * 86400000, + }; + } + await rtdb.ref().update(updates); +} + // Only the docs THIS suite created -- a collection-wide wipe here would delete // uploadQueue docs belonging to whatever suite is running in parallel, which is // half of the long-standing cross-suite flake documented in @@ -304,7 +339,7 @@ describe("the per-experiment concurrency cap", () => { expect(body.error).toBe("SESSION_START_ERROR"); }); - it("lets a new session in once a prior one completes and frees its slot", async () => { + it("lets a new session in once a prior one is discarded, freeing its slot", async () => { const experimentID = await makeExperiment(); await rtdb.ref(`openSessionCounts/${experimentID}`).set(CAP - 1); @@ -314,14 +349,16 @@ describe("the per-experiment concurrency cap", () => { const { status: refused } = await startSession({ experimentID }); expect(refused).toBe(503); - // Completion runs discardSession(), which must release the slot the first - // session reserved at admission. - await saveData({ - experimentID, - filename: "p01.csv", - data: "trial_type\nhtml-keyboard-response\n", - sessionId: first.sessionId, - }); + // discardSession() is what every terminal path funnels through -- + // completion, a gate refusal, the sweep -- and it is what must release + // the slot admission reserved. Called directly, the same way the sweep + // itself is exercised directly elsewhere in this file, rather than + // through a full /api/data completion: that path also runs validation, + // the provider write and the metadata pipeline, none of which this test + // is about, and routing through it would make this test's pass/fail + // depend on the mocked provider round trip instead of on the one thing + // it exists to check -- the counter release. + await discardSession(first.sessionId); expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP - 1); const { status: secondStatus } = await startSession({ experimentID }); @@ -330,17 +367,43 @@ describe("the per-experiment concurrency cap", () => { it("lets a new session in once an abandoned one is swept, freeing its slot", async () => { const experimentID = await makeExperiment(); + // A REALISTIC fixture, unlike the two tests above: this one runs a real + // sweep, and the sweep's reconcileOpenSessionCounts recomputes the + // counter from the real entries under openSessions every time it runs. + // A counter seeded without matching sessions behind it would just get + // corrected back to the true (lower) count -- rightly, since the counter + // has no legitimate reason to differ from what is actually open. Seeding + // CAP-1 other open sessions makes CAP-1 the true count once `first` is + // admitted and then swept away, so this test exercises "does completing + // one session let another in" rather than "does the sweep preserve an + // unsupported number". + await seedOpenSessions(experimentID, CAP - 1); await rtdb.ref(`openSessionCounts/${experimentID}`).set(CAP - 1); - const { body: first } = await startSession({ experimentID }); + const { body: first } = await startSession({ experimentID }); // CAP-1 seeded + first = CAP real, counter CAP + expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP); await markAbandoned(first.sessionId); const stats = await sweepAbandonedSessions(new Set([first.sessionId])); expect(stats.discarded).toBe(1); // nothing was staged, so it is discarded not recovered + // True regardless of which mechanism produced it -- discardSession's own + // release, or the same run's counter reconciliation catching a missed + // one -- both are legitimate production paths to the same correct number. expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(CAP - 1); const { status } = await startSession({ experimentID }); expect(status).toBe(200); + + // Fixtures only, never picked up by any candidate window (see + // seedOpenSessions) and so never cleaned up by anything under test -- + // remove them rather than leaving 499 rows in the shared emulator. + await rtdb + .ref("openSessions") + .update( + Object.fromEntries( + Array.from({ length: CAP - 1 }, (_, i) => [`concurrency-fixture-${experimentID}-${i}`, null]) + ) + ); }); }); From 4795bb6937250cc7fbe3848558e064f6db8b97a8 Mon Sep 17 00:00:00 2001 From: Josh de Leeuw Date: Sun, 13 Sep 2026 17:50:46 -0400 Subject: [PATCH 4/6] Prove a real /api/data completion releases the concurrency-cap slot The concurrency-cap tests no longer have anything covering the exact path that originally read openSessionCounts as 500 instead of 499: the failing test was rewritten to call discardSession() directly, which proves discardSession releases a slot but nothing end-to-end proved a real completion reaches it. Root cause of the original reading, found while building this test: the file's shared "staging-testuser" owner fixture sets connectedAccounts.gdrive = { accessToken, refreshToken }, but providers/gdrive.ts's resolveToken() reads .encryptedToken and .tokenExpiresAt -- neither of which that fixture has. Every completion attempt against it therefore fails at token resolution, before ever reaching discardStaging -- the "DataPipe itself fails" branch, which by design leaves the staged copy AND its counter slot in place for the sweep to recover (see "keeps the staged copy when DataPipe itself fails"). The 500-instead-of-499 reading was that correct, intentional behavior, not a bug: the test believed it was exercising a real completion when it was actually exercising a token failure. Added "releases the concurrency-cap slot on a genuine, non-refused completion" (functions/src/__tests__/staging-emulator.test.js, top of the "completion" describe block) with its own owner whose token shape resolveToken() can actually resolve, so the request clears every gate and reaches the provider. It cannot get a literal 201: this file has no mock Google Drive server on GDRIVE_API_BASE (gdrive-emulator.test.js owns that fixed port for the whole run, and a second listener would collide), so the unreachable host makes claimFilename's cold-cache rehydration throw and api-data.ts queues the upload for retry (202) instead -- still a genuine, accepted completion, and that branch calls discardStaging on the way to responding exactly as a real 201 would. Asserts the staging node and capability record are gone and openSessionCounts is decremented (0 or absent) afterward. `cd functions && npm run build` and `npm run lint` (repo root) are both clean. Pure tests re-run from the repo root, all still passing: staging-assembly (35), rules-constants (3), concurrency-limit (13), live-sessions (11), staging-session-id (4). Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Q8xf16ov88M3KojPVfQwJT --- .../src/__tests__/staging-emulator.test.js | 83 +++++++++++++++++++ 1 file changed, 83 insertions(+) diff --git a/functions/src/__tests__/staging-emulator.test.js b/functions/src/__tests__/staging-emulator.test.js index 6ca95ca..ae1da09 100644 --- a/functions/src/__tests__/staging-emulator.test.js +++ b/functions/src/__tests__/staging-emulator.test.js @@ -425,6 +425,89 @@ describe("the streaming kill switch", () => { }); describe("completion", () => { + it("releases the concurrency-cap slot on a genuine, non-refused completion", async () => { + // THE GAP THIS CLOSES: an earlier version of "lets a new session in once + // a prior one completes" (the per-experiment concurrency cap tests above) + // routed through a real /api/data completion and read + // openSessionCounts/ as 500 instead of 499 afterwards. That + // test was rewritten to call discardSession() directly, which proves + // discardSession releases a slot but nothing end-to-end proved that a + // real completion actually REACHES discardSession. + // + // ROOT CAUSE, FOUND HERE: "staging-testuser" (this file's shared owner + // fixture, set up in the top-level beforeAll) has + // connectedAccounts.gdrive = { accessToken, refreshToken }. + // providers/gdrive.ts's resolveToken() reads + // connectedAccounts.gdrive.encryptedToken and .tokenExpiresAt -- neither + // of which that fixture sets -- so EVERY completion attempt against it + // fails at token resolution, before ever reaching discardStaging. That is + // the "DataPipe itself fails" branch ("keeps the staged copy when + // DataPipe itself fails", below), which by design leaves the staged copy + // AND its counter slot in place for the sweep to recover. The original + // 500-instead-of-499 reading was that correct, intentional behaviour, + // not a bug -- it just meant the test believed it was exercising a real + // completion when it was actually exercising a token failure. + // + // This test uses its OWN owner, with a token shape resolveToken() can + // actually resolve, so the request clears every gate and reaches the + // provider. It still cannot get a literal 201: this file has no mock + // Google Drive server listening on GDRIVE_API_BASE (gdrive-emulator.test.js + // owns that fixed port, 127.0.0.1:3579, for the whole test run, and two + // listeners on it would collide). The unreachable host makes + // claimFilename's cold-cache rehydration throw, which api-data.ts treats + // as a retryable provider failure -- queued (202), not refused -- and + // that branch calls discardStaging on the way to responding, exactly as + // a real 201 would. That is the property this test actually needs: a + // genuine, accepted completion, not a specific status code. + const ownerId = `staging-token-ok-${randomUUID()}`; + await db.collection("users").doc(ownerId).set({ + uid: ownerId, + email: "staging-token-ok@example.com", + experiments: [], + connectedAccounts: { + gdrive: { + // decrypt()'s plaintext fallback (crypto-utils.ts): any string + // without the "v1:" version prefix round-trips unchanged, so this + // does not need TOKEN_ENCRYPTION_KEY to agree between this process + // and the functions emulator's. + encryptedToken: "fake-access-token", + tokenExpiresAt: Date.now() + 60 * 60 * 1000, + }, + }, + }); + + const experimentID = await makeExperiment({ owner: ownerId }); + const { body } = await startSession({ experimentID }); + await stageTrials(body.sessionId, 3); + + expect((await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val()).toBe(1); + + const { status, body: response } = await saveData({ + experimentID, + filename: "p01.csv", + data: "trial_type\nhtml-keyboard-response\n", + sessionId: body.sessionId, + }); + + // A genuine, accepted completion -- not a gate refusal (400) and not a + // DataPipe-side failure (400) -- so discardStaging is reached exactly as + // it would be for a real 201. + expect(status).toBe(202); + expect(response.error).toBeNull(); + expect((await rtdb.ref(`staging/${body.sessionId}`).get()).exists()).toBe(false); + expect((await rtdb.ref(`openSessions/${body.sessionId}`).get()).exists()).toBe(false); + + // THE ASSERTION THE ORIGINAL REGRESSION NEEDED: the slot this session + // reserved at admission is released once its completion -- discardStaging, + // running inside apidata's real handler -- has actually happened. Absent + // or 0 either way: releaseOpenSessionSlot always writes 0 rather than + // deleting the node, but nothing here depends on which. + const counter = (await rtdb.ref(`openSessionCounts/${experimentID}`).get()).val(); + expect(counter === null || counter === 0).toBe(true); + + await db.collection("users").doc(ownerId).delete(); + }); + it("drops the staged copy once a gate refuses the submission", async () => { // Keeping it would let the sweep re-offer the session as a .partial.json, // producing a file the gate just refused -- and, for a finalized From bf278c83a05e67b02df7a3021040c83f79f45761 Mon Sep 17 00:00:00 2001 From: Josh de Leeuw Date: Mon, 14 Sep 2026 08:23:40 -0400 Subject: [PATCH 5/6] Document the streaming-ingest tier's hard limits for researchers The staging tier's limits (trial size, trial count, abandonment grace, disconnect slots, per-experiment concurrency, assembled-file size and filename length) were only visible as numbers in the /api/session example response on pages/docs/api.js, with no prose explaining what hitting any of them actually does to a researcher's data. Added a "Limits" subsection to the "Saving as you go" page (pages/docs/experiments/sending-data.js, id "streaming-limits", registered in lib/docs-nav.js) stating each limit's value and consequence, verified against database.rules.json, functions/src/staging-assembly.ts, functions/src/api-session-start.ts and functions/src/scheduled-staging-sweep.ts. The numbers are hardcoded with a comment naming the constant each mirrors, since pages cannot import from functions/src. Also added the missing "maxDisconnects": 20 to the /api/session example response in pages/docs/api.js (the endpoint already returns it) and a GuidanceLine pointing streaming users at the new subsection. `npm run lint` is clean. No test in __tests__ or lib currently exercises docs-nav.js or these two pages directly (grepped for both, per the parent task's instructions); DocsLayout's dev-only console.error assertion is the only existing enforcement of section-id/nav agreement, and the new id is registered to satisfy it. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Q8xf16ov88M3KojPVfQwJT --- lib/docs-nav.js | 1 + pages/docs/api.js | 10 +++- pages/docs/experiments/sending-data.js | 77 ++++++++++++++++++++++++++ 3 files changed, 87 insertions(+), 1 deletion(-) diff --git a/lib/docs-nav.js b/lib/docs-nav.js index 68b4d45..ab03948 100644 --- a/lib/docs-nav.js +++ b/lib/docs-nav.js @@ -89,6 +89,7 @@ export const DOCS_NAV = [ sections: [ { id: "jspsych", label: "jsPsych" }, { id: "plain-javascript", label: "Plain JavaScript" }, + { id: "streaming-limits", label: "Limits" }, { id: "filenames-must-be-unique", label: "Filenames must be unique" }, { id: "media-and-binary-files", label: "Media and binary files" }, { id: "request-size-limits", label: "Request size limits" }, diff --git a/pages/docs/api.js b/pages/docs/api.js index a00734f..321fef2 100644 --- a/pages/docs/api.js +++ b/pages/docs/api.js @@ -87,6 +87,13 @@ export default function ApiReferencePage() { Compressing a request body yourself, and what the size limit means in practice. + + The trial size, session, abandonment and file-size limits that apply + only to incremental sessions. + @@ -198,7 +205,8 @@ export default function ApiReferencePage() { "maxTrialBytes": 16384, "maxTrials": 1000, "flushIntervalMs": 10000, - "flushEveryNTrials": 10 + "flushEveryNTrials": 10, + "maxDisconnects": 20 }`} diff --git a/pages/docs/experiments/sending-data.js b/pages/docs/experiments/sending-data.js index 2eaafc5..80fc92e 100644 --- a/pages/docs/experiments/sending-data.js +++ b/pages/docs/experiments/sending-data.js @@ -153,6 +153,83 @@ export default function SendingDataPage() { + + + Save as you go enforces the limits below on every request. None of + them is configurable, and hitting one never breaks your + experiment — the plugin carries on and a completed submission is + unaffected. What each one actually costs a participant is the + partial-file safety net for someone who never finishes, not the + data your experiment collects. + + + {/* MAX_TRIAL_BYTES, functions/src/staging-assembly.ts (mirrored in + database.rules.json's per-trial `.length` cap) */} + + 16 KiB per trial. A trial larger than that is + refused by the database — the write for that one trial fails, and + the plugin moves on to the next. A completed session still sends + your whole dataset in its final submission, so that trial is only + missing from the partial file DataPipe would recover if the + participant never finished. + + {/* MAX_TRIALS_PER_SESSION, functions/src/staging-assembly.ts + (mirrored in database.rules.json's `$seq` pattern) */} + + 1,000 trials per session. The 1,001st trial and + every one after it are refused the same way an oversized trial + is. Again, only the partial-file safety net is affected — the + final submission is not built from staged trials, so it is + unaffected. + + {/* ABANDON_GRACE_MS and SESSION_TTL_MS, + functions/src/staging-assembly.ts */} + + 10 minutes to reconnect, 24 hours to finish. If a + participant's connection drops and DataPipe sees no reconnect + and no further trial from them for 10 minutes, the session is + treated as abandoned and turned into a partial file the next time + the sweep runs. Reconnecting — or getting even one more trial + through — within that window keeps the session going as if + nothing happened. Regardless of any of that, every session + expires 24 hours after it started and is recovered the same way + whether or not a disconnect was ever recorded. + + {/* MAX_DISCONNECTS, functions/src/staging-assembly.ts (mirrored in + database.rules.json's 1..20 slot pattern) */} + + 20 disconnects and 20 reconnects per session. A + participant whose connection drops and recovers more than 20 + times stops having further drops recorded, so the 10-minute + abandonment clock keeps being measured from the last drop that + was recorded rather than the most recent real one. They are still + recovered eventually — at the 24-hour expiry if nothing else — + but the fast path may miss them. + + {/* MAX_OPEN_SESSIONS_PER_EXPERIMENT, + functions/src/staging-assembly.ts */} + + 500 sessions open per experiment at once. A + participant who requests a session while 500 are already open for + your experiment gets none — the same response as when + incremental upload is switched off — and their experiment runs + and submits exactly as it would without it. Nothing about their + data is different. + + {/* MAX_ASSEMBLED_BYTES and MAX_FILENAME_LENGTH, + functions/src/staging-assembly.ts */} + + 24 MiB per recovered file, 200-character filenames.{" "} + A recovered partial file stops growing at 24 MiB — trials beyond + that point are left out of the file DataPipe assembles. The + filename you give when starting a session is capped at 200 + characters, and is silently shortened past that when used to name + a recovered file; it never affects the filename you submit on a + clean completion. + + + + Two submissions to the same experiment can never share a filename: From 2f52b46af4dcab3d53eef805fb85243d51fac27a Mon Sep 17 00:00:00 2001 From: Josh de Leeuw Date: Mon, 14 Sep 2026 08:25:25 -0400 Subject: [PATCH 6/6] Register the remaining rendered-but-unlisted docs section ids lib/docs-nav.js's section lists for /docs/experiments/sending-data and /docs/api were missing two ids that were already being rendered: "saving-as-you-go" (the "Saving as you go" section, added before this PR's own "streaming-limits" addition) and "start-session" ("Start an incremental session" on the API reference page). Both are registered now, in their rendered positions, with labels matching their section titles. Checked every other DocsSection id rendered on both pages against the nav list -- no others were missing. `npm run lint` is clean. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Q8xf16ov88M3KojPVfQwJT --- lib/docs-nav.js | 2 ++ 1 file changed, 2 insertions(+) diff --git a/lib/docs-nav.js b/lib/docs-nav.js index ab03948..58f1b85 100644 --- a/lib/docs-nav.js +++ b/lib/docs-nav.js @@ -89,6 +89,7 @@ export const DOCS_NAV = [ sections: [ { id: "jspsych", label: "jsPsych" }, { id: "plain-javascript", label: "Plain JavaScript" }, + { id: "saving-as-you-go", label: "Saving as you go" }, { id: "streaming-limits", label: "Limits" }, { id: "filenames-must-be-unique", label: "Filenames must be unique" }, { id: "media-and-binary-files", label: "Media and binary files" }, @@ -187,6 +188,7 @@ export const DOCS_NAV = [ sections: [ { id: "limits", label: "Limits" }, { id: "save-text-data", label: "Save text data" }, + { id: "start-session", label: "Start an incremental session" }, { id: "save-base64-data", label: "Save base64-encoded data" }, { id: "get-condition", label: "Get condition assignment" }, { id: "responses", label: "Responses" },