Skip to content

fix(replication): coordinator local legs reach the replica instead of only advancing the offset - #859

Merged
TinDang97 merged 1 commit into
mainfrom
fix/815-local-leg-replication
Sep 8, 2026
Merged

TinDang97 merged 1 commit into
mainfrom
fix/815-local-leg-replication

Conversation

@TinDang97

Copy link
Copy Markdown
Collaborator

Fixes #815.

The issue reports that WAIT can 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:577 persist_local_leg → AofWriterPool::issue_append_lsn (aof/pool.rs:986) → ReplicationState::issue_lsn (replication/state.rs:262), bumping shard_offsets[shard] and master_repl_offset with no backlog append and no ReplicaLiveFanout. Seven call sites.

(b) A streamed write advances it via spsc_handler.rs:4447 (increment_shard_offset, AOF with lsn = 0) and handler_monoio/ft.rs:92.

(c) WAIT (replication/master.rs:931) compares Σ ack_offsets against master_repl_offset — which now counts bytes no replica can ever acknowledge.

Same class, found on the way: multi-shard SWAPDB was double-counted — it issued an AOF LSN (:3650) and emitted a replication record (:3696).

Semantics: replicate the leg (A), not stop counting it (B)

replicas receive primary counts breaks
A (chosen) the MSET/DEL/COPY/BITOP slice, SELECT-framed bytes actually sent nothing that was not already broken
B nothing (unchanged) streamed bytes only WAIT returns :1 while the replica is missing data — a hang becomes silent divergence; --appendonly no stays broken

A 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_leg now asks fanout_active_for(repl_state). If a replica is live: record_local_write_db_on (backlog append + offset advance, synchronous before any .await, fused SELECT <db> prefix when num_shards > 1, live delivery via the self-queue) and the AOF leg takes lsn = 0. Otherwise issue_append_lsn exactly as before. run_local_persist's no-AOF early exit also requires no live replica. SWAPDB AOF lsn = 0 — the replication record it already emits owns the advance.

record_local_write/record_local_write_db and their ctx-free _global twins now share one implementation in replication::state; the private serialize_select copy is gone (aof::serialize_select_record produces identical bytes).

Composes with #852 without touching it: local-leg payloads are deterministic (MSET/MSETNX/DEL/UNLINK verbatim, COPY/BITOP as synthesized SET dest/DEL dest) — no relative TTL, so no effect_rewrite seam is needed. The same serialized buffer feeds the AOF and the replica (one Bytes refcount 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=1

RED on ae6cd003 (binary built pre-edit, 231 crates):

--shards 4 --appendonly yes  -> WAIT after coordinator local legs: left: ":0" right: ":1"
--shards 4 --appendonly no   -> 15 replica divergences:
  {t3}:a (co-located MSET)  {t3}:c (DEL not applied)  {t3}:s (scattered MSET slice)
  {t5}:a ...                {t5}:c ...                {t5}:s ...

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 1 is 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-targets monoio rc=0 and tokio rc=0; clippy --all-targets -D warnings rc=0; fmt --check rc=0; replication::state unit tests 22/22 ×5 runs.

Also fixed: a pre-existing test race

test_ensure_backlogs_allocated_does_not_activate_fanout_hint reads the global FANOUT_HINT twice unlocked while a sibling test flips it. Latent on main; the new tests made it likelier, so it is closed here with a test-module HINT_LOCK serializing 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 Relaxed load 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_persist replicates through this same path — which is precisely what #815 warned about.

… 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-code-review

Copy link
Copy Markdown

ⓘ Qodo reviews are paused because the subscription is no longer active. Ask your workspace admin to reactivate the subscription to resume reviews. Manage billing

@coderabbitai

coderabbitai Bot commented Sep 7, 2026 •

Copy link
Copy Markdown

Warning

Review limit reached

Next included review available in 37 minutes.

Check out review usage here.

View limit details

Limit 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.

Learn how review limits work.

Review configuration:

⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Advanced

Run ID: 274e92f3-16e0-445a-aa7b-c143e2532f43

📥 Commits

Reviewing files that changed from the base of the PR and between ae6cd00 and 2f9f0d0.

📒 Files selected for processing (5)
  • CHANGELOG.md
  • src/replication/state.rs
  • src/server/conn/handler_monoio/ft.rs
  • src/shard/coordinator.rs
  • tests/replication_local_leg_815.rs

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@TinDang97

Copy link
Copy Markdown
Collaborator Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Sep 8, 2026 •

Copy link
Copy Markdown
⚠️ Action not completed

Review rate limited.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@TinDang97
TinDang97 merged commit 1f919d1 into main Sep 8, 2026
10 checks passed
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.

Coordinator local legs are never replicated, and still inflate master_repl_offset (WAIT breaks permanently)

1 participant