How the COPY Command Works in pgrust: Protocol Implementation and Data Streaming
The COPY command in pgrust is implemented through a multi-layered protocol stack where the frontend client handles message framing using COPY_DATA and COPY_DONE opcodes, while the backend processes these through destination receivers and command tags, mirroring PostgreSQL's native bulk data transfer mechanism.
The pgrust system replicates PostgreSQL's wire protocol and executor behavior, providing a Rust-based implementation of the COPY command that handles bulk data transfer between clients and the server. This implementation spans the frontend libpq interface, backend traffic cop (tcop) processing, and language-specific restrictions, maintaining compatibility with standard PostgreSQL COPY semantics while leveraging Rust's type safety and async capabilities.
Frontend Protocol Handling
The frontend implementation in pgrust manages the client-side lifecycle of COPY operations, from parsing the server's initiation message to streaming data frames across the wire.
Message Structure and Constants
Protocol constants defining COPY message types are centralized in crates/interfaces/libpq/fe/src/codec.rs. These constants map directly to PostgreSQL's frontend/backend protocol:
- B_COPY_IN_RESPONSE, B_COPY_OUT_RESPONSE, and B_COPY_BOTH_RESPONSE indicate the server is ready to receive or send data.
- F_COPY_DATA and F_COPY_DONE represent frames sent from client to server.
- B_COPY_DATA and B_COPY_DONE represent frames received from the server.
These constants enable the type-safe construction and parsing of protocol frames throughout the connection lifecycle.
Client-Side COPY Operations
In crates/interfaces/libpq/fe/src/client.rs, the Client struct provides the primary interface for COPY operations. The parse_copy_start function interprets the server's response to a COPY command, returning a PGRES_COPY_* status that signals the transition into bulk data mode.
For COPY TO operations (server-to-client), the client exposes a streaming interface that yields COPY_DATA frames until the server transmits a COPY_DONE message. The COPY FROM (client-to-server) path uses build_message with F_COPY_DATA opcodes to wrap payloads before transmission, terminating with a F_COPY_DONE frame to signal completion.
Backend Processing Pipeline
Once the backend receives COPY protocol messages, the system routes them through the traffic cop and destination receivers to execute the actual data transfer.
Opcode Dispatch and Routing
The main event loop in crates/backend/tcop/postgres/src/main_loop.rs recognizes three primary COPY-related opcodes: COPY_DATA, COPY_DONE, and COPY_FAIL. These opcodes trigger specific handlers that:
- Route incoming data to the appropriate table or query destination.
- Manage flow control during high-volume transfers.
- Discard unwanted data when the caller does not require the output (e.g., during non-COPY SELECT statements).
This dispatch mechanism ensures that COPY operations maintain protocol compliance while integrating with the broader query execution framework.
Destination Receivers
The DestReceiver abstraction in crates/backend/tcop/dest_seams/src/lib.rs implements the server-side data sink for COPY operations. When executing COPY TO, the destination receiver writes rows as COPY_DATA messages to the client stream, buffering output until the final COPY_DONE frame is dispatched. This abstraction decouples the protocol layer from the physical storage implementation, allowing the same COPY logic to support files, standard input/output, and network streams.
Language Restrictions and Command Metadata
The pgrust implementation enforces PostgreSQL's architectural constraints regarding where COPY commands can execute.
PL/pgSQL Restrictions
In crates/pl/plpgsql/src/plpgsql_exec_seams/src/lib.rs, the PL/pgSQL executor explicitly blocks COPY commands with the error message: "cannot COPY to/from client in PL/pgSQL". This restriction exists because COPY requires a direct client connection for data streaming, which is unavailable within the context of a stored procedure running on the server. The implementation mirrors PostgreSQL's behavior by returning ERRCODE_FEATURE_NOT_SUPPORTED when detecting COPY statements in PL/pgSQL blocks.
Command Tagging and Progress Tracking
The system defines CMDTAG_COPY = 56 in crates/backend/tcop/utility/src/consts.rs to identify COPY operations in the command tag system. This tag enables:
- Accurate row counting via
pg_command_countupdates. - Integration with
pg_stat_progress_copythrough the activity subsystem incrates/backend/utils/activity/backend_progress_seams/src/lib.rs. - Executor statistics tracking for rows processed and bytes transferred during the operation.
Practical Implementation Examples
The following examples demonstrate how to interact with the pgrust COPY implementation using the async client interface.
Streaming Data with COPY TO
This example executes a COPY TO STDOUT command and processes the resulting data stream:
use pgrust::interfaces::libpq::fe::client::Client;
/// Copy a whole table to STDOUT.
async fn copy_table(conn_str: &str) -> Result<(), Box<dyn std::error::Error>> {
let mut client = Client::connect(conn_str).await?;
// Issue a COPY command – the server will respond with a COPY‑OUT start message.
let mut copy = client.copy_out("COPY my_table TO STDOUT (FORMAT csv)").await?;
// Stream the data chunks until the server signals COPY‑DONE.
while let Some(chunk) = copy.next().await {
let bytes = chunk?;
// Here we simply write to stdout, but any sink (file, network) works.
std::io::stdout().write_all(&bytes)?;
}
// The COPY operation is automatically finalized when the stream ends.
Ok(())
}
Loading Data with COPY FROM
This example demonstrates bulk loading using COPY FROM STDIN:
use pgrust::interfaces::libpq::fe::client::Client;
use futures::stream::StreamExt;
async fn copy_from(conn_str: &str) -> Result<(), Box<dyn std::error::Error>> {
let mut client = Client::connect(conn_str).await?;
let mut copy = client.copy_in("COPY my_table FROM STDIN (FORMAT csv)").await?;
// Simulate a CSV source – replace with actual data source as needed.
let csv_rows = vec![
"1,alice,2024-01-01\n",
"2,bob,2024-01-02\n",
];
for row in csv_rows {
// Each call sends a COPY‑DATA frame.
copy.send(row.as_bytes()).await?;
}
// Signal end‑of‑copy.
copy.finish().await?;
Ok(())
}
Summary
- Protocol Constants:
crates/interfaces/libpq/fe/src/codec.rsdefinesF_COPY_DATA,F_COPY_DONE, and response opcodes that frame all COPY communication. - Client Implementation:
crates/interfaces/libpq/fe/src/client.rsprovidescopy_outandcopy_inmethods that handle the handshake and streaming lifecycle. - Backend Dispatch:
crates/backend/tcop/postgres/src/main_loop.rsroutesCOPY_DATAandCOPY_DONEopcodes to appropriate handlers. - Destination Abstraction:
crates/backend/tcop/dest_seams/src/lib.rsimplements theDestReceiverpattern for server-side data emission. - Language Constraints:
crates/pl/plpgsql/src/plpgsql_exec_seams/src/lib.rsexplicitly forbids COPY in PL/pgSQL contexts. - Metadata:
CMDTAG_COPY(value 56) and progress tracking in the activity subsystem provide visibility into operation status.
Frequently Asked Questions
How does pgrust handle the COPY protocol handshake?
The handshake begins when the client sends a COPY command and the server responds with one of the B_COPY_IN_RESPONSE, B_COPY_OUT_RESPONSE, or B_COPY_BOTH_RESPONSE messages. In crates/interfaces/libpq/fe/src/client.rs, the parse_copy_start function interprets these messages to transition the connection into COPY mode, allowing subsequent COPY_DATA frames to flow in the appropriate direction.
Can I use COPY commands inside PL/pgSQL functions in pgrust?
No. According to the source code in crates/pl/plpgsql/src/plpgsql_exec_seams/src/lib.rs, PL/pgSQL explicitly rejects COPY commands with the error "cannot COPY to/from client in PL/pgSQL". This restriction exists because COPY requires a direct client connection for data streaming, which is unavailable within the server-side execution context of a stored procedure.
What is the difference between copy_out and copy_in in the pgrust client?
copy_out initiates a COPY TO operation where the server sends data to the client, returning an async stream that yields COPY_DATA frames until the server sends COPY_DONE. Conversely, copy_in initiates a COPY FROM operation, providing a handle that accepts data via the send method (which wraps payloads in F_COPY_DATA frames) and terminates with finish (which sends F_COPY_DONE).
Where does pgrust track the progress of COPY operations?
Progress tracking occurs in crates/backend/utils/activity/backend_progress_seams/src/lib.rs, which exposes real-time metrics compatible with pg_stat_progress_copy. The system uses CMDTAG_COPY (defined as 56 in crates/backend/tcop/utility/src/consts.rs) to tag operations and update row counts and byte transfer statistics during the COPY lifecycle.
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 →