diff --git a/.github/workflows/ci-new-parser.yml b/.github/workflows/ci-new-parser.yml index 76b8131e9..2e2a11c38 100644 --- a/.github/workflows/ci-new-parser.yml +++ b/.github/workflows/ci-new-parser.yml @@ -25,9 +25,7 @@ jobs: id: pgdog-bin uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} - uses: Swatinem/rust-cache@v2 if: steps.pgdog-bin.outputs.cache-hit != 'true' @@ -36,16 +34,11 @@ jobs: - name: Build (debug) if: steps.pgdog-bin.outputs.cache-hit != 'true' run: cargo build --no-default-features --features new_parser --bin pgdog - - name: Build (release) - if: steps.pgdog-bin.outputs.cache-hit != 'true' - run: cargo build --release --no-default-features --features new_parser --bin pgdog - name: Save pgdog binaries if: steps.pgdog-bin.outputs.cache-hit != 'true' uses: actions/cache/save@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} ci: @@ -79,7 +72,7 @@ jobs: # cached pgdog binary. - { name: plugins, script: integration/plugins/run.sh, needs_rust_cache: true, continue_on_error: true } env: - PGDOG_BIN: ${{ github.workspace }}/target/release/pgdog + PGDOG_BIN: ${{ github.workspace }}/target/debug/pgdog PGDOG_PLUGIN_FEATURES: new_parser steps: - uses: actions/checkout@v6 @@ -89,8 +82,8 @@ jobs: id: cache-key run: echo "key=pgdog-bin-new-parser-${{ runner.os }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" # rust-cache must run before the binary restore: it lays down a - # stale target/ that can otherwise wipe target/release/pgdog when - # cargo reconciles fingerprints during plugin builds. + # stale target/ that can otherwise wipe target binaries when cargo + # reconciles fingerprints during plugin builds. - name: Restore Rust cache for plugin builds if: matrix.needs_rust_cache uses: Swatinem/rust-cache@v2 @@ -99,9 +92,7 @@ jobs: - name: Restore pgdog binaries uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} fail-on-cache-miss: true - name: Setup dependencies diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index c2544e4d9..6f803e378 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -7,8 +7,25 @@ on: types: [opened, synchronize, reopened] jobs: + # Build both pgdog artifacts in parallel: + # - release: plain release binary + # - coverage: llvm-cov instrumented debug binary used by the + # integration suites to produce line coverage for Codecov build: runs-on: blacksmith-4vcpu-ubuntu-2404 + strategy: + fail-fast: false + matrix: + include: + - artifact: release + path: target/release/pgdog + rust_prefix: build-v1 + - artifact: coverage + path: target/llvm-cov-target/debug/pgdog + # Separate prefix: instrumented artifacts live in + # target/llvm-cov-target and must not race with the plain + # build's cache. + rust_prefix: build-cov-debug-v1 steps: - uses: actions/checkout@v6 - name: Install CI deps @@ -17,35 +34,31 @@ jobs: # target/ build outputs that would otherwise change the key. Use the # split cache actions so pgdog binaries are saved immediately after the # build, before rust-cache's post-job cleanup prunes workspace binaries - # from target/{debug,release}. + # from target/. - name: Compute cache key id: cache-key - run: echo "key=pgdog-bin-${{ runner.os }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" - - name: Restore pgdog binaries + run: echo "key=pgdog-bin-${{ runner.os }}-${{ matrix.artifact }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" + - name: Restore pgdog binary id: pgdog-bin uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: ${{ matrix.path }} key: ${{ steps.cache-key.outputs.key }} - uses: Swatinem/rust-cache@v2 if: steps.pgdog-bin.outputs.cache-hit != 'true' with: - prefix-key: build-v1 - - name: Build (debug) - if: steps.pgdog-bin.outputs.cache-hit != 'true' - run: cargo build --bin pgdog + prefix-key: ${{ matrix.rust_prefix }} - name: Build (release) - if: steps.pgdog-bin.outputs.cache-hit != 'true' + if: steps.pgdog-bin.outputs.cache-hit != 'true' && matrix.artifact == 'release' run: cargo build --release --bin pgdog - - name: Save pgdog binaries + - name: Build (instrumented debug, for coverage) + if: steps.pgdog-bin.outputs.cache-hit != 'true' && matrix.artifact == 'coverage' + run: bash integration/ci/prepare-instrumented-pgdog.sh debug + - name: Save pgdog binary if: steps.pgdog-bin.outputs.cache-hit != 'true' uses: actions/cache/save@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: ${{ matrix.path }} key: ${{ steps.cache-key.outputs.key }} ci: @@ -79,30 +92,39 @@ jobs: # cached pgdog binary. - { name: plugins, script: integration/plugins/run.sh, needs_rust_cache: true, continue_on_error: true } env: - PGDOG_BIN: ${{ github.workspace }}/target/release/pgdog + # All suites run against the llvm-cov instrumented debug binary so we + # get line coverage from integration tests, not just unit tests. + PGDOG_BIN: ${{ github.workspace }}/target/llvm-cov-target/debug/pgdog + # NB: cargo llvm-cov report only globs profraw files at the top of + # target/llvm-cov-target, so LLVM_PROFILE_FILE must point there. + LLVM_PROFILE_FILE: ${{ github.workspace }}/target/llvm-cov-target/${{ matrix.name }}-%p-%m.profraw steps: - uses: actions/checkout@v6 - name: Install CI deps run: bash integration/ci/install-deps.sh - name: Compute cache key id: cache-key - run: echo "key=pgdog-bin-${{ runner.os }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" + run: echo "key=pgdog-bin-${{ runner.os }}-coverage-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" # rust-cache must run before the binary restore: it lays down a - # stale target/ that can otherwise wipe target/release/pgdog when - # cargo reconciles fingerprints during plugin builds. + # stale target/ that can otherwise wipe target binaries when cargo + # reconciles fingerprints during plugin builds. - name: Restore Rust cache for plugin builds if: matrix.needs_rust_cache uses: Swatinem/rust-cache@v2 with: prefix-key: build-v1 - - name: Restore pgdog binaries + - name: Restore instrumented pgdog binary uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/llvm-cov-target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} fail-on-cache-miss: true + # The Rust cache (plugins job) can lay down stale .profraw files that + # don't match the restored instrumented binary; start clean. + - name: Prepare coverage profile dir + run: | + mkdir -p target/llvm-cov-target + rm -f target/llvm-cov-target/*.profraw target/llvm-cov-target/*.profdata - name: Setup dependencies run: bash integration/ci/setup.sh --with-toxi - name: Run ${{ matrix.name }} @@ -110,3 +132,18 @@ jobs: - name: Ensure PgDog stopped if: always() run: bash integration/ci/ensure-pgdog-stopped.sh + - name: Generate coverage report + if: always() + run: cargo llvm-cov report --package pgdog --lcov --output-path ${{ matrix.name }}.lcov + # Codecov merges uploads per flag server-side, so each suite uploads + # its own lcov under the shared "integration" flag. + - name: Upload coverage to Codecov + if: always() + uses: codecov/codecov-action@v4 + env: + CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} + with: + files: ${{ matrix.name }}.lcov + flags: integration + name: integration-${{ matrix.name }} + fail_ci_if_error: false diff --git a/integration/ci/install-deps.sh b/integration/ci/install-deps.sh index 8716c7a25..b30cea9f0 100755 --- a/integration/ci/install-deps.sh +++ b/integration/ci/install-deps.sh @@ -54,3 +54,10 @@ if ! command -v cargo-llvm-cov >/dev/null; then curl -LsSf "https://github.com/taiki-e/cargo-llvm-cov/releases/download/v${LLVM_COV_VERSION}/cargo-llvm-cov-x86_64-unknown-linux-gnu.tar.gz" \ | tar zxf - -C "$CARGO_BIN" fi + +# cargo llvm-cov report needs llvm-profdata/llvm-cov from llvm-tools-preview. +if command -v rustup >/dev/null 2>&1; then + if ! rustup component list --installed 2>/dev/null | grep -q llvm-tools; then + rustup component add llvm-tools-preview || true + fi +fi diff --git a/integration/ci/prepare-instrumented-pgdog.sh b/integration/ci/prepare-instrumented-pgdog.sh index eeea163fe..38760a4d5 100755 --- a/integration/ci/prepare-instrumented-pgdog.sh +++ b/integration/ci/prepare-instrumented-pgdog.sh @@ -1,20 +1,36 @@ #!/usr/bin/env bash # Build an llvm-cov-instrumented pgdog binary for the integration job and # export its path via $GITHUB_ENV so later steps can invoke it. +# +# Usage: prepare-instrumented-pgdog.sh [debug|release] (default: release) set -euo pipefail -export RUSTFLAGS="-C link-dead-code" +PROFILE="${1:-release}" +if [[ "${PROFILE}" != "debug" && "${PROFILE}" != "release" ]]; then + echo "Usage: $0 [debug|release]" >&2 + exit 1 +fi + +PROFILE_FLAG="" +if [[ "${PROFILE}" == "release" ]]; then + PROFILE_FLAG="--release" +fi + +# RUSTFLAGS overrides .cargo/config.toml rustflags entirely, so +# tokio_unstable (set there) must be repeated here. +export RUSTFLAGS="--cfg tokio_unstable -C link-dead-code" cargo llvm-cov clean --workspace mkdir -p target/llvm-cov-target/profiles -cargo llvm-cov run --no-report --release --package pgdog --bin pgdog -- --help +# shellcheck disable=SC2086 +cargo llvm-cov run --no-report ${PROFILE_FLAG} --package pgdog --bin pgdog -- --help rm -f target/llvm-cov-target/profiles/*.profraw rm -f target/llvm-cov-target/profiles/.last_snapshot rm -rf target/llvm-cov-target/reports -BIN_PATH=$(find target/llvm-cov-target -type f -path '*/release/pgdog' | head -n 1) +BIN_PATH=$(find target/llvm-cov-target -type f -path "*/${PROFILE}/pgdog" | head -n 1) if [ -z "$BIN_PATH" ]; then - echo "Instrumented PgDog binary not found" >&2 + echo "Instrumented PgDog binary (${PROFILE}) not found" >&2 exit 1 fi echo "Using instrumented binary at $BIN_PATH" diff --git a/pgdog/src/backend/databases.rs b/pgdog/src/backend/databases.rs index f3d6b762b..11d553d94 100644 --- a/pgdog/src/backend/databases.rs +++ b/pgdog/src/backend/databases.rs @@ -467,6 +467,9 @@ impl Databases { cluster.user(), cluster.name() ); + // No boot-time maintenance will run, don't block + // readiness waiters. Checkouts will fail instead. + cluster.mark_ready(); } else { cluster.launch(); } diff --git a/pgdog/src/backend/pool/cluster.rs b/pgdog/src/backend/pool/cluster.rs index 78f39af2e..ff843a61c 100644 --- a/pgdog/src/backend/pool/cluster.rs +++ b/pgdog/src/backend/pool/cluster.rs @@ -6,22 +6,13 @@ use pgdog_config::{ LoadSchema, PreparedStatements, QueryParser, QueryParserEngine, QueryParserLevel, Rewrite, RewriteMode, users::PasswordKind, }; -use std::{ - sync::{ - Arc, - atomic::{AtomicBool, Ordering}, - }, - time::Duration, -}; -use tracing::error; +use std::{sync::Arc, time::Duration}; use crate::frontend::router::sharding::ShardedTable; -use crate::tasks; use crate::{ backend::{ Schema, ShardedTables, databases::{User as DatabaseUser, databases}, - pool::ee::schema_changed_hook, replication::{ReplicationConfig, ShardedSchemas}, }, config::{ @@ -31,7 +22,10 @@ use crate::{ net::{Query, messages::FrontendPid}, }; -use super::{Address, Config, Error, Guard, MirrorStats, Request, Shard, ShardConfig}; +use super::{ + Address, Config, Error, Guard, MirrorStats, Request, Shard, ShardConfig, + cluster_launch::Readiness, +}; use crate::config::LoadBalancingStrategy; #[derive(Clone, Debug, Default)] @@ -43,11 +37,6 @@ pub struct PoolConfig { pub(crate) config: Config, } -#[derive(Default, Debug)] -struct Readiness { - online: AtomicBool, -} - /// A collection of sharded replicas and primaries /// belonging to the same database cluster. #[derive(Clone, Default, Debug)] @@ -67,7 +56,7 @@ pub struct Cluster { cross_shard_disabled: bool, two_phase_commit: bool, two_phase_commit_auto: bool, - readiness: Arc, + pub(super) readiness: Arc, rewrite: Rewrite, prepared_statements: PreparedStatements, dry_run: bool, @@ -561,7 +550,7 @@ impl Cluster { self.reload_schema_on_ddl && self.load_schema() } - fn load_schema(&self) -> bool { + pub(super) fn load_schema(&self) -> bool { match self.load_schema { LoadSchema::On => true, LoadSchema::Off => false, @@ -635,65 +624,6 @@ impl Cluster { self.resharding_replication_retry_min_delay } - /// Launch the connection pools. - pub(crate) fn launch(&self) { - for shard in self.shards() { - shard.launch(); - } - - self.readiness.online.store(true, Ordering::Relaxed); - - if !self.load_schema() { - for shard in &self.shards { - shard.schema_not_needed(); - } - return; - } - - for shard in self.shards() { - let identifier = self.identifier(); - let shard = shard.clone(); - - tasks::spawn("cluster schema sync", async move { - use tokio::time::sleep; - - loop { - match shard.load_schema().await { - Ok(true) => { - schema_changed_hook(&shard.schema(), &identifier, &shard); - return; - } - Ok(false) => return, - Err(err) => { - if shard.online() { - error!( - "error loading schema for shard {}: {}", - shard.number(), - err - ); - sleep(Duration::from_millis(100)).await; - } else { - // Cluster is shutting down: unblock any - // wait_schema_loaded callers. - shard.schema_not_needed(); - return; - } - } - } - } - }); - } - } - - /// Shutdown the connection pools. - pub(crate) fn shutdown(&self) { - for shard in self.shards() { - shard.shutdown(); - } - - self.readiness.online.store(false, Ordering::Relaxed); - } - /// Send a cancellation request for all running queries. pub(crate) async fn cancel_all(&self) -> Result<(), Error> { let pools: Vec<_> = self @@ -709,22 +639,6 @@ impl Cluster { Ok(()) } - /// Is the cluster online? - pub(crate) fn online(&self) -> bool { - self.readiness.online.load(Ordering::Relaxed) - } - - /// Schema loaded for all shards? - pub(crate) async fn wait_schema_loaded(&self) { - if !self.load_schema() { - return; - } - - for shard in &self.shards { - shard.wait_schema_loaded().await; - } - } - /// Execute a query on every primary in the cluster. pub async fn execute( &self, @@ -1075,83 +989,107 @@ mod test { } #[tokio::test] - async fn test_launch_schema_loading_idempotent() { - use tokio::time::{Duration, sleep}; - + async fn test_launch_marks_ready() { let config = ConfigAndUsers::default(); - let mut cluster = Cluster::new_test(&config); - cluster.sharded_schemas = ShardedSchemas::default(); - - assert!(cluster.load_schema()); + let cluster = Cluster::new_test(&config); - // Pre-populate per-shard schemas so launch's spawned loaders become - // no-ops (Shard::load_schema returns Ok(false) when already initialized). - // This avoids touching a real database in the test. + // Pre-populate per-shard schemas so launch doesn't touch a real database. for shard in &cluster.shards { shard.schema_not_needed(); } + assert!(!cluster.ready()); cluster.launch(); - cluster.wait_schema_loaded().await; - - // Second launch must be safe: per-shard OnceCell prevents any reload. - cluster.launch(); - sleep(Duration::from_millis(50)).await; - cluster.wait_schema_loaded().await; + cluster.wait_ready().await; + assert!(cluster.ready()); } #[tokio::test] - async fn test_wait_schema_loaded_returns_immediately_when_not_needed() { + async fn test_shutdown_releases_readiness_waiters() { + use tokio::time::{Duration, timeout}; + let config = ConfigAndUsers::default(); - let cluster = Cluster::new_test_single_shard(&config); + let cluster = Cluster::new_test(&config); - // load_schema() returns false for single shard without multi_tenant - assert!(!cluster.load_schema()); + assert!(!cluster.ready()); + + let waiter = cluster.clone(); + let handle = tokio::spawn(async move { + waiter.wait_ready().await; + }); + + cluster.shutdown(); - // Should return immediately without waiting - cluster.wait_schema_loaded().await; + let result = timeout(Duration::from_millis(200), handle).await; + assert!(result.is_ok()); } #[tokio::test] - async fn test_wait_schema_loaded_fast_path_when_already_loaded() { + async fn test_launch_schema_loading_idempotent() { + use tokio::time::{Duration, sleep}; + let config = ConfigAndUsers::default(); let mut cluster = Cluster::new_test(&config); cluster.sharded_schemas = ShardedSchemas::default(); assert!(cluster.load_schema()); - // Mark every shard's schema as already loaded. + // Pre-populate per-shard schemas so launch's spawned loaders become + // no-ops (Shard::load_schema returns Ok(false) when already initialized). + // This avoids touching a real database in the test. for shard in &cluster.shards { shard.schema_not_needed(); } - // Each shard's wait_schema_loaded() takes the fast path and returns - // immediately because schema.initialized() is true. - cluster.wait_schema_loaded().await; + cluster.launch(); + cluster.wait_ready().await; + + // Second launch must be safe: per-shard OnceCell prevents any reload. + cluster.launch(); + sleep(Duration::from_millis(50)).await; + cluster.wait_ready().await; } #[tokio::test] - async fn test_wait_schema_loaded_waits_for_notification() { - use tokio::time::{Duration, timeout}; + async fn test_wait_ready_waits_for_schema_notification() { + use tokio::time::{Duration, sleep, timeout}; let config = ConfigAndUsers::default(); - let mut cluster = Cluster::new_test(&config); - cluster.sharded_schemas = ShardedSchemas::default(); + let cluster = Cluster::new_test(&config); - assert!(cluster.load_schema()); + cluster.launch(); - // Trigger schema_not_needed on each shard after a short delay so the - // waiter wakes up via the per-shard schema_waiter notification. + // Schemas not loaded yet, readiness is pending. + assert!(!cluster.ready()); + + // Simulate schema load finishing on each shard: the readiness + // monitor wakes up via the per-shard schema_waiter notification. let shards: Vec<_> = cluster.shards.to_vec(); tokio::spawn(async move { - tokio::time::sleep(Duration::from_millis(10)).await; + sleep(Duration::from_millis(10)).await; for shard in &shards { shard.schema_not_needed(); } }); - let result = timeout(Duration::from_millis(200), cluster.wait_schema_loaded()).await; + let result = timeout(Duration::from_millis(500), cluster.wait_ready()).await; assert!(result.is_ok()); + assert!(cluster.ready()); + } + + #[tokio::test] + async fn test_wait_ready_returns_immediately_when_schema_not_needed() { + let config = ConfigAndUsers::default(); + let cluster = Cluster::new_test_single_shard(&config); + + // load_schema() returns false for single shard without multi_tenant + assert!(!cluster.load_schema()); + + cluster.launch(); + + // Should return without waiting: no schema to load. + cluster.wait_ready().await; + assert!(cluster.ready()); } #[test] diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs new file mode 100644 index 000000000..a8536913f --- /dev/null +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -0,0 +1,167 @@ +//! Cluster startup and shutdown primitives. +//! +//! Launching and shutting down the connection pools, and gating +//! traffic until boot-time maintenance has completed. + +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; + +use tokio::{select, time::sleep}; +use tokio_util::sync::CancellationToken; +use tracing::error; + +use crate::backend::pool::ee::schema_changed_hook; +use crate::tasks; + +use super::Cluster; + +/// Cluster readiness state. +#[derive(Default, Debug)] +pub(super) struct Readiness { + online: AtomicBool, + launch_waiter: CancellationToken, +} + +impl Readiness { + fn set_online(&self, online: bool) { + self.online.store(online, Ordering::Relaxed); + } + + fn online(&self) -> bool { + self.online.load(Ordering::Relaxed) + } + + fn mark_ready(&self) { + self.launch_waiter.cancel(); + } + + #[cfg(test)] + fn ready(&self) -> bool { + self.launch_waiter.is_cancelled() + } + + async fn wait_ready(&self) { + self.launch_waiter.cancelled().await; + } +} + +impl Cluster { + /// Launch the connection pools. + pub(crate) fn launch(&self) { + for shard in self.shards() { + shard.launch(); + } + + self.readiness.set_online(true); + + self.launch_schema_sync(); + self.launch_readiness_monitor(); + } + + /// Shutdown the connection pools. + pub(crate) fn shutdown(&self) { + self.readiness.set_online(false); + + for shard in self.shards() { + shard.shutdown(); + } + + // Release readiness waiters, the cluster is going away. + self.mark_ready(); + } + + /// Is the cluster online? + pub(crate) fn online(&self) -> bool { + self.readiness.online() + } + + /// Wait until boot-time maintenance is done and the cluster can serve traffic. + pub(crate) async fn wait_ready(&self) { + self.readiness.wait_ready().await; + } + + /// Boot-time maintenance is done and the cluster can serve traffic. + #[cfg(test)] + pub(crate) fn ready(&self) -> bool { + self.readiness.ready() + } + + /// Mark the cluster as ready to serve traffic. + pub(crate) fn mark_ready(&self) { + self.readiness.mark_ready(); + } + + /// Mark the cluster ready schema loading is done. + fn launch_readiness_monitor(&self) { + let cluster = self.clone(); + tasks::spawn("cluster readiness monitor", async move { + let shutdown = tasks::shutdown_signal(); + select! { + _ = cluster.wait_schema_loaded() => {} + _ = shutdown.cancelled() => {} + } + cluster.mark_ready(); + }); + } + + /// Wait for the schema to load on all shards. + async fn wait_schema_loaded(&self) { + if !self.load_schema() { + return; + } + + for shard in self.shards() { + shard.wait_schema_loaded().await; + } + } + + /// Load database schema in the background, retrying + /// until success or shutdown. + fn launch_schema_sync(&self) { + if !self.load_schema() { + for shard in self.shards() { + shard.schema_not_needed(); + } + return; + } + + for shard in self.shards() { + let identifier = self.identifier(); + let shard = shard.clone(); + let shutdown = tasks::shutdown_signal(); + + tasks::spawn("shard schema sync", async move { + loop { + let loader = shard.load_schema(); + let result = select! { + _ = shutdown.cancelled() => break, + result = loader => { result }, + }; + + match result { + Ok(true) => { + schema_changed_hook(&shard.schema(), &identifier, &shard); + return; + } + Ok(false) => return, + Err(err) => { + if shard.online() { + error!( + "error loading schema for shard {}: {}", + shard.number(), + err + ); + sleep(Duration::from_millis(100)).await; + } else { + // Cluster is shutting down: unblock any + // wait_schema_loaded callers. + shard.schema_not_needed(); + return; + } + } + } + } + }); + } + } +} diff --git a/pgdog/src/backend/pool/connection/binding.rs b/pgdog/src/backend/pool/connection/binding.rs index cf3fb0d77..a762b248d 100644 --- a/pgdog/src/backend/pool/connection/binding.rs +++ b/pgdog/src/backend/pool/connection/binding.rs @@ -3,7 +3,10 @@ use crate::{ frontend::{ ClientRequest, - client::query_engine::{TwoPcPhase, two_pc::statement::phase_control}, + client::query_engine::{ + TwoPcPhase, + two_pc::{TwoPcTransaction, statement::phase_control}, + }, }, net::{FrontendPid, ProtocolMessage, Query, parameter::Parameters}, state::State, @@ -361,14 +364,14 @@ impl Binding { pub(crate) async fn two_pc_on_guards( servers: &mut [Guard], - name: &str, + transaction: TwoPcTransaction, phase: TwoPcPhase, ) -> Result<(), Error> { let skip_missing = matches!(phase, TwoPcPhase::Phase2 | TwoPcPhase::Rollback); let mut futures = Vec::new(); for (shard, server) in servers.iter_mut().enumerate() { - let query = phase_control(name, shard, phase); + let query = phase_control(transaction, shard, phase); futures.push(server.execute(query)); } @@ -394,9 +397,15 @@ impl Binding { } /// Execute two-phase commit transaction control statements. - pub async fn two_pc(&mut self, name: &str, phase: TwoPcPhase) -> Result<(), Error> { + pub(crate) async fn two_pc( + &mut self, + transaction: TwoPcTransaction, + phase: TwoPcPhase, + ) -> Result<(), Error> { match self { - Binding::MultiShard(servers, _) => Self::two_pc_on_guards(servers, name, phase).await, + Binding::MultiShard(servers, _) => { + Self::two_pc_on_guards(servers, transaction, phase).await + } _ => Err(Error::TwoPcMultiShardOnly), } diff --git a/pgdog/src/backend/pool/connection/binding_test.rs b/pgdog/src/backend/pool/connection/binding_test.rs index 8fb3f300e..0a6258b6e 100644 --- a/pgdog/src/backend/pool/connection/binding_test.rs +++ b/pgdog/src/backend/pool/connection/binding_test.rs @@ -8,7 +8,7 @@ mod tests { server::test::test_server, }, frontend::{ - client::query_engine::TwoPcPhase, + client::query_engine::{TwoPcPhase, two_pc::TwoPcTransaction}, router::{ Route, parser::{Shard, ShardWithPriority}, @@ -76,7 +76,9 @@ mod tests { let guard = crate::backend::pool::Guard::new(pool, server, Instant::now()); let mut binding = Binding::Direct(guard, 0); - let result = binding.two_pc("test", TwoPcPhase::Phase1).await; + let result = binding + .two_pc(TwoPcTransaction::new(), TwoPcPhase::Phase1) + .await; // Should fail with TwoPcMultiShardOnly error assert!(result.is_err()); @@ -94,7 +96,9 @@ mod tests { let admin_server = AdminServer::default(); let mut binding = Binding::Admin(admin_server); - let result = binding.two_pc("test", TwoPcPhase::Phase1).await; + let result = binding + .two_pc(TwoPcTransaction::new(), TwoPcPhase::Phase1) + .await; // Should fail with TwoPcMultiShardOnly error assert!(result.is_err()); @@ -106,10 +110,10 @@ mod tests { #[tokio::test] async fn test_two_pc_phase1_prepare() { let mut binding = create_multishard_binding().await; - let transaction_name = "test_txn"; + let transaction = TwoPcTransaction::new(); // Test Phase1 - PREPARE TRANSACTION - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; // Should succeed if let Err(ref error) = result { @@ -118,60 +122,60 @@ mod tests { assert!(result.is_ok()); // Cleanup: Rollback the prepared transaction to avoid leaving dangling transactions - let _cleanup = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let _cleanup = binding.two_pc(transaction, TwoPcPhase::Rollback).await; } #[tokio::test] async fn test_two_pc_phase2_commit() { let mut binding = create_multishard_binding().await; - let transaction_name = "test_commit_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then commit it - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!(result.is_ok()); } #[tokio::test] async fn test_two_pc_rollback() { let mut binding = create_multishard_binding().await; - let transaction_name = "test_rollback_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then rollback - let result = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await; assert!(result.is_ok()); } #[tokio::test] async fn test_two_pc_commit_after_prepare_and_commit() { let mut binding = create_multishard_binding().await; - let transaction_name = "committed_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then commit it binding - .two_pc(transaction_name, TwoPcPhase::Phase2) + .two_pc(transaction, TwoPcPhase::Phase2) .await .expect("Phase2 should succeed"); // Try to commit again - should succeed because skip_missing is true for Phase2 - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!( result.is_ok(), "Committing non-existent prepared transaction should be skipped" @@ -181,22 +185,22 @@ mod tests { #[tokio::test] async fn test_two_pc_rollback_after_prepare_and_rollback() { let mut binding = create_multishard_binding().await; - let transaction_name = "rolled_back_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then rollback it binding - .two_pc(transaction_name, TwoPcPhase::Rollback) + .two_pc(transaction, TwoPcPhase::Rollback) .await .expect("Rollback should succeed"); // Try to rollback again - should succeed because skip_missing is true for Rollback - let result = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await; assert!( result.is_ok(), "Rolling back non-existent prepared transaction should be skipped" @@ -209,22 +213,22 @@ mod tests { #[tokio::test] async fn test_two_pc_transaction_lifecycle() { let mut binding = create_multishard_binding().await; - let transaction_name = "lifecycle_test"; + let transaction = TwoPcTransaction::new(); // 1. Prepare transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; assert!(result.is_ok(), "Phase1 preparation should succeed"); // 2. Try to prepare the same transaction again - PostgreSQL behavior may vary - let _result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let _result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; // Note: PostgreSQL behavior for duplicate PREPARE TRANSACTION can vary depending on context // 3. Commit the prepared transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!(result.is_ok(), "Phase2 commit should succeed"); // 4. Try to commit again - should succeed (skip_missing = true) - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!( result.is_ok(), "Committing non-existent transaction should be skipped" @@ -234,18 +238,18 @@ mod tests { #[tokio::test] async fn test_two_pc_prepare_then_rollback() { let mut binding = create_multishard_binding().await; - let transaction_name = "prepare_rollback_test"; + let transaction = TwoPcTransaction::new(); // 1. Prepare transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; assert!(result.is_ok(), "Phase1 preparation should succeed"); // 2. Rollback the prepared transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await; assert!(result.is_ok(), "Rollback should succeed"); // 3. Try to commit after rollback - should succeed (skip_missing = true) - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!( result.is_ok(), "Committing rolled back transaction should be skipped" diff --git a/pgdog/src/backend/pool/mod.rs b/pgdog/src/backend/pool/mod.rs index ce5635800..ef7ecce01 100644 --- a/pgdog/src/backend/pool/mod.rs +++ b/pgdog/src/backend/pool/mod.rs @@ -3,6 +3,7 @@ pub mod address; pub mod cleanup; pub mod cluster; +pub mod cluster_launch; pub mod comms; pub mod config; pub mod connection; diff --git a/pgdog/src/backend/replication/logical/orchestrator.rs b/pgdog/src/backend/replication/logical/orchestrator.rs index 598c1fb82..e09a20497 100644 --- a/pgdog/src/backend/replication/logical/orchestrator.rs +++ b/pgdog/src/backend/replication/logical/orchestrator.rs @@ -121,7 +121,7 @@ impl Orchestrator { self.destination = databases().schema_owner(&self.destination.identifier().database)?; self.source = databases().schema_owner(&self.source.identifier().database)?; - self.destination.wait_schema_loaded().await; + self.destination.wait_ready().await; self.refresh_publisher(); diff --git a/pgdog/src/backend/replication/logical/subscriber/copy.rs b/pgdog/src/backend/replication/logical/subscriber/copy.rs index eeeccdb86..6b7c5eff0 100644 --- a/pgdog/src/backend/replication/logical/subscriber/copy.rs +++ b/pgdog/src/backend/replication/logical/subscriber/copy.rs @@ -296,14 +296,14 @@ impl CopySubscriber { async { let _guard_phase_1 = manager - .transaction_state(&txn, &identifier, TwoPcPhase::Phase1) + .transaction_state(txn, &identifier, TwoPcPhase::Phase1) .await?; - self.two_pc_on_shards(&txn, TwoPcPhase::Phase1).await?; + self.two_pc_on_shards(txn, TwoPcPhase::Phase1).await?; let _guard_phase_2 = manager - .transaction_state(&txn, &identifier, TwoPcPhase::Phase2) + .transaction_state(txn, &identifier, TwoPcPhase::Phase2) .await?; - self.two_pc_on_shards(&txn, TwoPcPhase::Phase2).await?; + self.two_pc_on_shards(txn, TwoPcPhase::Phase2).await?; manager.done(&txn).await?; Ok(()) @@ -317,7 +317,7 @@ impl CopySubscriber { async fn two_pc_on_shards( &mut self, - txn: &TwoPcTransaction, + txn: TwoPcTransaction, phase: TwoPcPhase, ) -> Result<(), Error> { let mut futures = Vec::new(); @@ -329,7 +329,7 @@ impl CopySubscriber { // via binding.rs using the same phase_control() helper. let query = match phase { TwoPcPhase::Rollback => unreachable!(), - phase => phase_control(&txn.to_string(), shard, phase), + phase => phase_control(txn, shard, phase), }; futures.push(Self::send_and_confirm(server, Query::new(query).into())); } diff --git a/pgdog/src/config/mod.rs b/pgdog/src/config/mod.rs index 22cc79d10..4aa140ae9 100644 --- a/pgdog/src/config/mod.rs +++ b/pgdog/src/config/mod.rs @@ -154,6 +154,16 @@ pub fn load_test() { #[cfg(test)] pub fn load_test_with_pooler_mode(pooler_mode: PoolerMode) { + load_test_with_user_and_pooler_mode("pgdog", pooler_mode, Role::default()) +} + +#[cfg(test)] +pub fn load_test_with_user(user: &str) { + load_test_with_user_and_pooler_mode(user, PoolerMode::Transaction, Role::Primary) +} + +#[cfg(test)] +fn load_test_with_user_and_pooler_mode(user: &str, pooler_mode: PoolerMode, role: Role) { use crate::backend::databases::init; let mut config = ConfigAndUsers::default(); @@ -162,10 +172,11 @@ pub fn load_test_with_pooler_mode(pooler_mode: PoolerMode) { host: "127.0.0.1".into(), port: 5432, pooler_mode: Some(pooler_mode), + role, ..Default::default() }]; config.users.users = vec![User { - name: "pgdog".into(), + name: user.into(), database: "pgdog".into(), password: Some("pgdog".into()), pooler_mode: Some(pooler_mode), diff --git a/pgdog/src/frontend/client/query_engine/end_transaction.rs b/pgdog/src/frontend/client/query_engine/end_transaction.rs index fc74b149f..affcc9053 100644 --- a/pgdog/src/frontend/client/query_engine/end_transaction.rs +++ b/pgdog/src/frontend/client/query_engine/end_transaction.rs @@ -106,17 +106,17 @@ impl QueryEngine { } let identifier = cluster.identifier(); - let name = self.two_pc.transaction().to_string(); + let transaction = self.two_pc.transaction(); // If interrupted here, the transaction must be rolled back. let _guard_phase_1 = self.two_pc.phase_one(&identifier).await?; - self.backend.two_pc(&name, TwoPcPhase::Phase1).await?; + self.backend.two_pc(transaction, TwoPcPhase::Phase1).await?; debug!("[2pc] phase 1 complete"); // If interrupted here, the transaction must be committed. let _guard_phase_2 = self.two_pc.phase_two(&identifier).await?; - self.backend.two_pc(&name, TwoPcPhase::Phase2).await?; + self.backend.two_pc(transaction, TwoPcPhase::Phase2).await?; debug!("[2pc] phase 2 complete"); diff --git a/pgdog/src/frontend/client/query_engine/route_query.rs b/pgdog/src/frontend/client/query_engine/route_query.rs index 292dfde7f..74d6be3e1 100644 --- a/pgdog/src/frontend/client/query_engine/route_query.rs +++ b/pgdog/src/frontend/client/query_engine/route_query.rs @@ -50,15 +50,14 @@ impl QueryEngine { }; if let Ok(ClusterCheck::Ok) = res { - // Make sure schema is loaded before we throw traffic - // at it. This matters for sharded deployments only. + // Wait for boot-time maintenance before we throw traffic at the cluster. if let Ok(cluster) = self.backend.cluster() { safe_timeout( context.timeouts.query_timeout(&State::Active), - cluster.wait_schema_loaded(), + cluster.wait_ready(), ) .await - .map_err(|_| Error::SchemaLoad)?; + .map_err(|_| Error::ClusterStart)?; } res } else { diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index 94dd9fa36..ff22ffb23 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -201,14 +201,14 @@ impl Manager { /// the WAL exists to prevent. pub(crate) async fn transaction_state( &self, - transaction: &TwoPcTransaction, + transaction: TwoPcTransaction, identifier: &Arc, phase: TwoPcPhase, ) -> Result { let prior = { let mut guard = self.inner.lock(); - let prior = guard.transactions.get(transaction).cloned(); - let entry = guard.transactions.entry(*transaction).or_default(); + let prior = guard.transactions.get(&transaction).cloned(); + let entry = guard.transactions.entry(transaction).or_default(); entry.identifier = identifier.clone(); entry.phase = phase; prior @@ -218,13 +218,13 @@ impl Manager { let result = match phase { TwoPcPhase::Phase1 => { wal.append_begin( - *transaction, + transaction, identifier.user.clone(), identifier.database.clone(), ) .await } - TwoPcPhase::Phase2 => wal.append_committing(*transaction).await, + TwoPcPhase::Phase2 => wal.append_committing(transaction).await, TwoPcPhase::Rollback => { unreachable!("rollback is not a state transition; it's the cleanup direction") } @@ -233,10 +233,10 @@ impl Manager { let mut guard = self.inner.lock(); match prior { Some(prior) => { - guard.transactions.insert(*transaction, prior); + guard.transactions.insert(transaction, prior); } None => { - guard.transactions.remove(transaction); + guard.transactions.remove(&transaction); } } warn!( @@ -248,7 +248,7 @@ impl Manager { } Ok(TwoPcGuard { - transaction: *transaction, + transaction, manager: Self::get(), }) } @@ -309,7 +309,7 @@ impl Manager { r#"[2pc] cleaning up transaction "{}""#, transaction.to_string() ); - match manager.cleanup_phase(&transaction).await { + match manager.cleanup_phase(transaction).await { Err(err) => { error!( r#"[2pc] error cleaning up "{}" transaction: {}"#, @@ -345,8 +345,8 @@ impl Manager { } /// Reconnect to cluster if available and rollback the two-phase transaction. - async fn cleanup_phase(&self, transaction: &TwoPcTransaction) -> Result<(), Error> { - let state = match self.inner.lock().transactions.get(transaction).cloned() { + async fn cleanup_phase(&self, transaction: TwoPcTransaction) -> Result<(), Error> { + let state = match self.inner.lock().transactions.get(&transaction).cloned() { Some(state) => state, _ => { return Ok(()); @@ -379,7 +379,7 @@ impl Manager { &Route::write(ShardWithPriority::new_override_transaction(Shard::All)), ) .await?; - connection.two_pc(&transaction.to_string(), phase).await?; + connection.two_pc(transaction, phase).await?; connection.disconnect(); Ok(()) diff --git a/pgdog/src/frontend/client/query_engine/two_pc/mod.rs b/pgdog/src/frontend/client/query_engine/two_pc/mod.rs index 1024d12ab..b0cba45aa 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/mod.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/mod.rs @@ -18,6 +18,7 @@ pub use guard::TwoPcGuard; pub use manager::Manager; pub use phase::TwoPcPhase; pub(crate) use server_transactions::TwoPcTransactions; +pub(crate) use statement::TwoPcTransactionOnShard; pub use stats::TwoPcStats; pub use transaction::TwoPcTransaction; @@ -58,7 +59,7 @@ impl TwoPc { pub(super) async fn phase_one(&mut self, cluster: &Arc) -> Result { let transaction = self.transaction(); self.manager - .transaction_state(&transaction, cluster, TwoPcPhase::Phase1) + .transaction_state(transaction, cluster, TwoPcPhase::Phase1) .await } @@ -68,7 +69,7 @@ impl TwoPc { pub(super) async fn phase_two(&mut self, cluster: &Arc) -> Result { let transaction = self.transaction(); self.manager - .transaction_state(&transaction, cluster, TwoPcPhase::Phase2) + .transaction_state(transaction, cluster, TwoPcPhase::Phase2) .await } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs index 8e95be6f0..b2a3445ef 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs @@ -1,16 +1,24 @@ -use std::{ops::Deref, str::FromStr}; +use std::ops::Deref; use crate::{ backend::{Error, Server}, net::DataRow, }; -use super::TwoPcTransaction; +use super::{TwoPcTransaction, TwoPcTransactionOnShard}; #[derive(Debug, Clone)] pub enum TwoPcServerTransaction { - Ours(TwoPcTransaction), - Other { name: String }, + Ours { + txn: TwoPcTransaction, + user: String, + database: String, + }, + Other { + name: String, + user: String, + database: String, + }, } pub(crate) struct TwoPcTransactions { @@ -28,20 +36,32 @@ impl Deref for TwoPcTransactions { impl TwoPcTransactions { /// Load two phase transactions from the server. pub(crate) async fn load(server: &mut Server) -> Result { - let records: Vec = server.fetch_all("SELECT * FROM pg_prepared_xacts").await?; + let records: Vec = server + .fetch_all("SELECT gid, owner, database FROM pg_prepared_xacts") + .await?; let mut transactions = vec![]; for record in records { - let transaction = record.get_text(1).map(|name| { - if let Ok(ours) = TwoPcTransaction::from_str(&name) { - TwoPcServerTransaction::Ours(ours) + let gid = record.get_text(0); + let user = record.get_text(1).unwrap_or_default(); + let database = record.get_text(2).unwrap_or_default(); + + if let Some(gid) = gid { + let txn = if let Ok(txn) = gid.parse::() { + TwoPcServerTransaction::Ours { + txn: txn.transaction(), + user, + database, + } } else { - TwoPcServerTransaction::Other { name } - } - }); + TwoPcServerTransaction::Other { + name: gid, + user, + database, + } + }; - if let Some(transaction) = transaction { - transactions.push(transaction); + transactions.push(txn); } } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs index 4456db029..9c59824e4 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs @@ -1,21 +1,63 @@ //! Per-shard two-phase commit transaction names and control statements. +use std::{fmt::Display, str::FromStr}; + +use crate::frontend::client::query_engine::two_pc::TwoPcTransaction; + use super::TwoPcPhase; -/// Prepared transaction name for a coordinator transaction on one shard. -pub fn shard_name(transaction: &str, shard: usize) -> String { - format!("{transaction}_{shard}") +/// 2pc transaction executed on a shard. We +/// make them unique per shard in case two or more +/// shards are located on the same postgres server. +pub(crate) struct TwoPcTransactionOnShard { + transaction: TwoPcTransaction, + shard: usize, +} + +impl TwoPcTransactionOnShard { + /// Create new 2pc transaction on shard x. + pub(crate) fn new(transaction: TwoPcTransaction, shard: usize) -> Self { + Self { transaction, shard } + } + + /// Get the coordinator transaction. + pub(crate) fn transaction(&self) -> TwoPcTransaction { + self.transaction + } +} + +impl Display for TwoPcTransactionOnShard { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{}_{}", self.transaction, self.shard) + } +} + +impl FromStr for TwoPcTransactionOnShard { + type Err = (); + + fn from_str(s: &str) -> Result { + let (transaction, shard) = s.rsplit_once('_').ok_or(())?; + + Ok(Self { + transaction: transaction.parse()?, + shard: shard.parse().map_err(|_| ())?, + }) + } } /// Build `PREPARE TRANSACTION`, `COMMIT PREPARED`, or `ROLLBACK PREPARED` /// for a shard participant. -pub fn phase_control(transaction: &str, shard: usize, phase: TwoPcPhase) -> String { - let name = shard_name(transaction, shard); +pub(crate) fn phase_control( + transaction: TwoPcTransaction, + shard: usize, + phase: TwoPcPhase, +) -> String { + let txn = TwoPcTransactionOnShard::new(transaction, shard); match phase { - TwoPcPhase::Phase1 => format!("PREPARE TRANSACTION '{name}'"), - TwoPcPhase::Phase2 => format!("COMMIT PREPARED '{name}'"), - TwoPcPhase::Rollback => format!("ROLLBACK PREPARED '{name}'"), + TwoPcPhase::Phase1 => format!("PREPARE TRANSACTION '{txn}'"), + TwoPcPhase::Phase2 => format!("COMMIT PREPARED '{txn}'"), + TwoPcPhase::Rollback => format!("ROLLBACK PREPARED '{txn}'"), } } @@ -24,24 +66,56 @@ mod test { use super::*; #[test] - fn shard_name_appends_index() { - assert_eq!(shard_name("__pgdog_2pc_42", 0), "__pgdog_2pc_42_0"); - assert_eq!(shard_name("test_txn", 3), "test_txn_3"); + fn transaction_on_shard_appends_index() { + let transaction = TwoPcTransaction::new(); + + assert_eq!( + TwoPcTransactionOnShard::new(transaction, 0).to_string(), + format!("{transaction}_0") + ); + assert_eq!( + TwoPcTransactionOnShard::new(transaction, 3).to_string(), + format!("{transaction}_3") + ); + } + + #[test] + fn parse_transaction_on_shard() { + let transaction = TwoPcTransaction::new(); + let parsed: TwoPcTransactionOnShard = format!("{transaction}_3") + .parse() + .expect("valid transaction on shard"); + + assert_eq!(parsed.transaction, transaction); + assert_eq!(parsed.shard, 3); + } + + #[test] + fn reject_invalid_transaction_on_shard() { + assert!("invalid".parse::().is_err()); + assert!("invalid_0".parse::().is_err()); + assert!( + "__pgdog_2pc_1_invalid" + .parse::() + .is_err() + ); } #[test] fn phase_control_statements() { + let transaction = TwoPcTransaction::new(); + assert_eq!( - phase_control("test", 1, TwoPcPhase::Phase1), - "PREPARE TRANSACTION 'test_1'" + phase_control(transaction, 1, TwoPcPhase::Phase1), + format!("PREPARE TRANSACTION '{transaction}_1'") ); assert_eq!( - phase_control("test", 1, TwoPcPhase::Phase2), - "COMMIT PREPARED 'test_1'" + phase_control(transaction, 1, TwoPcPhase::Phase2), + format!("COMMIT PREPARED '{transaction}_1'") ); assert_eq!( - phase_control("test", 1, TwoPcPhase::Rollback), - "ROLLBACK PREPARED 'test_1'" + phase_control(transaction, 1, TwoPcPhase::Rollback), + format!("ROLLBACK PREPARED '{transaction}_1'") ); } } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/test.rs b/pgdog/src/frontend/client/query_engine/two_pc/test.rs index 8d6e8fb32..0ad10f27d 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/test.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/test.rs @@ -38,9 +38,7 @@ async fn test_cleanup_transaction_phase_one() { let info = Manager::get().transaction(&transaction).unwrap(); assert_eq!(info.phase, TwoPcPhase::Phase1); - conn.two_pc(&transaction.to_string(), TwoPcPhase::Phase1) - .await - .unwrap(); + conn.two_pc(transaction, TwoPcPhase::Phase1).await.unwrap(); let two_pc = conn .execute("SELECT * FROM pg_prepared_xacts") @@ -110,9 +108,7 @@ async fn test_cleanup_transaction_phase_two() { let info = Manager::get().transaction(&transaction).unwrap(); assert_eq!(info.phase, TwoPcPhase::Phase1); - conn.two_pc(&transaction.to_string(), TwoPcPhase::Phase1) - .await - .unwrap(); + conn.two_pc(transaction, TwoPcPhase::Phase1).await.unwrap(); let txns = conn .execute("SELECT * FROM pg_prepared_xacts") diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index be09e9d00..4466fa650 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -2,6 +2,8 @@ use rand::{Rng, rng}; use serde::{Deserialize, Serialize}; use std::{fmt::Display, str::FromStr}; +use crate::util::{deployment_id, instance_id}; + #[derive(Debug, Clone, Copy, PartialEq, Hash, Eq, Serialize, Deserialize)] pub struct TwoPcTransaction(usize); @@ -13,11 +15,25 @@ impl TwoPcTransaction { // so multiple instances of PgDog don't create an identical transaction. Self(rng().random_range(0..usize::MAX)) } + + /// A prefix to identify two-phase commit transactions generated + /// by this PgDog process. + fn global_prefix() -> String { + format!( + "{PREFIX}{}{}_", + if let Some(cluster_id) = deployment_id() { + format!("{}_", cluster_id) + } else { + "".into() + }, + instance_id(), + ) + } } impl Display for TwoPcTransaction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{PREFIX}{}", self.0) + write!(f, "{}{}", Self::global_prefix(), self.0) } } @@ -25,7 +41,7 @@ impl FromStr for TwoPcTransaction { type Err = (); fn from_str(s: &str) -> Result { - let id = s.split(PREFIX).last().map(|id| id.parse()); + let id = s.rsplit("_").next().map(|id| id.parse()); if let Some(Ok(id)) = id { Ok(Self(id)) @@ -37,6 +53,8 @@ impl FromStr for TwoPcTransaction { #[cfg(test)] mod test { + use crate::test_utils::set_env_var; + use super::*; #[test] @@ -46,4 +64,24 @@ mod test { let reverse = TwoPcTransaction::from_str(transaction.to_string().as_str()).unwrap(); assert_eq!(reverse.0, transaction.0); } + + #[test] + fn test_instance_id() { + for id in [1024, 11111111, usize::MAX, usize::MIN] { + let transaction = TwoPcTransaction(id); + let instance_id = instance_id(); // It's a singleton. + assert_eq!( + format!("__pgdog_2pc_{instance_id}_{id}"), + transaction.to_string() + ); + } + } + + #[test] + fn test_deployment_id() { + let _guard = set_env_var("DEPLOYMENT_ID", "1"); + let txn = TwoPcTransaction(1678); + let instance_id = instance_id(); // It's a singleton. + assert_eq!(format!("__pgdog_2pc_1_{instance_id}_1678"), txn.to_string()); + } } diff --git a/pgdog/src/frontend/error.rs b/pgdog/src/frontend/error.rs index 7469fbf7b..f0325badc 100644 --- a/pgdog/src/frontend/error.rs +++ b/pgdog/src/frontend/error.rs @@ -46,8 +46,8 @@ pub enum Error { #[error("query timeout")] Timeout(#[from] tokio::time::error::Elapsed), - #[error("schema load timeout")] - SchemaLoad, + #[error("cluster start timeout")] + ClusterStart, #[error("join error")] Join(#[from] tokio::task::JoinError), diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index f7a995023..684c87e55 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -105,6 +105,10 @@ fn main() -> Result<(), Box> { } async fn pgdog(command: Option) -> Result<(), Box> { + // Run atexit handlers on SIGTERM (e.g. llvm-cov profile flushing). + #[cfg(unix)] + install_sigterm_handler(); + // Preload TLS. Resulting primitives // are async, so doing this after Tokio launched seems prudent. net::tls::load()?; @@ -236,6 +240,24 @@ async fn pgdog(command: Option) -> Result<(), Box std::io::Result { match workers { 0 => Builder::new_current_thread() diff --git a/pgdog/src/util.rs b/pgdog/src/util.rs index c7097b021..8ea220448 100644 --- a/pgdog/src/util.rs +++ b/pgdog/src/util.rs @@ -138,11 +138,27 @@ pub fn instance_id() -> &'static str { /// Get an externally assigned, unique, node identifier /// for this instance of PgDog. +/// +/// This assumes the NODE ID follows the following format: +/// +/// - +/// pub fn node_id() -> Result { // split always returns at least one element. instance_id().split("-").last().unwrap().parse() } +static DEPLOYMENT_ID: Lazy> = Lazy::new(|| env::var("DEPLOYMENT_ID").ok()); + +/// Get the ID of this PgDog deployment. +/// +/// This should be _globally_ unique +/// and is used to differentiate 2pc transactions. +/// +pub(crate) fn deployment_id() -> Option<&'static str> { + DEPLOYMENT_ID.as_deref() +} + static HOSTNAME: Lazy = Lazy::new(|| { let hostname = env::var("HOSTNAME").unwrap_or_default(); let host = env::var("HOST").unwrap_or_default();