Skip to content

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.

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.

The assignment durability policy decides what a confirmed publish waits for. It is set per cluster through coordination.ganglion.assignment_durability (see configuration):

ModeA confirmed publish returns once…
local_durablethe owner has durably written the append (the default, same as a single node)
replica_acceptedN assigned nodes (including the owner) have accepted the append (weaker than fsync)
replica_durableN assigned nodes (including the owner) have durably written the append
majority_durablea 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.

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.

INSIDE FIBRIL / THE MESSAGE PATHCurrent main

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.

01

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

  1. A publish enters the owner

    The owner must preserve both the payload and its enqueue event.

  2. Local persistence is in progress

    Appending and fsyncing establish the owner’s durable boundary.

  3. The follower can read the durable batch

    This animation shows the data response. The current transport is follower-pull.

  4. The follower persists its own copy

    Both logs must complete and apply safely before reporting durable progress.

  5. Durable progress reaches the owner

    The exact payload/enqueue dependency and assignment epoch are checked.

  6. The durability requirement is satisfied

    Delivery and publisher confirmation can proceed independently.

  7. The worker acknowledges processing

    Processing ACK and publisher confirmation are different signals.

Early replication · Unmerged experiment

  1. A publish enters the owner

    The durable confirmation contract stays the same.

  2. Local I/O starts

    The owner’s fsync is still pending.

  3. A written batch becomes readable early

    The experimental path can begin replica work before local fsync finishes.

  4. Two persistence windows overlap

    Overlap can remove serial waiting. Network, batching and progress-report costs remain.

  5. The follower has durable evidence

    The owner still verifies both dependencies and its own durability.

  6. Confirm and deliver after the durable boundary

    Early replication changes scheduling, not the required durable copies.

  7. The worker sends a processing ACK

    This signal does not retroactively replace the durability policy.

Local speculation · Unmerged prototype · no replication

  1. A receiver has an open slot

    The prototype preserves delivery ordering and has bounded speculative capacity.

  2. Local persistence starts

    The data has not crossed the fsync boundary yet.

  3. The worker gets an early delivery

    Internal speculative identity can help identify repeats. It is not exactly-once processing.

  4. Processing finishes before fsync

    A processing ACK may satisfy the experimental confirmation contract.

  5. The publisher receives its confirmation

    This path promises processing acknowledgement, not a durable surviving copy.

  6. The durability path completes later

    Crash/error, fallback, expiry and duplicate cases remain adoption gates.

Replicated speculation · Proposed · correctness work pending

  1. A possible replicated path

    This is design exploration, not available production behavior.

  2. Delivery and replica transfer begin early

    An open receiver slot and ordering checks remain necessary.

  3. Persistence overlaps with processing

    Early work increases concurrency without proving a safe confirmation shortcut.

  4. An ACK can precede the replica’s enqueue

    Failover must preserve ACK/enqueue dependencies, identity and producer outcomes.

  5. Durable replica progress arrives

    This conservative illustrated path still uses required durability.

  6. The publisher can be confirmed

    A different ACK-or-durability contract requires separate proof before adoption.

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.

INSIDE FIBRIL / AN AGREED BASEImplemented · opt-in checkpoints

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.

01

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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 = true
eager_failover_grace_ms = 1000

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

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.

Recovery stages and overlapping work · read-only sample · fixed view Open full dashboard ↗

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_replicas is 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. 1 disables the floor.
  • confirm_timeout_ms bounds 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.

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.

  • Internal reads, writes, checkpoint operations and streaming controls require the @node principal 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_followers greater 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_durable queues 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.

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