Repository navigation
fix(replication): coordinator local legs reach the replica instead of only advancing the offset - #859
Conversation
… master_repl_offset (#815) On a multi-shard master the in-process leg of a multi-key write -- the MSET/MSETNX slice, DEL/UNLINK, COPY or BITOP result that the connection's own shard owns -- was persisted by persist_local_leg, which did exactly two things: issue_append_lsn and send_append_group. It never appended to the replication backlog and never queued a ReplicaLiveFanout, so the leg never reached a replica. issue_append_lsn still advanced shard_offsets[shard] and master_repl_offset for it, the same atomics increment_shard_offset bumps on the streamed paths, so the master counted bytes no replica could ever ACK. wait_for_replicas targets total_offset(), so after the first local leg WAIT answered :0 forever with --appendonly yes. With --appendonly no the leg short-circuited on the missing pool and neither counted nor shipped: WAIT answered :1 while the replica silently diverged. The fix gives the local leg the two-plane contract every other write already has (handler is_write block, blocking-pop path, wal_append_and_fanout): when the replication plane is live, record to the backlog and advance the offset synchronously before the first .await -- fused SELECT <db> prefix on num_shards > 1, emit-on-change at 1 -- queue the live delivery on the shard self-queue, and hand the AOF leg lsn = 0 so the offset is not advanced twice; otherwise advance through the AOF LSN exactly as before, so a primary that never had a replica attach is unchanged. run_local_persist's "no AOF pool" early return now also requires no live replica. The multi-shard SWAPDB leg, which issued an AOF LSN and then emitted a replication record through record_local_write_global, was counting itself twice; its AOF leg now carries lsn = 0 as well. record_local_write / record_local_write_db (handler_monoio::ft) and the ctx-free record_local_write_global / record_local_write_db_global now share one implementation, replication::state::record_local_write_on / record_local_write_db_on, plus fanout_active_for as the ctx-free twin of replication_fanout_active. The global db-aware twin gains the multi-shard branch it previously documented as not implemented. Red/green: tests/replication_local_leg_815.rs (#[ignore]d like the other replication suites) runs a real primary at --shards 1 and --shards 4, both AOF modes, plus a single-shard replica. On pristine main (ae6cd00) the --shards 4 legs fail: appendonly=yes "WAIT after coordinator local legs: left :0 right :1"; appendonly=no "10 replica divergences" (5 of 24 groups, the connection's own shard's share). On this tree all four pass, and master_repl_offset equals the replica's acked offset. Three unit tests cover the new state.rs helpers. Refs: #815 author: Tin Dang
|
ⓘ Qodo reviews are paused because the subscription is no longer active. Ask your workspace admin to reactivate the subscription to resume reviews. Manage billing |
|
Warning Review limit reachedNext included review available in 37 minutes. View limit detailsLimit details: You’ve used the included review currently available. You've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository. Review configuration: ⚙️ Run configurationConfiguration used: defaults Review profile: CHILL Plan: Advanced Run ID: 📒 Files selected for processing (5)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
@coderabbitai review |
|
Fixes #815.
The issue reports that
WAITcan never be satisfied. That is the visible symptom; the cause is worse — the replica is silently missing data.Root cause (pre-fix line numbers)
(a) The local leg advances the offset without sending anything.
coordinator.rs:577persist_local_leg→AofWriterPool::issue_append_lsn(aof/pool.rs:986) →ReplicationState::issue_lsn(replication/state.rs:262), bumpingshard_offsets[shard]andmaster_repl_offsetwith no backlog append and noReplicaLiveFanout. Seven call sites.(b) A streamed write advances it via
spsc_handler.rs:4447(increment_shard_offset, AOF withlsn = 0) andhandler_monoio/ft.rs:92.(c)
WAIT(replication/master.rs:931) comparesΣ ack_offsetsagainstmaster_repl_offset— which now counts bytes no replica can ever acknowledge.Same class, found on the way: multi-shard
SWAPDBwas double-counted — it issued an AOF LSN (:3650) and emitted a replication record (:3696).Semantics: replicate the leg (A), not stop counting it (B)
SELECT-framedWAITreturns:1while the replica is missing data — a hang becomes silent divergence;--appendonly nostays brokenA is what the repo's own contracts require: #285 (every state writer feeds both planes), the write-time offset advance the inline PSYNC snapshot cut depends on, and R2's exactly-once as an offset cut (so the offset must equal delivered bytes). B contradicts all three.
Offset accounting for a primary with no replica is unchanged — the AOF LSN still advances it — so the PSYNC baseline (
Σ shard_offsets) is untouched.Fix
persist_local_legnow asksfanout_active_for(repl_state). If a replica is live:record_local_write_db_on(backlog append + offset advance, synchronous before any.await, fusedSELECT <db>prefix whennum_shards > 1, live delivery via the self-queue) and the AOF leg takeslsn = 0. Otherwiseissue_append_lsnexactly as before.run_local_persist's no-AOF early exit also requires no live replica.SWAPDBAOFlsn = 0— the replication record it already emits owns the advance.record_local_write/record_local_write_dband their ctx-free_globaltwins now share one implementation inreplication::state; the privateserialize_selectcopy is gone (aof::serialize_select_recordproduces identical bytes).Composes with #852 without touching it: local-leg payloads are deterministic (
MSET/MSETNX/DEL/UNLINKverbatim,COPY/BITOPas synthesizedSET dest/DEL dest) — no relative TTL, so noeffect_rewriteseam is needed. The sameserializedbuffer feeds the AOF and the replica (oneBytesrefcount clone).Red/green
tests/replication_local_leg_815.rs,#[ignore]d per repo convention:MOON_BIN=<bin> cargo test --test replication_local_leg_815 -- --ignored --test-threads=1RED on
ae6cd003(binary built pre-edit, 231 crates):GREEN on this branch (
Compiling moon, 35s, distinct sha): 4/4 at shards 1 and 4, both AOF modes,master_repl_offset == slave0 offset.--shards 1is green on main because the coordinator is gated out there — a harness control, and the reason shards=1 alone proves nothing (the self-SPSC gap). The affected key set shifts between runs (t4/9/12/18/20 vs t3/5/6/11/13) because it follows the listener's shard assignment, not test flakiness.Gates
cargo check --all-targetsmonoio rc=0 and tokio rc=0;clippy --all-targets -D warningsrc=0;fmt --checkrc=0;replication::stateunit tests 22/22 ×5 runs.Also fixed: a pre-existing test race
test_ensure_backlogs_allocated_does_not_activate_fanout_hintreads the globalFANOUT_HINTtwice unlocked while a sibling test flips it. Latent on main; the new tests made it likelier, so it is closed here with a test-moduleHINT_LOCKserializing the four hint-touching tests. Same family as #856 and #822.Not done, and why
No kill -9 / parity / VM A/B: the AOF bytes per leg are unchanged — only the LSN tag becomes 0 when a replica is live, which every streamed write already does. No benchmark; one
Relaxedload is added on the no-replica path.Landmine worth knowing
The tokio path is unchanged only because a tokio master refuses PSYNC. When that changes,
run_local_persistreplicates through this same path — which is precisely what #815 warned about.