How pgrust Handles Physical Replication: A Deep Dive into the Rust Implementation

pgrust implements PostgreSQL physical replication through a client-side driver that initiates COPY-BOTH streams and a server-side wal-sender process that reads WAL segments and emits them as copy-protocol messages, faithfully porting PostgreSQL's C logic to safe Rust.

pgrust is a Rust implementation of PostgreSQL internals that replicates the database's physical streaming replication semantics while leveraging Rust's memory safety guarantees. Understanding pgrust physical replication requires examining how the codebase mirrors PostgreSQL's architecture across the client driver, wal-sender process, and replication helper functions located in specific crate paths.

The Three-Layer Stack Behind pgrust Physical Replication

The implementation spans three distinct layers that mirror PostgreSQL's architecture:

  • Client-side driver (crates/interfaces/libpq/fe/src/client.rs): Opens replication connections via START_REPLICATION and exposes the copy-protocol API including start_replication, copy_receive, and copy_send.
  • Server-side wal-sender (crates/backend/replication/walsender/src/physical.rs): Runs as a background process that determines safe streaming boundaries, reads WAL data, and emits copy-protocol messages.
  • Replication helpers (crates/backend/replication/walsender/src/start_replication/*): Utility functions such as GetStandbyFlushRecPtr and xlogreader_close_if_open that manage timeline switches and flush pointers.

Initiating Replication: The Client-Side Implementation

The client driver exposes PgClientConn::start_replication to establish a physical replication session. When a client calls:

let mut conn = PgClientConn::connect(...)?;
let result = conn.start_replication("START_REPLICATION 0/0")?;

The driver automatically adds the startup option replication=true (defined in protocol3.rs) and transmits the START_REPLICATION command. The server treats this as a COPY-BOTH stream, switching the connection into replication mode where both client and server can send data asynchronously.

The client receives WAL data through copy_receive, which returns an enum distinguishing between data chunks (CopyRecv::Data), stream completion (CopyRecv::Done), and blocking states (CopyRecv::WouldBlock).

The Wal-Sender Loop: XLogSendPhysical in Detail

The core server-side logic resides in crates/backend/replication/walsender/src/physical.rs, specifically in the XLogSendPhysical function. This Rust implementation ports PostgreSQL's C logic directly, performing six critical operations per iteration:

  1. Graceful shutdown detection – Checks proc.got_STOPPING to initiate a stopping state when the server is shutting down.
  2. Send request pointer calculation – Determines streaming boundaries using GetStandbyFlushRecPtr to fetch the latest flushed LSN, accounting for cascading standbys and historic timelines.
  3. Historic timeline completion – When streaming a historic timeline reaches its fork point, the sender closes the xlog reader via xlogreader_close_if_open and emits a CopyDone message (type 'c').
  4. Data window calculation – Respects MAX_SEND_SIZE for message boundaries, aligns reads to page boundaries, and updates WalSndCaughtUp status.
  5. WAL reading and message emission – Calls xlog::wal_read::call to fetch raw WAL slices, then packs them into messages beginning with 'w' followed by dataStart, walEnd, sendtime, and the WAL bytes. These are sent via pq_putmessage_noblock_output_message.
  6. Shared memory updates – Advances sentPtr, updates the process title, and records replication lag via LagTrackerWrite.

Timeline Switches and Replication Slot Management

Physical replication must handle complex scenarios including cascading standbys and timeline switches:

  • Cascading standbys – The sender detects promotion states through !xlog::recovery_in_progress and switches timelines accordingly.
  • Historic timelines – When the wal-sender streams a historic timeline, it stops at the fork point, closes the xlog reader, and sends CopyDone (type 'c') to signal completion.
  • Replication slots – While slot initialization occurs elsewhere (replication_slot_initialize, sync_replication_slots), the wal-sender reads slot LSN values through GetStandbyFlushRecPtr to determine synchronization points.

Permission checks for replication are handled in crates/backend/utils/init/miscinit_seams/src/lib.rs, which verifies the rolreplication role attribute. Configuration variables like max_replication_slots and wal_sender_timeout are defined in crates/backend/utils/misc/guc_tables/src/vars.rs.

End-to-End Data Flow in pgrust Physical Replication

The complete data flow follows this path:


PgClientConn::start_replication()  [client.rs]
            │
            ▼
WalSender (XLogSendPhysical)       [physical.rs]
   ├── Reads WAL via xlog::wal_read
   ├── Packs into 'w' messages
   └── Emits via copy-protocol
            │
            ▼
PgClientConn::copy_receive()         [client.rs]

This architecture ensures that pgrust physical replication maintains wire compatibility with PostgreSQL while providing Rust's safety guarantees against memory errors in the wal-sender process.

Practical Implementation Examples

Starting a Physical Replication Stream (Client)

use pgrust::interfaces::libpq::fe::client::PgClientConn;
use pgrust::interfaces::libpq::fe::transport::TcpTransport;

fn start_physical_replication() -> Result<(), Box<dyn std::error::Error>> {
    // Connect to the primary with replication=true
    let transport = TcpTransport::connect("127.0.0.1:5432")?;
    let mut conn = PgClientConn::connect(
        transport,
        &protocol3::StartupParams::replication(),
        None,
    )?;

    // Issue START_REPLICATION (physical)
    let _res = conn.start_replication("START_REPLICATION 0/0")?;

    // Loop receiving WAL data
    loop {
        match conn.copy_receive()? {
            CopyRecv::Data(buf) => {
                // `buf` contains raw WAL bytes
                println!("Received {} bytes of WAL", buf.len());
            }
            CopyRecv::Done => break, // End of stream
            CopyRecv::WouldBlock => continue, // Wait for more data
        }
    }

    Ok(())
}

Server-Side Physical Streaming Loop

use pgrust::backend::replication::walsender::physical::XLogSendPhysical;

// In the wal-sender background worker's main loop:
loop {
    // Periodically called by the scheduler
    XLogSendPhysical();

    // Wait for new WAL or client acknowledgment
    std::thread::sleep(std::time::Duration::from_millis(100));
}

Summary

  • pgrust physical replication implements PostgreSQL's streaming semantics through a three-layer architecture: client driver, wal-sender process, and replication helpers.
  • The client initiates streams via PgClientConn::start_replication in crates/interfaces/libpq/fe/src/client.rs, using the COPY-BOTH protocol.
  • The server-side wal-sender runs XLogSendPhysical from crates/backend/replication/walsender/src/physical.rs, handling WAL reads, timeline switches, and graceful shutdowns.
  • Helper functions like GetStandbyFlushRecPtr manage replication slots and cascading standby logic.
  • The implementation preserves PostgreSQL's wire protocol compatibility while leveraging Rust's memory safety for the wal-sender background process.

Frequently Asked Questions

What is the COPY-BOTH protocol and how does pgrust use it for physical replication?

The COPY-BOTH protocol is a PostgreSQL wire protocol mode that allows bidirectional data streaming between client and server. In pgrust physical replication, the client initiates this mode by sending START_REPLICATION with the replication=true startup option. The server then treats the connection as a COPY-BOTH stream, enabling the wal-sender to push WAL data (prefixed with 'w') to the client while receiving standby feedback. The client uses copy_receive to read these messages and CopyRecv::Done to detect stream completion.

How does pgrust handle timeline switches during physical replication?

Timeline switches are handled in crates/backend/replication/walsender/src/physical.rs within the XLogSendPhysical function. The wal-sender detects promotion events by checking !xlog::recovery_in_progress to identify cascading standbys. For historic timelines, the sender tracks the fork point and stops streaming at that boundary, closing the xlog reader via xlogreader_close_if_open and sending a CopyDone message (type 'c') to signal the timeline switch to the client.

What limits the amount of WAL data sent per message in pgrust?

The wal-sender respects MAX_SEND_SIZE when calculating the data window for each message. As implemented in crates/backend/replication/walsender/src/physical.rs, the XLogSendPhysical function aligns reads to page boundaries and ensures no single message exceeds this size limit. The function updates WalSndCaughtUp to track whether the sender has reached the end of available WAL.

How does pgrust ensure memory safety in the wal-sender process?

Unlike PostgreSQL's C implementation, pgrust physical replication leverages Rust's ownership and borrowing rules to prevent memory errors in the wal-sender. The XLogSendPhysical function and its helpers (such as GetStandbyFlushRecPtr) are implemented in safe Rust within crates/backend/replication/walsender/src/, eliminating risks of buffer overflows or use-after-free bugs while maintaining the exact control flow and semantics required for PostgreSQL compatibility.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →