# How Does the pgrust Executor Manage Parallel Query Execution?

> Discover how the pgrust executor manages parallel query execution using ParallelContextHandle for leader-worker coordination and Rust error handling.

- Repository: [Michael Malis/pgrust](https://github.com/malisper/pgrust)
- Tags: internals
- Published: 2026-07-13

---

**The pgrust executor orchestrates parallel query execution by establishing a **ParallelContextHandle** that manages leader-worker coordination through dynamic shared memory, enabling workers to share tuplestores and filesets while using Rust's Result-based error handling instead of PostgreSQL's longjmp mechanism.**

Pgrust is a Rust implementation of PostgreSQL's query engine that reimagines the executor's parallel query capabilities through the dedicated **execparallel** crate. Unlike the traditional C-based backend, pgrust utilizes type-safe handles and shared memory abstractions to manage parallel workers, allowing multiple processes to execute plan nodes concurrently while maintaining memory safety and deterministic error propagation.

## The execparallel Crate: Core Parallel Primitives

The foundation of pgrust's parallel execution resides in the `execparallel` support crate, which provides the type definitions for dynamic shared memory (DSM) management. According to the source code in [[`crates/_support/types/execparallel/src/lib.rs`](https://github.com/malisper/pgrust/blob/main/crates/_support/types/execparallel/src/lib.rs)](https://github.com/malisper/pgrust/blob/main/crates/_support/types/execparallel/src/lib.rs), this module defines **ParallelContextHandle**, **ParallelWorkerContextHandle**, and shared resource handles including **SharedTuplestoreHandle** and **SharedFileSetHandle**. These opaque structures map directly to PostgreSQL's DSM layout but provide Rust's memory safety guarantees through cloneable, thread-safe handles.

### Shared Memory Abstractions

Parallel execution requires structures that persist across process boundaries. The `execparallel` crate implements **SharedTuplestoreHandle** for exchanging intermediate rows between workers and **SharedFileSetHandle** for spilling data that exceeds memory capacity. These handles allow the leader process to distribute references to worker contexts without unsafe pointer manipulation, ensuring that all participants in a parallel query share the same logical memory space.

## Leader-Worker Coordination Model

Pgrust implements a leader-worker pattern where the executor process spawns parallel workers to execute portions of the query plan. The coordination begins when the executor detects a parallel-aware plan node and initializes the parallel context.

### Query Planning and Worker Allocation

During executor startup, the system records the intended number of parallel workers in the query descriptor, specifically tracking `es_parallel_workers_to_launch` and `es_parallel_workers_launched`. As implemented in [[`crates/contrib/pg_stat_statements/src/store.rs`](https://github.com/malisper/pgrust/blob/main/crates/contrib/pg_stat_statements/src/store.rs)](https://github.com/malisper/pgrust/blob/main/crates/contrib/pg_stat_statements/src/store.rs), these counters hook into the executor start mechanism to track resource allocation and ensure the planner's parallel degree requests are honored.

### Worker Spawning and Context Distribution

The leader process creates a **ParallelContextHandle** and spawns workers using **ParallelWorkerContextHandle** instances. Each worker receives a copy of the shared context and maps the DSM area to access the same tuplestores and filesets. The parallel hash join implementation in [[`crates/backend/executor/nodeHash/src/parallel.rs`](https://github.com/malisper/pgrust/blob/main/crates/backend/executor/nodeHash/src/parallel.rs)](https://github.com/malisper/pgrust/blob/main/crates/backend/executor/nodeHash/src/parallel.rs) demonstrates how workers obtain their execution context and begin processing assigned batches of the build relation.

## Inter-Worker Communication via Shared Tuplestores

When parallel workers need to exchange data—such as during the build phase of a hash join or parallel aggregation—they use the shared tuplestore API rather than copying data between processes.

### Parallel Scan API

The shared tuplestore exposes a coordinated scanning interface defined in [[`crates/backend/utils/sort/sort_storage_seams/src/lib.rs`](https://github.com/malisper/pgrust/blob/main/crates/backend/utils/sort/sort_storage_seams/src/lib.rs)](https://github.com/malisper/pgrust/blob/main/crates/backend/utils/sort/sort_storage_seams/src/lib.rs). Workers invoke `sts_begin_parallel_scan` to initialize their access, `sts_parallel_scan_next` to retrieve tuples in a round-robin fashion, and `sts_end_parallel_scan` to release resources. This ensures that multiple workers can read from the same intermediate result set without race conditions or duplicate processing.

### Shared Filesets for Spill Data

For operations requiring disk spillage, **SharedFileSetHandle** provides a unified namespace for temporary files. Workers write to logically separate portions of the same fileset, allowing the leader to coordinate merge operations or sequential scans across spilled data without creating separate temporary files per worker.

## Error Handling in Parallel Workers

Traditional PostgreSQL uses `longjmp` for error recovery, which is unsafe across process boundaries. Pgrust replaces this with **Result-threaded** error propagation throughout the executor.

### Result-Based Error Propagation

As shown in [[`crates/pl/plpgsql/src/exec/src/lib.rs`](https://github.com/malisper/pgrust/blob/main/crates/pl/plpgsql/src/exec/src/lib.rs)](https://github.com/malisper/pgrust/blob/main/crates/pl/plpgsql/src/exec/src/lib.rs), all parallel worker functions return `Result<PgResult<...>, PgError>`. When a worker encounters an error, it returns `Err(PgError)` rather than invoking a long jump. The leader monitors these results and triggers cleanup of all shared resources—including DSM segments and filesets—if any worker fails, ensuring no resource leaks occur across the parallel context.

## Practical Implementation Examples

### Initializing the Parallel Context

```rust
// Pattern derived from nodeHash/src/parallel.rs
use execparallel::{ParallelContextHandle, SharedTuplestoreHandle};

// Leader creates the parallel context
let pcxt = ParallelContextHandle::new();
// Allocate shared tuplestore for hash table build
let shared_ts = SharedTuplestoreHandle::new(&pcxt);

// Launch workers with access to shared context
for worker_id in 0..num_workers {
    let worker_ctx = pcxt.create_worker_handle();
    spawn_parallel_worker(worker_ctx, shared_ts.clone());
}

```

### Worker Implementation with Shared Tuplestore

```rust
// Worker entry point pattern using sort_storage_seams API
use execparallel::{
    sts_begin_parallel_scan, 
    sts_parallel_scan_next, 
    sts_end_parallel_scan
};

fn parallel_worker_main(
    worker_ctx: ParallelWorkerContextHandle,
    shared_ts: SharedTuplestoreHandle
) -> Result<(), PgError> {
    // Initialize scan access
    let mut accessor = worker_ctx.get_tuplestore_accessor(shared_ts);
    sts_begin_parallel_scan(&mut accessor)?;
    
    // Process tuples in round-robin fashion
    while let Some((tuple, meta)) = sts_parallel_scan_next(&mut accessor)? {
        process_tuple(tuple, meta)?;
    }
    
    // Cleanup
    sts_end_parallel_scan(&mut accessor)?;
    Ok(())
}

```

### Error Propagation Across Workers

```rust
// Error handling pattern from plpgsql/src/exec/src/lib.rs
match parallel_worker_main(ctx, ts) {
    Ok(_) => {/* worker succeeded */},
    Err(e) => {
        // Abort entire parallel context and propagate
        execparallel::abort_parallel_context(pcxt, e)?;
        return Err(e.into());
    }
}

```

## Summary

- The **execparallel** crate provides the foundational handles (`ParallelContextHandle`, `SharedTuplestoreHandle`) for DSM-based parallel execution.
- Parallel workers communicate through **shared tuplestores** using the `sts_*` API defined in `sort_storage_seams`.
- **SharedFileSetHandle** manages temporary disk storage across worker processes for spill operations.
- Error handling uses **Result-threaded** propagation rather than PostgreSQL's `longjmp`, ensuring memory safety across process boundaries.
- The executor tracks worker allocation via counters (`es_parallel_workers_to_launch`) hooked into the planner statistics in `pg_stat_statements`.

## Frequently Asked Questions

### What is the role of ParallelContextHandle in pgrust?

The **ParallelContextHandle** is the central coordination object that manages the dynamic shared memory (DSM) segment for a parallel query. It maintains the shared state between the leader and worker processes, including references to shared tuplestores and filesets, and provides the mechanism for spawning and synchronizing parallel workers through the `execparallel` crate.

### How does pgrust handle errors in parallel workers differently from PostgreSQL?

While PostgreSQL uses `longjmp` for error recovery—which is unsafe across process boundaries—pgrust implements **Result-threaded** error handling. Parallel workers return `Result<T, PgError>` types, allowing the leader to gracefully handle failures through Rust's type system. This ensures that shared resources are properly cleaned up when any worker encounters an error, preventing memory leaks and undefined behavior.

### What are shared tuplestores used for in parallel execution?

**SharedTuplestoreHandle** structures serve as the inter-process communication channel for intermediate query results. During parallel operations like hash joins or sorts, workers write their partial results to the shared tuplestore, and the leader or other workers read from it using the `sts_begin_parallel_scan`, `sts_parallel_scan_next`, and `sts_end_parallel_scan` API. This eliminates the need to merge results through the client socket or temporary files.

### How does pgrust differ from PostgreSQL's native parallel executor?

Pgrust reimplements PostgreSQL's parallel executor in Rust, replacing raw pointer manipulations and `longjmp` error handling with type-safe handles and Result-based error propagation. The architecture uses the **execparallel** crate to abstract DSM operations, providing memory safety guarantees while maintaining compatibility with PostgreSQL's parallel query semantics. The shared resource handles (`SharedTuplestoreHandle`, `SharedFileSetHandle`) are managed through Rust's ownership model rather than manual memory management.