Skip to content

fix(crdt): redeliver late-stamped merges and hub-pulled rows to every peer - #37

Merged
juicycleff merged 6 commits into
mainfrom
fix/crdt-late-stamp-main
Oct 8, 2026
Merged

juicycleff merged 6 commits into
mainfrom
fix/crdt-late-stamp-main

Conversation

@juicycleff

Copy link
Copy Markdown
Contributor

This puts the late-stamp fix on main. #36 was merged into fix/crdt-sync-defects after #35 had already been squashed into main, so none of it landed. This branch carries #36's two commits, rebased onto main, plus four more that came out of reviewing it.

It fixes edits that never reach other devices. Say device A goes offline, bumps a counter, and comes back a day later. The server merges A's change fine, but B already pulled past A's old timestamp, so B never sees the new total. In the Dart convergence test B sat at 7 when it should have been 10. With this branch it reaches 10.

Running it

go test -race ./crdt/...

There are fifteen new tests. They cover late stamps for sets, lists, text, documents and counters, an LWW loser and winner, tombstones, cursor monotonicity, the multi-table window from #35, the stream, time travel on a restamped row, a counter touched by many nodes, and a hub's Syncer. Every one that checks a fix fails without it. The Dart client in #34 passes its full conformance suite against this branch (37 tests on SQLite), with the skipped late-stamp test turned on.

What was wrong

The pull cursor pages by each shadow row's stored HLC. When a push merged a late-stamped change into a row, the state changed but the row kept its older HLC, so any peer whose cursor was already past it skipped the row for good. That hit counters, sets, lists, text and documents. LWW mostly escaped, because an older write loses anyway.

Counters had a second problem. A counter row merged from another node was pulled back as that row's own node delta only, so a peer pulling from scratch still undercounted.

The Go Syncer had the same bug, and worse. A hub that pulls from upstream with a Syncer and serves its own clients wrote each pulled row at the row's own clock. Ordinary polling lag puts that clock behind the clients' cursors, so a hub's clients could miss almost any upstream row, late stamp or not.

What changed

A shadow row now keeps its own clock (in crdt_state, which is what pulls deliver) apart from its cursor position (hlc_ts, hlc_counter). When a push or a Syncer merge changes a row's stored state, the server restamps the cursor with a fresh server-clock HLC, so every earlier cursor sorts before it and the row gets redelivered. The row's own clock doesn't move, so LWW and tombstone results come out the same on every replica. A merge that changes nothing writes nothing: an LWW loser, an older tombstone, a retried push.

A counter row is pulled as one record. It carries the full counter state in state, the way sets, lists, text and documents already do, plus counter_delta for the row's own node. Go, crdt-js and Dart all merge state before they look at the delta, so you get every node's totals in a single record, however many devices have touched the counter.

ReadStateAt and ReadFieldHistory filter on each row's own clock, not its cursor position. Without that, time travel lost the whole state of any row that had been restamped.

Pushes and Syncer merges take one lock per plugin around the read, merge and write. Each push reads a table's highest stored position once, not once per row.

Things you should know before merging

  • A client that ignores state sees only the row node's counter totals. That's what it saw before this PR. All three clients in this repo read state.
  • Rolling back is lossy for LWW fields. An older binary reads a restamped row's cursor position as its clock. If a winning value stamped T2 was restamped to T10 and someone later wrote at T6, the old binary keeps the T2 value, thinking it was written at T10, and redelivers it with that clock, so the T6 write is lost on every replica.
  • The lock rules out a lost update between two pushes, or a push and a Syncer merge, on one server. It doesn't cover local writes through AfterMutation, which still read and write the server's own rows without it. That predates this PR.
  • BeforeMerge, AfterMerge and BeforeMetadataWrite run under the lock, so a slow one stalls every push on the plugin. AfterMetadataWrite and the sync hooks run outside it. A hook that calls back into HandlePush or Syncer.Sync on the same plugin would deadlock.
  • If several instances share one database, a peer can still miss a row when two instances commit out of order. Reading the table's highest position on each push keeps an instance from landing behind the other's committed rows, but closing the race fully needs a database-side sequence.
  • Rows merged from late stamps before this ships aren't restamped. Peers get them on the next change to the row, or on a full resync.
  • Change HLCs now sit at or below cursor positions. Clients should resume from LatestHLC. One that resumes strictly from the highest change HLC and loops until empty would keep pulling the newest restamped row. Dart and crdt-js already do the right thing.
  • The new SQL has run against the in-memory fake and SQLite. It hasn't run on Postgres yet.

…ll counters per node

A shadow row now has two clocks. The state's own clock, stored in crdt_state, is what merges and LWW compare and what a pull hands out as the change HLC. The hlc_ts and hlc_counter columns are the cursor position pulls page by. Every writer so far kept them equal, and rows written that way read back exactly as before, so nothing moves until a writer asks for a different position through WriteFieldStateAt or WriteTombstoneAt.

Tombstone rows store their delete clock in crdt_state as well, and ReadState keeps the latest delete clock across nodes, not whichever row it happened to read last.

A counter row is pulled as one change per node it holds. The server merges an increment into whichever node's row has the newer clock, so a pull that only sent the row's own node dropped the other node's totals. That is why the Dart convergence test left device B at 7.
…ps reach every peer

A device that was offline pushes changes stamped hours ago. The server merged them, but the row kept the old clock as its cursor position, so a peer that had already pulled past that clock never got the merged state.

Now a push that changes a row's stored state moves the row to a fresh cursor position: past the plugin clock, past every position this process handed out, past the highest one in the table, and past the row's own clock. Only the position moves. The change keeps its own clock on the wire, so you get the record you would have got by pulling earlier, and every client picks the winner it would have picked then. An LWW value that loses writes nothing. Neither does an older tombstone or a retried push.

Pulls and streams resume from cursor positions, including the multi-table window from #35. Allocation and the write share one lock per plugin, which also stops two concurrent pushes to one field from losing each other's merge.
…sition

ReadStateAt filtered rows on hlc_ts and hlc_counter. Since a push now
restamps the position of every row it changes, a row merged at T+10s with
its own clock at T+1s vanished from every cut before T+10s, taking all of
its field state with it. You'd ask for the record at T+5s and get no views
field at all.

We now read the record's rows and keep those whose own clock (the field
state HLC, or the tombstone's delete clock) is at or before the cut. That
is what the cut meant before restamping existed.

The in-memory shadow fake now honours a position cut when a query has one,
so the test fails on the old query instead of passing by accident.

ReadFieldHistory gets a doc note: its since, order and limit use positions
while each entry reports the row's own clock.
A counter row went out as one record per node it held. One +2 on a
counter that 500 devices had touched sent 500 records (about 89 KB) to
every client on its next pull, and the Go Syncer paid a read and a write
for each of them.

Now a counter row is one record with State set to the stored
PNCounterState, the same carrier sets, lists, text and documents already
use. Go ApplyChange, crdt-js mergeFieldState and Dart applyChange all
merge a state carrier before they look at counter_delta, so every current
client still gets every node's totals and still converges after a late
increment. The record keeps counter_delta for the row's own node, which is
what a client that predates state carriers saw before.
…lock

Three smaller fixes to the push path, all under the cursor write lock.

Each changed row used to cost a maxCursor query on top of its read and its
upsert, serialized across every client by one mutex. A push now reads each
table's stored maximum once (cursorBatch). Within one process the
allocator's high mark already covers every position it handed out, so the
stored maximum only matters after a restart or for another instance, and
reading it once per push still keeps a restarted server past its old rows.

AfterMetadataWrite ran under the lock because of a deferred unlock. It now
runs after the lock is released, so a slow hook stalls only its own push.

A record delete read the whole record first, so a field row that failed to
decode failed the delete. It now reads only the record's tombstone rows.

The cursorAllocator comment now says what the lock does not cover: local
writes through AfterMutation (the lost-update guarantee is push against
push only), other instances sharing the database, and which plugin hooks
still run under it.
A hub pulls from upstream with a Syncer and serves its own clients with a
SyncController. The Syncer wrote each pulled row at its own clock, and
ordinary polling lag puts that clock behind the hub clients' cursors, so
on a hub almost any upstream row could be missed by its clients. This was
true before restamping too; it is the same bug as the late stamp.

mergeRemoteChange now does what a push does. It takes the plugin's cursor
write lock around the read, merge and write, writes nothing when the merge
leaves the state unchanged, and writes a changed row at a fresh position
from the shared allocator. A delete is written only when it is newer than
the stored one, also at a fresh position. The AfterInboundChange hook runs
after the lock is released.
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