Skip to content

feat(sync): pull a stream's journal - verify, stage once, apply (T12343 S5-1) - #1959

Open
kryptobaseddev wants to merge 5 commits into
feat/T12343-pushfrom
feat/T12343-pull
Open

kryptobaseddev wants to merge 5 commits into
feat/T12343-pushfrom
feat/T12343-pull

Conversation

@kryptobaseddev

Copy link
Copy Markdown
Owner

Task: T12343 (S5-1: pull)

Stacked on #1957. Nothing calls pullSyncStream yet (cleo cloud sync, T12996, will).

What

pullStream(db, …) (store/sync/pull.ts) pages through the stream after the store's persisted _sync_cursor. The journal client has already verified order, gaps, the replica-to-device pin, trust, hashes and decryption. Each page is decoded and verified in full before anything is written, then staged in one BEGIN IMMEDIATE:

  • Decode: deflate-raw canonical JSON of LedgerTxn[], parsed with the ledger schema.
    • The vault-delta skip (hard requirement): a cleo-vault-delta/v1 body (plain JSON, a count correction a vault push appended) carries no ops, is passed over, and is counted in vaultDeltas.
    • Anything else that doesn't decode throws SegmentRefusedError.
  • Signatures: every txn is checked against the segment's device keys (§2.8; the TxnVerifier port returns the first bad index, or null). One bad signature refuses the segment and nothing of its page is staged.
  • Txn-id dedupe (hard requirement): INSERT … _sync_seen_txn ON CONFLICT DO NOTHING per txn. A txn id already staged from the stream under any earlier segment is skipped and counted in redelivered, so a counter delta is never applied twice.
  • Stage and advance: stageTxns stages the rest (schema-ahead ones too; the applier holds them refused-schema), and the cursor advances in the same transaction.
  • Apply: applyStagedTxns applies, own echoes included (fast path and rebase, §3.5).

pullSyncStream(opts) (cloud/nexus-vault.ts):

  • refuses a store bound to another replica than the link (as push does, T13304);
  • a store with no pull position starts after the checkpoint it last synced (its genesis, or the snapshot it restored), verified and seeded by cursorFromCheckpoint;
  • the puller is journal.pull with the trusted signers, and the verifier tries each trusted key of the device.

Migration t12343-pull-cursor (local-only): _sync_cursor(stream PK, cursor_json, updated_at) and _sync_seen_txn(stream, txn PK, seq).

  • Both are classified local-only in both scopes, and pinned in table-classification-gate and nexus-attach.
  • Both are cleared from vault bundles, so a restored store resumes from its own checkpoint.

A note on $inc: SYNC_COUNTER_COLUMNS is empty today, since no counter-bearing table has joined the sync set yet (brain_patterns joins with T12894). So the re-delivery test proves the guard that protects $inc: a re-delivered txn is skipped at staging, never staged, applied or conflicted twice. A $inc-specific assertion lands with the first counter column.

Watermark (R7-7, T13256): not in this PR. It changes the segment plaintext format and belongs with T13256.

Tests

pull.test.ts (5), two real stores (A authors signed segments with the real builder; a fake stream serves them; B pulls):

  1. Stage and apply in stream order, with the cursor persisted; a rerun pulls nothing.
  2. Re-delivery: the same txn in a later segment is skipped (staged 0, redelivered 1), with one inbox row, one task and no conflict.
  3. A vault delta is passed over, and the ledger segment after it is applied.
  4. A txn signed by another device refuses the segment: nothing staged, no cursor.
  5. A body that is neither is refused.

nexus-vault.test.ts, "S5-1", against the fake server: enable push, write, push, then pull. The pull starts after the genesis checkpoint, verifies and stages the own echo, and sequences it (one _sync_sequenced row). A rerun pulls nothing.

Mutations (4, all killed): dedupe off, the vault-delta skip off, verification off, the cursor not persisted.

Checks

  • core tsc;
  • tests: table-classification-gate, nexus-attach, pull, push, nexus-vault, sync-store and capture pass (224/224);
  • gates 28 (pull.ts exemptions), 36 (--base: forward files checked, released files unchanged), 37 (--base), 40, 4, and 38 --base origin/feat/T12343-push;
  • skill coverage and lint-changesets.

🤖 Generated with Claude Code

kryptobaseddev and others added 2 commits October 6, 2026 18:15
…43 S5-1)

pullStream (store/sync/pull.ts) pages through the stream after the store's
persisted cursor and, per page in one transaction, decodes each segment
(a cleo-vault-delta/v1 segment carries no ops and is passed over),
verifies every transaction's device signature before writing anything,
skips transaction ids already staged from the stream (_sync_seen_txn: a
re-delivered transaction is never staged or applied twice), stages the
rest and advances _sync_cursor; then the applier applies, own echoes
included.

pullSyncStream (cloud/nexus-vault.ts) wires it to the journal client and
the trusted device keys, refuses a store bound to another replica than
the link, and starts a store with no position after the checkpoint it
last synced. New local-only migration t12343-pull-cursor; both tables are
cleared from vault bundles.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
pullStream ran regardless of the sync.pull flag; it now refuses (refused:
'sync.pull is off', nothing pulled or staged) unless the flag is on, as
push does for sync.push. PullStreamReport gains `refused` and its apply
report is null on a refusal.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…recedes a refused segment (T13306, T13307)

review-hotfix on #1959:
- T13306: a store whose last synced checkpoint is a vault snapshot from
  before the stream's journal genesis was seeded from it and pulled across
  the genesis, missing everything folded into it: silent divergence. It is
  now refused, nothing pulled, with the remedy to restore the journal
  checkpoint and join with `cleo sync enable push` (T12999); #1955's
  journaled-stream remedy says the same instead of "join by pulling".
- T13307: a refused segment threw before apply, stranding earlier staged
  pages and sticking every retry. The pull now stops at the refused
  segment: what precedes it (in the same page too) is staged and applied,
  the cursor stops just before it, and the report's `refused` names its
  seq, replica and device.
- Pruning: _sync_seen_txn rows first staged at or below the latest
  verified checkpoint's coversSeq (capped at the store's own position) are
  dropped at the start of each pull.
- Nit: with sync.pull off, pullSyncStream refuses before any network call.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…l floor exists (T13256)

A seen txn can return in a NEW segment above any checkpoint (a retired
replica's late segment, a rebind re-pushing a stored-but-unrecorded
upload); with its first-delivery row pruned by checkpoint coversSeq it
would stage and apply twice (an $inc column would double). pullSyncStream
no longer passes pruneSeenUpTo; pruneSeenTxns stays, documented as safe
only below a floor the applier refuses under (the receive watermark).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant