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:

  1. Route incoming data to the appropriate table or query destination.
  2. Manage flow control during high-volume transfers.
  3. 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:

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

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:

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 →