How Logical and Physical Replication Work in pgrust: A Deep Dive into the CopyBoth API
pgrust implements logical replication through a CopyBoth streaming API that mirrors PostgreSQL's libpq protocol, exposing type-safe Rust methods for connection setup, slot management, and bidirectional data flow between clients and the replication server.
The malisper/pgrust project provides a Rust-native implementation of PostgreSQL's client-server protocol, including comprehensive support for logical and physical replication. By exposing a CopyBoth streaming interface through PgClientConn, pgrust enables Rust applications to receive logical decoding messages and send replication feedback using the same control flow as PostgreSQL's native libpq client, while maintaining memory safety and type guarantees.
Establishing the Replication Connection
Replication begins with a specialized connection handshake that negotiates the replication protocol. In crates/interfaces/libpq/fe/src/client.rs, the PgClientConn::connect method handles this initialization, building a startup packet that includes the replication flag.
Startup Parameters and the Replication Flag
The connection process starts in crates/interfaces/libpq/fe/src/protocol3.rs, where StartupParams::new() constructs the parameter list. To enable replication mode, the code sets the replication startup option to true using the replication method:
let params = pgrust::interfaces::libpq::fe::protocol3::StartupParams::new()
.user("replicator")
.replication("true");
This flag signals the PostgreSQL server to enter replication mode, allowing subsequent START_REPLICATION commands. Connection options are managed through the registry system defined in crates/interfaces/libpq/fe/src/registry.rs, which validates and processes these parameters before transmission.
Initiating the Logical Replication Stream
Once connected, the client initiates streaming through the PgClientConn::start_replication method found in client.rs (around line 425). This method sends the START_REPLICATION command to the server, which responds by entering CopyBoth mode—a bidirectional streaming state where the server continuously sends WAL (Write-Ahead Log) data and the client can respond with feedback messages.
The method returns an execution result with status == ExecStatusType::CopyBoth, indicating that the connection is now in streaming mode and ready for the copy data loop.
Bidirectional Data Flow with CopyBoth
The CopyBoth protocol creates a persistent stream where the client alternates between receiving decoded WAL records and sending standby progress updates. This loop is managed through three core methods on PgClientConn.
Receiving Logical Decoding Messages
The copy_receive method (defined around lines 39-50 in client.rs) wraps the underlying PQgetCopyData call and returns a CopyRecv enum with three variants:
- CopyRecv::Data(buf) – Contains raw bytes of an
XLogDatamessage from the logical decoding output plugin - CopyRecv::Done – Signals the server has finished sending data
- CopyRecv::WouldBlock – Indicates no data is available and the client should wait
match conn.copy_receive()? {
pgrust::interfaces::libpq::fe::client::CopyRecv::Data(buf) => {
// Process XLogData record from the output plugin
println!("Received {} bytes", buf.len());
}
pgrust::interfaces::libpq::fe::client::CopyRecv::Done => break,
pgrust::interfaces::libpq::fe::client::CopyRecv::WouldBlock => {
// Poll or await the socket
}
}
Sending Replication Feedback
To maintain the replication slot and prevent WAL accumulation on the server, the client must periodically send standby progress updates. The copy_send method transmits feedback messages, while copy_done signals the end of the COPY operation:
conn.copy_send(b"standby feedback ...")?;
// ... after processing completes ...
conn.copy_done()?;
conn.end_copy()?;
This bidirectional flow mirrors exactly the control flow found in PostgreSQL's C-based libpq implementation, but exposes safe Rust bindings that handle buffer management and error propagation.
Server-Side Logical Decoding Infrastructure
While the client handles the streaming protocol, the server side manages logical decoding through output plugins and configuration variables.
The Test Decoding Output Plugin
The repository includes a reference implementation in crates/contrib/test_decoding/src/lib.rs that demonstrates how logical decoding plugins process replication events. This plugin receives:
- LogicalDecodingContext – The runtime context for decoding operations
- TxnHandle – Represents transactions being decoded
- RelationHandle – References to table schemas
- ChangeHandle – Individual row changes (INSERT, UPDATE, DELETE)
These types from the types_logical crate form the foundation of the logical decoding pipeline, converting raw WAL entries into application-readable formats.
Replication Slot Management
Logical replication requires persistent slots to track client progress. Server-side functions such as pg_create_logical_replication_slot are declared in crates/backend/utils/fmgr/builtin_canonical.rs (function IDs 3786 and 3787). These functions manage the metadata that associates a named slot with a specific LSN (Log Sequence Number) position in the WAL stream.
Configuration GUCs
Replication behavior is controlled through GUC (Grand Unified Configuration) variables defined in crates/backend/utils/misc/guc_tables/src/vars.rs. Key settings include:
- max_logical_replication_workers – Limits concurrent logical replication processes
- wal_receiver_create_temp_slot – Controls temporary slot creation for WalReceiver processes
These configuration knobs expose runtime control over replication resource usage and slot persistence.
Complete Implementation Example
The following example demonstrates the full replication client workflow, from connection setup through the streaming loop:
use pgrust::interfaces::libpq::fe::client::PgClientConn;
use pgrust::interfaces::libpq::fe::transport::TcpTransport;
fn main() -> Result<(), Box<dyn std::error::Error>> {
// Open a TCP transport to the server
let transport = TcpTransport::connect("localhost:5432")?;
// Build startup parameters with replication=true
let params = pgrust::interfaces::libpq::fe::protocol3::StartupParams::new()
.user("replicator")
.replication("true");
// Connect and authenticate
let mut conn = PgClientConn::connect(transport, ¶ms, None)?;
// Start the replication stream
let result = conn.start_replication("START_REPLICATION 0/0")?;
assert!(result.status == pgrust::interfaces::libpq::fe::result::ExecStatusType::CopyBoth);
// Receive logical-decoding messages
loop {
match conn.copy_receive()? {
pgrust::interfaces::libpq::fe::client::CopyRecv::Data(buf) => {
println!("Received {} bytes", buf.len());
// Process the logical change and send feedback
conn.copy_send(b"standby feedback ...")?;
}
pgrust::interfaces::libpq::fe::client::CopyRecv::Done => break,
pgrust::interfaces::libpq::fe::client::CopyRecv::WouldBlock => {
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
}
// End the COPY-Both transaction
conn.copy_done()?;
conn.end_copy()?;
Ok(())
}
This implementation follows the exact sequence defined in the source: PgClientConn::connect handles the startup negotiation, start_replication initiates the CopyBoth stream, and copy_receive/copy_send manage the bidirectional data flow.
Summary
- pgrust implements logical replication through a type-safe Rust API that mirrors PostgreSQL's libpq CopyBoth protocol.
- The connection process in
client.rsandprotocol3.rsnegotiates replication mode via startup parameters. - CopyBoth mode enables bidirectional streaming where
copy_receivedeliversXLogDatamessages andcopy_sendtransmits standby feedback. - Server-side infrastructure includes logical decoding types (
TxnHandle,RelationHandle), slot management functions inbuiltin_canonical.rs, and configuration GUCs invars.rs. - The test decoding plugin in
crates/contrib/test_decoding/src/lib.rsdemonstrates how to process replication events using thetypes_logicalcrate.
Frequently Asked Questions
What is the CopyBoth mode in pgrust replication?
CopyBoth is a streaming protocol mode that opens a bidirectional data channel between the client and server. In crates/interfaces/libpq/fe/src/client.rs, entering CopyBoth mode (confirmed by ExecStatusType::CopyBoth) allows the client to continuously receive logical decoding data via copy_receive while simultaneously sending replication feedback through copy_send, maintaining the streaming state until copy_done is called.
How does pgrust handle logical replication slots?
Replication slots are managed through server-side functions declared in crates/backend/utils/fmgr/builtin_canonical.rs, specifically pg_create_logical_replication_slot (function IDs 3786 and 3787). These slots persistently track the client's LSN position in the WAL stream, ensuring that the server retains WAL segments until the client confirms receipt via the CopyBoth feedback mechanism.
What types represent logical decoding transactions in pgrust?
The logical decoding system uses types from the types_logical crate, including LogicalDecodingContext for the decoding runtime, TxnHandle for transaction boundaries, RelationHandle for schema information, and ChangeHandle for individual row modifications. These types are demonstrated in the test decoding plugin at crates/contrib/test_decoding/src/lib.rs.
How does pgrust differ from PostgreSQL's libpq for replication?
While pgrust mirrors the exact C control flow of PostgreSQL's libpq logical replication client, it exposes a memory-safe, type-safe Rust API. The PgClientConn struct in client.rs wraps the underlying protocol implementation with Rust's ownership and error handling guarantees, while still supporting the same START_REPLICATION commands and CopyBoth streaming behavior as the native C library.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →