Replication
By default a queue partition lives on a single broker. That broker fsyncs every durable write, so the data survives a process restart, but it does not survive losing that node, and while the node is down the partition is unavailable.
Replication keeps copies of a partition on other brokers so a partition can survive and keep serving when its owner fails.
This is experimental and only active in Ganglion coordination mode. Standalone brokers own every queue and do not replicate.
Follow a failover for an animated walkthrough of owner loss, recovery proof and the return to service. The architecture stories also cover partition placement, delivery and agreed checkpoints.
What Fibril does
Section titled “What Fibril does”When coordination is enabled and a partition is assigned followers, each partition has one owner and one or more followers:
- Followers catch up by pulling. A follower worker reads the owner’s durable message and event records over the protocol and applies them durably to its own log. If a follower falls too far behind the owner’s retained log, it installs an owner checkpoint and resumes from there.
- Replicated failover waits for recovery proof. The controller records a proposed replacement and retains the previous assignment and replication source. Fresh Unix cluster queues enroll automatically; a worker seals replicas, verifies a source, installs the recovered state and activates a new quorum. Legacy histories have no supported migration into this recovery path. Heartbeat tails can suggest a candidate but cannot authorize it to serve.
- Replica-durable publishes wait for replicas. When the assignment’s durability policy requires more than the owner, a confirmed publish does not return until enough followers have reported the required progress, subject to a timeout and an in-sync floor. Each counted follower must have the queue payload batch and its exact enqueue-event frontier; the complete payload group is required when one enqueue record names several messages. Delivery visibility uses the same dependency proof.
Durability levels
Section titled “Durability levels”The assignment durability policy decides what a confirmed publish waits for. It
is set per cluster through coordination.ganglion.assignment_durability (see
configuration):
| Mode | A confirmed publish returns once… |
|---|---|
local_durable | the owner has durably written the append (the default, same as a single node) |
replica_accepted | N assigned nodes (including the owner) have accepted the append (weaker than fsync) |
replica_durable | N assigned nodes (including the owner) have durably written the append |
majority_durable | a durable majority of the assigned replica set has the append |
N includes the owner, so replica_durable with N = 2 means the owner plus
one durable follower.
Watch the message path
Section titled “Watch the message path”Follow a publish through the owner, persistence and the worker. Switch modes to see which waits overlap and which confirmation boundary remains. The badges separate current behavior, unmerged experiments and proposed work. Pause or scrub to inspect a frame. Save frame exports the current view for a presentation.
Follow one message. Watch the waits overlap.
Illustrative time units, not benchmark results. The replicated examples require owner + one durable follower. Payload and enqueue-event dependencies both matter. Polling/wakeup and batching delays are simplified.
Follow the story
Illustrative time units, not benchmark results. The replicated examples require owner + one durable follower. Payload and enqueue-event dependencies both matter. Polling/wakeup and batching delays are simplified.
Illustrative time units, not benchmark results. The replicated examples require owner + one durable follower. Payload and enqueue-event dependencies both matter. Polling/wakeup and batching delays are simplified.
Read the story, step by step
Durable replication · Current main
- A publish enters the owner
The owner must preserve both the payload and its enqueue event.
- Local persistence is in progress
Appending and fsyncing establish the owner’s durable boundary.
- The follower can read the durable batch
This animation shows the data response. The current transport is follower-pull.
- The follower persists its own copy
Both logs must complete and apply safely before reporting durable progress.
- Durable progress reaches the owner
The exact payload/enqueue dependency and assignment epoch are checked.
- The durability requirement is satisfied
Delivery and publisher confirmation can proceed independently.
- The worker acknowledges processing
Processing ACK and publisher confirmation are different signals.
Early replication · Unmerged experiment
- A publish enters the owner
The durable confirmation contract stays the same.
- Local I/O starts
The owner’s fsync is still pending.
- A written batch becomes readable early
The experimental path can begin replica work before local fsync finishes.
- Two persistence windows overlap
Overlap can remove serial waiting. Network, batching and progress-report costs remain.
- The follower has durable evidence
The owner still verifies both dependencies and its own durability.
- Confirm and deliver after the durable boundary
Early replication changes scheduling, not the required durable copies.
- The worker sends a processing ACK
This signal does not retroactively replace the durability policy.
Local speculation · Unmerged prototype · no replication
- A receiver has an open slot
The prototype preserves delivery ordering and has bounded speculative capacity.
- Local persistence starts
The data has not crossed the fsync boundary yet.
- The worker gets an early delivery
Internal speculative identity can help identify repeats. It is not exactly-once processing.
- Processing finishes before fsync
A processing ACK may satisfy the experimental confirmation contract.
- The publisher receives its confirmation
This path promises processing acknowledgement, not a durable surviving copy.
- The durability path completes later
Crash/error, fallback, expiry and duplicate cases remain adoption gates.
Replicated speculation · Proposed · correctness work pending
- A possible replicated path
This is design exploration, not available production behavior.
- Delivery and replica transfer begin early
An open receiver slot and ordering checks remain necessary.
- Persistence overlaps with processing
Early work increases concurrency without proving a safe confirmation shortcut.
- An ACK can precede the replica’s enqueue
Failover must preserve ACK/enqueue dependencies, identity and producer outcomes.
- Durable replica progress arrives
This conservative illustrated path still uses required durability.
- The publisher can be confirmed
A different ACK-or-durability contract requires separate proof before adoption.
Agreed recovery checkpoints
Section titled “Agreed recovery checkpoints”Unix cluster queues can periodically agree on a durable recovery snapshot. Set
runtime_seed.replication.agreed_checkpoint_interval_ms at startup, or change
replication.agreed_checkpoint_interval_ms in the dashboard’s cluster settings. Optional
agreed_checkpoint_max_events and agreed_checkpoint_max_bytes thresholds in the
same replication settings start attempts earlier under load. The interval must
remain enabled. Zero thresholds preserve periodic-only behavior. Starts are
limited to once per second per queue and periodic work is staggered across queues.
Byte counts are scheduling hints that reset on reopen, not durability evidence.
The Cluster page shows checkpoint age, event cut, local suffix and retained message
counts, agreement/install progress and failures. Retained counts are logical
ranges, not physical disk usage.
The default is 0. Positive intervals range from 1,000 to 86,400,000 ms.
Each admitted replica verifies and persists the same applied cut before consensus accepts it. Recovery verifies the retained snapshot and later event suffix, plus all required payloads. Large live backlogs still require payload reads. A missing replica delays replacement while ordinary replication continues.
Setting the interval to zero stops new attempts. Existing attempts finish and the last accepted checkpoint remains pinned. Retained disk can grow while replacement is delayed. This policy is opt-in pending broader workload and retention tuning. See the checkpoint plan.
Choose Missing replica and retry to follow delayed agreement, or Conflicting evidence to see why a new proposal must remain unaccepted. The previous accepted base stays in force until a valid replacement is accepted.
Agree on a base. Verify what comes next.
One queue, three admitted replicas. Event cut 120 is an illustrative exclusive boundary. Event positions and message offsets are different coordinates. This story does not measure time or storage size.
Follow the story
One queue, three admitted replicas. Event cut 120 is an illustrative exclusive boundary. Event positions and message offsets are different coordinates. This story does not measure time or storage size.
One queue, three admitted replicas. Event cut 120 is an illustrative exclusive boundary. Event positions and message offsets are different coordinates. This story does not measure time or storage size.
Read the story, step by step
Checkpoint agreement
- 01 / Capture an exact applied cut
A captures queue state after applying events before 120. The snapshot describes ready, inflight, delayed and settled state at that exact boundary. A snapshot file alone is not an accepted recovery checkpoint.
- 02 / Every admitted replica verifies
B and C reconstruct the same applied cut and verify the required evidence. Each persists its checkpoint and returns a bound receipt. A missing participant delays agreement while ordinary replication continues.
- 03 / Consensus accepts the checkpoint
Acceptance binds the verified base to the queue history and assignment. The accepted checkpoint stays pinned until a replacement is accepted. The agreed cut covers events before 120.
- 04 / New events form a suffix
Publishing, acknowledgements and timer transitions add events from 120 onward. Some messages created before the cut remain live. Their payloads are still needed, regardless of how old the checkpoint is.
- 05 / The owner disappears
The replacement remains fenced. Recovery obtains sealed evidence from the surviving replicas and checks that the accepted base is available and bound to the relevant history.
- 06 / Replay the suffix and validate live payloads
Recovery verifies the snapshot, reconstructs the later event suffix and checks required live payload identities and bytes. A large settled history can become cheap. A large live backlog still requires reads. Missing or incompatible evidence keeps the queue fenced.
- 07 / Evidence feeds the normal recovery barriers
The verified base and suffix reduce reconstruction work. They do not grant ownership. Installation, exact quorum receipts, consensus activation and local admission still have to succeed before the replacement serves.
Missing replica and retry
- 01 / Capture an exact applied cut
A captures queue state after applying events before 120. The snapshot describes ready, inflight, delayed and settled state at that exact boundary. A snapshot file alone is not an accepted recovery checkpoint.
- 02 / A required replica is unavailable
B has returned a durable receipt, but C is unavailable. The proposal waits without becoming the accepted base. Existing replication and queue traffic can continue when their own durability requirements are satisfied.
- 03 / Retry with fresh verification
C returns. A later attempt obtains verification for the exact proposed history, cut and current assignment. Old receipts are not assumed valid across changed identity or authority. Every required participant must pass before consensus can accept a checkpoint.
- 04 / Every admitted replica verifies
B and C reconstruct the same applied cut and verify the required evidence. Each persists its checkpoint and returns a bound receipt. A missing participant delays agreement while ordinary replication continues.
- 05 / Consensus accepts the checkpoint
Acceptance binds the verified base to the queue history and assignment. The accepted checkpoint stays pinned until a replacement is accepted. The agreed cut covers events before 120.
- 06 / New events form a suffix
Publishing, acknowledgements and timer transitions add events from 120 onward. Some messages created before the cut remain live. Their payloads are still needed, regardless of how old the checkpoint is.
- 07 / The owner disappears
The replacement remains fenced. Recovery obtains sealed evidence from the surviving replicas and checks that the accepted base is available and bound to the relevant history.
- 08 / Replay the suffix and validate live payloads
Recovery verifies the snapshot, reconstructs the later event suffix and checks required live payload identities and bytes. A large settled history can become cheap. A large live backlog still requires reads. Missing or incompatible evidence keeps the queue fenced.
- 09 / Evidence feeds the normal recovery barriers
The verified base and suffix reduce reconstruction work. They do not grant ownership. Installation, exact quorum receipts, consensus activation and local admission still have to succeed before the replacement serves.
Conflicting evidence
- 01 / Capture an exact applied cut
A captures queue state after applying events before 120. The snapshot describes ready, inflight, delayed and settled state at that exact boundary. A snapshot file alone is not an accepted recovery checkpoint.
- 02 / Evidence does not match
C returns evidence that does not match the proposed applied cut. The new checkpoint cannot be accepted. A receipt from B cannot substitute for C's required verification.
- 03 / Preserve the previous recovery base
The previous accepted checkpoint, if any, remains in force. Repeating the same contradictory evidence cannot make this proposal valid. Investigate the mismatch or repair the replica through an authorized path. This failed checkpoint attempt does not itself activate a replacement or imply that ordinary queue traffic has stopped.
Eager failover
Section titled “Eager failover”The default detector uses broker heartbeat expiry. An optional policy lets the active metadata controller exclude a peer from placement after repeated explicit Raft connection failures, with failed reconnects spanning a configurable grace. Enable it in the dashboard’s Settings → Replication section, or seed a fresh cluster with:
[runtime_seed.replication]eager_failover = trueeager_failover_grace_ms = 1000eager_failover defaults to false; the grace defaults to 1,000 ms and accepts
0–60,000 ms. The saved cluster runtime-settings document takes precedence over
startup seeds. Keep seeds consistent across nodes. Runtime policy changes restart
pending suspicion; a newly elected controller evaluates its own transport observations.
A successful Raft RPC, a fresh broker heartbeat or a changed broker process
identity resets suspicion. Client disconnects and replication-stream restarts do
not trigger it. Timeouts and silent packet loss continue to use heartbeat expiry.
Explicit RPC disconnects trigger one immediate retry, bounded to 200 ms. Setting
eager_failover_grace_ms = 0 also monitors idle Raft connections on the active
controller and checks reconnection immediately after an explicit close. An explicit
close followed by a refused reconnect can trigger placement without waiting for
another heartbeat. A successful dial alone does not establish Raft health, and a
probe timeout supplies no second explicit-failure observation. Positive grace
values still require failed contact across that interval. Raft heartbeat/election
timing and broker heartbeat/TTL settings remain unchanged.
Excluding a peer starts the existing placement/recovery process. Verified history, fencing and the configured confirmation threshold still gate activation. A crash that also removes the metadata leader first needs a Raft election; recovery may then dominate the outage. Shorter detection therefore does not promise service restoration within the grace interval. Network interruptions can also cause extra recovery work, so eager detection remains opt-in and experimental.
The controller logs the peer, error kind, failed attempts and elapsed suspicion
time without payloads. /admin/api/topology exposes current exclusions in
consensus.controller.eager_suspects.
Recovery observations
Section titled “Recovery observations”The Cluster dashboard records the recovery worker’s stages, peers and outcomes. Overlapping bars show concurrent work; failure detection and client reconnect happen outside this timing window. Select the failed or cancelled attempts below to inspect earlier retries of the same transition. Grey phases have no recorded stages; a retry can reuse work or stop before reaching them. These are illustrative timings, and the live view holds bounded history for the current process.
In-sync replicas
Section titled “In-sync replicas”A follower counts as in sync when it has reported durable progress recently
enough (runtime_seed.replication.isr_timeout_ms). Two runtime settings gate
replica-durable acceptance:
min_in_sync_replicasis a floor. When fewer replicas are recently in sync than the floor, replica-durable publishes fail fast with a clear error instead of blocking until timeout.1disables the floor.confirm_timeout_msbounds how long a replica-durable confirm waits before failing.
The follower read budget and poll intervals
(runtime_seed.replication.*) tune catch-up throughput versus idle confirm
latency. See the configuration replication settings.
Follower retention and repair
Section titled “Follower retention and repair”Owner compaction respects each active follower’s reported durable message and event positions. An unknown cursor starts at zero. A follower without fresh progress reports has a 60-second retention grace, so a disconnected follower does not pin that prefix indefinitely. Idle followers with fresh reports remain protected. These are catch-up protections and do not change confirmation counts.
If a follower needs records the owner has already compacted, consensus suspends its acknowledgement role and the owner invalidates its old replication sessions. The follower builds a replacement in a separate storage generation. Its original copy remains available as recovery evidence throughout copying. The owner gives repair a bounded 120-second retention grace while the replacement catches up.
The switch requires durable payloads, completely applied events and coverage of both original log tails. It publishes a durable route only after these checks. Re-admission uses a fresh transport identity. A lost admission reply resumes from the installed completion receipt without resetting the replacement. A recovery seal prevents a concurrent repair switch.
A healthy majority can continue serving during repair. A policy requiring every replica continues waiting for every required replica. Repair does not lower that threshold. A slow copy can exceed its retention grace and retry from a newer checkpoint. Retained original generations consume disk space, and automatic reclamation and byte-based retention budgets remain future work. This path applies to enrolled queues with matching broker revisions, with directory durability support on Unix. Stream recovery remains separate.
Activation and conditions
Section titled “Activation and conditions”-
Internal reads, writes, checkpoint operations and streaming controls require the
@nodeprincipal authenticated on the current connection. Ordinary users cannot access them, including when a logical session resumes. Production peers use the configured cluster secret and authenticate again when reconnecting. -
Replication requires Ganglion coordination mode and a follower target (
coordination.ganglion.target_followersgreater than zero). It is a cluster-level placement decision, not a per-queue client option. -
A partition replicates only after the controller has assigned it followers.
-
local_durablequeues behave exactly like single-node queues even in a cluster. Replica-durable confirms only mean something when the durability policy requires more than the owner. -
Follower application is durable: message and event writes overlap, and both completions drain before queue state is applied and progress is reported. Failed or interrupted application blocks promotion until recovery or resync.
-
Progress reports are tied to an assignment epoch and an ordered transport session. Replacement sessions and assignments invalidate earlier reports; resets replace both reported frontiers together.
-
Run matching broker revisions across a cluster. Older reports without an assignment epoch remain readable but do not count toward replicated confirms. The client publish, delivery and acknowledgment frames are unchanged.
Tradeoffs and limits
Section titled “Tradeoffs and limits”- Replica-durable confirms add latency: a publish waits for follower progress, bounded by the follower poll interval and the confirm timeout.
- This surface is experimental. A follower’s local tails can omit a batch confirmed by the old owner and another follower. Recovery therefore requires accepted-history witnesses and fencing before selecting and installing a source. Enrolled queues retain their configured confirmation threshold after recovery; unsupported or insufficient evidence leaves the partition fenced. Checkpoint backfill is checked before promotion and after restart. Unix checkpoint installation uses a durable journal to resume interrupted log/state replacement before ordinary replay; completed retries preserve later backfill. Linux fault and process-kill tests cover this local path. Non-Unix checkpoint installation is unsupported pending durable metadata support.
- The controller persists pending recovery metadata for owner replacement and follower-set or durability-policy changes involving replicated confirmation. Existing healthy owners retain their active configuration until recovery can safely activate the proposed assignment. Recovery of enrolled queues requires a complete source and an available proposed owner; replacing an unavailable proposed owner and reconstructing crossed histories remain rollout gates. Pending requests appear in the admin topology’s controller status. Do not discard a surviving replica or clear a request to bypass recovery proof.
- Cross-broker replication-lag aggregation into a single cluster view is still pending. A broker’s own follower workers and their progress are visible on the admin queues page.
See also
Section titled “See also”- Clustering for ownership, epochs, and coordination modes.
- Recovery sealing for explicit seal authorization, retained identity and remaining recovery gates.
- Recovery quarantine for how a node handles a damaged log on restart.
- Configuration for the replication and durability settings.
- Project status and implemented surface for what is wired.