Skip to content

Module replication Architecture

github-actions[bot] edited this page Sep 28, 2026 · 25 revisions

Navigation: Home > Modules

Architecture - Replication Module

Overview

The replication module composes replication orchestration, consensus/failover behavior, conflict resolution, logical replication/CDC streaming, and replication observability into a bounded high-availability subsystem.

Main Execution Planes

  1. Core orchestration plane
  • replication manager lifecycle and mode control
  • leader promotion/failover and topology management
  1. Data propagation and conflict plane
  • WAL/logical propagation and slot/event stream behavior
  • HLC/LWW/CRDT conflict detection and merge behavior
  1. Observability and policy plane
  • lag/health/topology diagnostics and export behavior
  • replication policy validation and assignment behavior

Core Contracts

Contract Behavior
replication contract deterministic init/replicate/promote semantics
consensus contract explicit election/promotion transitions
conflict contract deterministic conflict resolver outcomes per strategy
observability contract explicit lag/health/topology visibility

Failure Semantics

  • initialization and promotion failures are explicit.
  • slot/stream/CDC path faults surface deterministic outcomes.
  • conflict-resolution edge cases remain explicit and non-silent.
  • degraded replica lag/health is observable via module surfaces.

Lock Hierarchy (Wave A Block 2 Hardening)

To prevent circular deadlocks and ensure bounded lock contention, the replication module implements a strict multi-level lock hierarchy across all components:

Level 1: Manager Collection Locks (Outermost)

  • Files: replication_slot.cpp, multi_tier_replication.cpp, logical_replication.cpp
  • Locks: slots_mutex_, collection_tiers_mutex_, slots_mutex_ (shared_mutex)
  • Scope: Slot/tier collection access (create/lookup/list)
  • Hold Time: MINIMAL (~microseconds for map operations only)
  • Pattern: Acquire shared/unique β†’ map access β†’ release β†’ I/O outside
  • Mutex Type: std::mutex or std::shared_mutex depending on read/write ratio

Level 2: Per-Resource State Locks

  • Files: replication_slot.cpp, raft_v2.cpp, event_stream.cpp, logical_replication.cpp
  • Locks: state_mutex_, config_mutex_, subs_mutex_, SlotRuntime::mutex
  • Scope: Individual slot/config state (pause/resume/query)
  • Hold Time: MINIMAL (~microseconds for state copy)
  • Pattern: Copy state under lock β†’ release β†’ I/O on copy
  • Mutex Type: std::mutex or std::lock_guard

Level 3: Background Worker and I/O Locks (Variable Hold Time)

  • Files: async_wal_shipper.cpp
  • Locks: queue_mutex_, callback_mutex_, stats_mutex_
  • Scope: Queues, callbacks, metrics (external I/O operations)
  • Hold Time: VARIABLE (depends on network, 10ms-1000ms typical)
  • Pattern: Acquire β†’ quick operation β†’ release β†’ invoke handler outside
  • Mutex Type: std::unique_lock with optional timeout support

Lock Ordering Diagram (ASCII)

Level 1 (Manager Collection)
  β”œβ”€β†’ slots_mutex_ (ReplicationSlotManager)
  β”œβ”€β†’ collection_tiers_mutex_ (MultiTierReplicationManager)
  └─→ slots_mutex_ (LogicalReplicationManager, shared_mutex)
      ↓
Level 2 (Per-Resource State)
  β”œβ”€β†’ state_mutex_ (ReplicationSlot)
  β”œβ”€β†’ config_mutex_ (RaftV2ClusterConfig)
  β”œβ”€β†’ subs_mutex_ (ReplicationEventStream)
  └─→ SlotRuntime::mutex (LogicalReplicationManager)
      ↓
Level 3 (Background I/O)
  β”œβ”€β†’ queue_mutex_ (AsyncWalShipper)
  β”œβ”€β†’ callback_mutex_ (AsyncWalShipper)
  └─→ stats_mutex_ (AsyncWalShipper)
      ↓
[Blocking Operations - NO LOCKS HELD]
  β”œβ”€β†’ File I/O (persist slots, state)
  β”œβ”€β†’ WAL append (wal_->append())
  β”œβ”€β†’ Network I/O (ship handler)
  └─→ Callbacks (event listeners)

Critical Invariants

  1. Acquire-Only Forward Pattern: Always acquire locks in increasing level order (1β†’2β†’3β†’I/O)
  2. No Backward Locks: Never acquire Level N lock while holding Level N-1 lock
  3. Lock-Free I/O: All blocking operations execute OUTSIDE all acquired locks
  4. State Copy Pattern: Copy mutable state while holding lock, release, then use copy for I/O
  5. Timeout Guards: Long-running operations (condition variable waits) use timeouts

Deadlock Prevention Examples

βœ… Safe Pattern (Lock-Free I/O)

bool ReplicationSlot::pause() {
    SlotState state_copy;
    {
        std::lock_guard<std::mutex> lock(state_mutex_);  // Level 2
        state_.status = SlotStatus::PAUSED;
        state_copy = state_;
    }  // ← Lock RELEASED
    persistStateImpl(state_copy);  // ← I/O happens OUTSIDE lock
    return true;
}

❌ Unsafe Pattern (I/O Under Lock) - FIXED in Wave A Block 2

// BEFORE (UNSAFE - Lock ordering violation):
MembershipChangeEntry MembershipChangeManager::writeEntry(...) {
    entry = ...;
    wal_->append(wal_entry);  // ← I/O UNDER lock_guard
    return entry;
}

// AFTER (SAFE - I/O outside lock):
auto entry = writeEntry(...);  // ← Create entry under lock
{
    std::lock_guard<std::mutex> lock(mutex_);  // ← Release before I/O
    // ...
}
wal_->append(wal_entry);  // ← I/O outside lock

Timeout Support

All long-running blocking operations use timeouts to ensure bounded wait times:

  • Condition Variable Waits: cv.wait_for(lock, timeout, predicate)
  • Default Timeout: 1-5 seconds depending on operation
  • Configuration: Via replication.timeout_ms and related config keys

Module Dependencies

Direct Upstream Dependencies (this module uses)

Module Interface / File Purpose
cdc include/replication/schema_cdc.h (via cdc/schema_registry) Schema change events consumed during logical replication
utils include/utils/ Utility helpers (encoding, error codes, logging)

Direct Downstream Consumers (modules that use this module)

Module Via Notes
server include/replication/replication_manager.h Replication admin and status API endpoints
storage include/replication/ (WAL shipping to replicas via CDC/WAL path) Storage WAL events propagated to replica nodes
failover include/replication/replication_manager.h, include/replication/raft_v2.h Failover module consults replication state for leader election; auto_failover_manager and disaster_recovery_manager both import replication_manager.h
temporal include/replication/multi_master_replication.h temporal_conflict_resolver imports multi-master replication state to resolve concurrent write conflicts across temporal branches (include/temporal/temporal_conflict_resolver.h:34)

Integration Points

Critical Integration: Replication β†’ CDC / Schema Registry

Files: src/replication/logical_replication.cpp ↔ include/cdc/schema_registry.h (via include/replication/schema_cdc.h) Contract: Logical replication subscribes to schema change events from CDC schema registry to maintain schema-aware replication slots. Schema version is embedded in each logical event for consumer compatibility validation. Thread Safety: Schema registry reads are concurrent-safe; logical replication uses shared_mutex (Level 1 in lock hierarchy) for slot access. Failure Mode: Schema registry unavailable β†’ logical replication pauses the affected slot and emits a lag alert; does not corrupt data.

Critical Integration: Replication β€” Raft V2 Leader Election

Files: src/replication/raft_v2.cpp ↔ include/replication/raft_v2.h Contract: RaftV2 manages leader election and log replication; leader promotion triggers replication_manager failover callbacks. Config changes require quorum before being applied. Thread Safety: Config state uses Level-2 config_mutex_; election state uses Level-2 per-resource mutex. I/O (WAL append) executes outside all locks per Wave A Block 2 hardening. Failure Mode: Loss of quorum β†’ cluster read-only; stale leader detection β†’ automatic step-down; WAL failure β†’ abort with explicit error (not silent).

Critical Integration: Replication β€” Async WAL Shipping

Files: src/replication/async_wal_shipper.cpp ↔ include/replication/async_wal_shipper.h Contract: WAL segments are queued and shipped to replica endpoints asynchronously. Queue is bounded; backpressure is applied on overflow. All callbacks and network I/O execute outside Level-3 locks. Thread Safety: Level-3 locks (queue_mutex_, callback_mutex_, stats_mutex_) guard queue and stats; callbacks invoked without any lock held. Failure Mode: Network failure β†’ retried with exponential backoff up to configured budget; queue overflow β†’ oldest entry dropped with sequence gap marker; lag alert fired.

Critical Integration: Replication β€” Conflict Resolution (Multi-Master)

Files: src/replication/conflict_resolution.cpp ↔ include/replication/conflict_resolution.h, include/replication/crdt_types.h Contract: On concurrent writes from multiple masters, conflict resolver applies strategy (HLC/LWW/CRDT) to produce deterministic merge outcome. Resolver is called synchronously on the apply path. Thread Safety: CRDT merge operations are stateless and concurrent-safe; LWW comparisons are lock-free. Failure Mode: Unresolvable conflict (strategy mismatch) β†’ operation flagged for manual review in audit log; data not silently overwritten.


Sourcecode Verification (Module: replication/architecture)

  • Verified files (with lock hierarchy annotations):

    • src/replication/replication_slot.cpp (Level 1β†’2, lock-free I/O)
    • src/replication/raft_v2.cpp (Level 1β†’2, fixed WAL lock violation)
    • src/replication/event_stream.cpp (Level 1β†’2, callbacks outside locks)
    • src/replication/async_wal_shipper.cpp (Level 3, timeout-guarded worker)
    • src/replication/logical_replication.cpp (Level 1β†’2, shared_mutex)
    • src/replication/multi_tier_replication.cpp (Level 1β†’2, scope-optimized)
  • Verified architecture claims:

    • orchestration + propagation/conflict + observability/policy plane split
    • explicit failure boundaries for init/promotion/slot/conflict behaviors
    • module-local ownership of replication-domain behavior surfaces
    • NEW (Wave A Block 2): strict 3-level lock hierarchy enforced
    • NEW (Wave A Block 2): all blocking I/O executes lock-free
    • NEW (Wave A Block 2): timeout guards on all waits
    • NEW (Wave A Block 2): zero circular lock ordering scenarios

ThemisDB 1.9.0-beta Β· Home Β· Module-Index Β· GitHub Β· Issues

ThemisDB Wiki

🏠 Overview

πŸ“š Compendium

πŸš€ Getting Started

πŸ“– Tutorials

πŸ“— User Guide

βš™οΈ Operations & Security

πŸ“Ÿ Ops Runbooks

πŸ—οΈ Architecture

πŸ“ ADRs

πŸ”§ Contributing

πŸ“‹ Governance

πŸ” Audit

🧩 Plugins

πŸ”Œ Adapters

πŸ’‘ Examples

πŸ“¦ Client SDKs

πŸŽ“ Training

πŸ› οΈ Tools

πŸ€– Developer LLM Wiki

Clone this wiki locally