How pgrust Handles Logical Replication: A Deep Dive into the PostgreSQL Rust Client
pgrust implements logical replication by exposing a CopyBoth streaming API that mirrors PostgreSQL's libpq logical-replication protocol, using PgClientConn to manage bidirectional streams between Rust clients and PostgreSQL servers.
The pgrust project (malisper/pgrust) provides a Rust-native implementation of the PostgreSQL wire protocol and client libraries. When handling pgrust logical replication, the codebase exposes a type-safe abstraction over PostgreSQL's libpq logical-replication protocol, enabling Rust applications to act as replication consumers through the CopyBoth streaming mechanism.
The CopyBoth Replication Protocol
PostgreSQL logical replication relies on the COPY BOTH protocol mode, which creates a bidirectional stream between client and server. In pgrust, this is implemented through the PgClientConn struct in crates/interfaces/libpq/fe/src/client.rs, which manages the stateful connection required for replication slots and logical decoding.
Establishing a Replication Connection
Startup Parameters and Connection Flags
Replication connections require special handling during the startup phase. In crates/interfaces/libpq/fe/src/protocol3.rs, the StartupParams struct provides a replication() method that sets the replication parameter to "true", signaling the server to enter replication mode.
To negotiate a replication connection, configure the startup parameters before calling PgClientConn::connect:
let params = pgrust::interfaces::libpq::fe::protocol3::StartupParams::new()
.user("replicator")
.replication("true"); // Enables logical replication mode
Starting and Managing the Replication Stream
Initiating START_REPLICATION
Once connected, the client initiates logical replication by calling PgClientConn::start_replication (defined at line 425 in client.rs). This sends the START_REPLICATION command as part of a CopyBoth transaction, opening the bidirectional stream.
The method returns a result with ExecStatusType::CopyBoth, confirming the connection has entered streaming mode.
Receiving XLogData Messages
The PgClientConn::copy_receive method (lines 39-50 in client.rs) wraps the underlying PQgetCopyData call and returns a CopyRecv enum with variants:
Data– Contains the logical decoding payload (XLogDatamessages)Done– Signals the stream has completedWouldBlock– Indicates no data is currently available
Sending Standby Feedback
Replication slots require periodic feedback to maintain progress and prevent WAL accumulation. The PgClientConn::copy_send method allows the client to push standby progress messages back to the server, while copy_done and end_copy properly terminate the CopyBoth lifecycle.
Logical Decoding Plugin Architecture
The Test Decoding Plugin Implementation
The repository includes a reference implementation in crates/contrib/test_decoding/src/lib.rs, demonstrating how logical decoding plugins receive transactions through LogicalDecodingContext, TxnHandle, RelationHandle, and ChangeHandle. This plugin uses the types_logical crate types that underpin the entire replication data path, processing changes into textual output.
Server-Side Infrastructure
Replication Slot Management
Server-side functions for slot management, such as pg_create_logical_replication_slot, are declared in crates/backend/utils/fmgr/builtin_canonical.rs (function IDs 3786, 3787). These provide the server-side counterpart to the client's replication commands, enabling Rust applications to interact with PostgreSQL's replication slot infrastructure.
Configuration GUCs
Runtime configuration for logical replication workers is controlled through GUCs defined in crates/backend/utils/misc/guc_tables/src/vars.rs, including max_logical_replication_workers and wal_receiver_create_temp_slot. These variables expose runtime knobs for tuning replication behavior.
Complete Implementation Example
The following example demonstrates the complete flow for establishing a pgrust logical replication connection, receiving XLogData messages, and sending feedback:
use pgrust::interfaces::libpq::fe::client::PgClientConn;
use pgrust::interfaces::libpq::fe::transport::TcpTransport;
fn main() -> Result<(), Box<dyn std::error::Error>> {
// 1️⃣ Open a TCP transport to the server
let transport = TcpTransport::connect("localhost:5432")?;
// 2️⃣ Build startup parameters with replication=true
let params = pgrust::interfaces::libpq::fe::protocol3::StartupParams::new()
.user("replicator")
.replication("true"); // <‑‑ logical‑replication flag
// 3️⃣ Connect and authenticate
let mut conn = PgClientConn::connect(transport, ¶ms, None)?;
// 4️⃣ Start the replication stream
let result = conn.start_replication("START_REPLICATION 0/0")?;
assert!(result.status == pgrust::interfaces::libpq::fe::result::ExecStatusType::CopyBoth);
// 5️⃣ Receive logical‑decoding messages
loop {
match conn.copy_receive()? {
pgrust::interfaces::libpq::fe::client::CopyRecv::Data(buf) => {
// `buf` contains an XLogData record from the output plugin
println!("Received {} bytes", buf.len());
// …process the logical change…
// (optionally send feedback)
conn.copy_send(b"standby feedback ...")?;
}
pgrust::interfaces::libpq::fe::client::CopyRecv::Done => break,
pgrust::interfaces::libpq::fe::client::CopyRecv::WouldBlock => {
// In a real app you would poll / await the socket here
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
}
// 6️⃣ End the COPY‑Both transaction
conn.copy_done()?;
conn.end_copy()?;
Ok(())
}
Summary
- CopyBoth Protocol: pgrust uses
PgClientConnto manage bidirectionalCopyBothstreams that mirror libpq's logical replication behavior. - Startup Configuration: Set
replication("true")inStartupParamsto negotiate a replication connection during protocol negotiation inprotocol3.rs. - Stream Management: Use
start_replication,copy_receive, andcopy_sendto handle the lifecycle of logical decoding messages. - Type Safety: The
CopyRecvenum provides safe handling ofXLogDatamessages withData,Done, andWouldBlockvariants. - Plugin Support: The
test_decodingplugin incrates/contrib/test_decoding/src/lib.rsdemonstrates the server-side logical decoding API usingtypes_logicalprimitives. - Server Integration: Replication slots and GUCs in
builtin_canonical.rsandguc_tables/src/vars.rscomplete the implementation stack.
Frequently Asked Questions
What is the CopyBoth protocol in pgrust?
The CopyBoth protocol is PostgreSQL's bidirectional streaming mode used for logical replication. In pgrust, it is implemented through the PgClientConn struct in client.rs, which sends START_REPLICATION commands and manages the bidirectional flow of XLogData messages and standby feedback between Rust clients and PostgreSQL servers.
How do I initiate a logical replication connection in pgrust?
You initiate a logical replication connection by calling PgClientConn::connect with StartupParams that include .replication("true"), as defined in protocol3.rs. This sets the replication flag in the startup packet, causing the server to negotiate a replication connection rather than a standard query session.
What does the CopyRecv enum represent in pgrust logical replication?
The CopyRecv enum in client.rs represents the possible states of data reception during logical replication. It has three variants: Data containing the payload buffer, Done indicating the stream has ended, and WouldBlock signaling that no data is currently available and the client should poll or wait.
Does pgrust provide server-side logical decoding capabilities?
Yes, pgrust includes server-side infrastructure for logical decoding, demonstrated by the test_decoding output plugin in crates/contrib/test_decoding/src/lib.rs. This plugin implements the LogicalDecodingContext trait and handles transactions through TxnHandle, RelationHandle, and ChangeHandle types, processing changes into textual output using the types_logical crate.
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 →