Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 57 additions & 13 deletions __tests__/database-rules.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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', () => {
Expand Down Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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 () => {
Expand Down Expand Up @@ -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,
})
);
});
Expand All @@ -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,
})
);
});
Expand Down
76 changes: 66 additions & 10 deletions database.rules.json
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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": {
Expand All @@ -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.
Expand Down Expand Up @@ -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
Expand Down
66 changes: 62 additions & 4 deletions docs/streaming-ingest-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<name>.partial.json`, do not increment
`sessions`, and do not arm the upload-failure notifier. This answers open
question 2.
3. **Partial sessions** upload as `<name>-<hash>.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
Expand Down Expand Up @@ -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): `<name>.partial.json`, uncounted, no notification.**
2. **ANSWERED (decision 3): `<name>-<hash>.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`
Expand Down Expand Up @@ -419,6 +421,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
Expand Down
7 changes: 7 additions & 0 deletions functions/.env.datapipe-test
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading