diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index ddf4022c9..f5f654afb 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -60,3 +60,49 @@ jobs: with: channel: engineering webhook_url: ${{ secrets.SLACK_NOTIFICATION_WEBHOOK_URL }} + + burn-in: + name: Encrypted burn-in (PostgreSQL 17) + runs-on: blacksmith-16vcpu-ubuntu-2204 + timeout-minutes: 15 + env: + PG_VERSION: 17 + CS_ZEROKMS_HOST: https://us-east-1.aws.zerokms.cipherstashmanaged.net + CS_CTS_HOST: https://ap-southeast-2.aws.cts.cipherstashmanaged.net + RUST_BACKTRACE: "1" + + steps: + - uses: actions/checkout@v4 + - uses: ./.github/actions/setup-test + + - name: Decrypt secrets + uses: cipherstash/secrets-action@main + with: + secrets-file: .github/secrets.env.encrypted + env: + CS_CLIENT_ID: ${{ secrets.CS_VAULT_CLIENT_ID }} + CS_CLIENT_KEY: ${{ secrets.CS_VAULT_CLIENT_KEY }} + CS_CLIENT_ACCESS_KEY: ${{ secrets.CS_VAULT_CLIENT_ACCESS_KEY }} + CS_WORKSPACE_CRN: ${{ secrets.CS_VAULT_WORKSPACE_CRN }} + + - name: Start bare PostgreSQL and download EQL + run: | + mise run postgres:up --extra-args "--detach --wait" + mise run eql:download + mise run postgres:eql:teardown + + - name: Run encrypted burn-in + run: mise run test:burn-in + + - name: Upload burn-in RSS report + if: ${{ !cancelled() }} + uses: actions/upload-artifact@v4 + with: + name: proxy-burn-in-pg17 + path: target/burn-in/soak-report.json + if-no-files-found: warn + + - uses: ./.github/actions/send-slack-notification + with: + channel: engineering + webhook_url: ${{ secrets.SLACK_NOTIFICATION_WEBHOOK_URL }} diff --git a/Cargo.lock b/Cargo.lock index 269cfe0dd..d3f166ca4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -869,6 +869,18 @@ dependencies = [ "x509-parser", ] +[[package]] +name = "cipherstash-proxy-burn-in" +version = "0.1.0" +dependencies = [ + "anyhow", + "clap", + "serde", + "serde_json", + "tokio", + "tokio-postgres", +] + [[package]] name = "cipherstash-proxy-integration" version = "0.1.0" @@ -3243,7 +3255,7 @@ version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" dependencies = [ - "zerocopy 0.8.24", + "zerocopy 0.8.56", ] [[package]] @@ -6306,11 +6318,11 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.24" +version = "0.8.56" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2586fea28e186957ef732a5f8b3be2da217d65c5969d4b1e17f973ebbe876879" +checksum = "556764e583adb45a9f8d413c2a147fa7e8d821e48e12b14fd560b607998b75eb" dependencies = [ - "zerocopy-derive 0.8.24", + "zerocopy-derive 0.8.56", ] [[package]] @@ -6326,9 +6338,9 @@ dependencies = [ [[package]] name = "zerocopy-derive" -version = "0.8.24" +version = "0.8.56" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a996a8f63c5c4448cd959ac1bab0aaa3306ccfd060472f85943ee0750f0169be" +checksum = "f2ab42fc20575779bd240faa45f94a74256f755c0fa9e89f0ede20d91d0cdfc1" dependencies = [ "proc-macro2", "quote", diff --git a/mise.toml b/mise.toml index 44c6902cc..568775892 100644 --- a/mise.toml +++ b/mise.toml @@ -173,6 +173,22 @@ run = """ cargo nextest run --no-fail-fast --nocapture -p cipherstash-proxy-integration """ +[tasks."test:burn-in"] +description = "Run a bounded encrypted CRUD soak against a release Proxy" +run = """ +set -e +duration="${BURN_IN_DURATION_SECONDS:-30}" +concurrency="${BURN_IN_CONCURRENCY:-4}" + +# The burn-in owns its Proxy process. Disable the optional metrics listener so +# shared developer and CI environments cannot collide on its separate port. +CS_PROMETHEUS__ENABLED=false cargo run --locked -p cipherstash-proxy-burn-in -- \ + soak \ + --duration-seconds "${duration}" \ + --concurrency "${concurrency}" \ + --output target/burn-in/soak-report.json +""" + [tasks."test:integration:setup:tls"] description = "Setup for TLS integration tests: preflight, postgres, proxy" run = """ @@ -493,8 +509,8 @@ run = """ #!/bin/bash cd tests mise run postgres:fail_if_not_running -cat sql/schema-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- -cat ../cipherstash-encrypt-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- +cat sql/schema-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql -v ON_ERROR_STOP=1 postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- +cat ../cipherstash-encrypt-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql -v ON_ERROR_STOP=1 postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- cat ../cipherstash-encrypt.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- cat sql/schema.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- """ @@ -506,8 +522,8 @@ run = """ #!/bin/bash cd tests mise run postgres:fail_if_not_running -cat sql/schema-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- -cat ../cipherstash-encrypt-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- +cat sql/schema-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql -v ON_ERROR_STOP=1 postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- +cat ../cipherstash-encrypt-uninstall.sql | docker exec -i postgres${CONTAINER_SUFFIX} psql -v ON_ERROR_STOP=1 postgresql://${CS_DATABASE__USERNAME}:${CS_DATABASE__PASSWORD_ESCAPED_FOR_TESTS}@${CS_DATABASE__HOST}:${CS_DATABASE__PORT}/${CS_DATABASE__NAME} -f- """ [tasks."postgres:up"] diff --git a/packages/cipherstash-proxy-burn-in/Cargo.toml b/packages/cipherstash-proxy-burn-in/Cargo.toml new file mode 100644 index 000000000..efceac113 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/Cargo.toml @@ -0,0 +1,13 @@ +[package] +name = "cipherstash-proxy-burn-in" +version = "0.1.0" +edition.workspace = true +publish = false + +[dependencies] +anyhow = "1" +clap = { version = "4.5", features = ["derive", "env"] } +serde = { version = "1", features = ["derive"] } +serde_json = "1" +tokio = { workspace = true } +tokio-postgres = { version = "0.7", features = ["with-serde_json-1"] } diff --git a/packages/cipherstash-proxy-burn-in/README.md b/packages/cipherstash-proxy-burn-in/README.md new file mode 100644 index 000000000..51fd5a8ff --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/README.md @@ -0,0 +1,41 @@ +# CipherStash Proxy burn-in + +This package drives deterministic conformance checks and a timed mixed CRUD workload through a +real Proxy into PostgreSQL. The fixture schema and seed migration are adapted from pg-proto's +burn-in package so results use the same type-lab and commerce model while exercising EQL domains. + +Start the test PostgreSQL service and configure the CipherStash credentials used by Proxy in +`mise.local.toml`: + +```toml +[env] +CS_WORKSPACE_CRN = "crn:region:workspace-id" +CS_CLIENT_ACCESS_KEY = "your-access-key" +CS_DEFAULT_KEYSET_ID = "your-keyset-id" +CS_CLIENT_ID = "your-client-id" +CS_CLIENT_KEY = "your-client-key" +``` + +The commands inherit these values from the environment when they launch Proxy. +The target database must also have EQL installed. If it does not, the burn-in installs +`cipherstash-encrypt.sql` automatically; run `mise run eql:download` first or provide a different +file with `--eql-path` / `BURN_IN_EQL_PATH`. + +```bash +cargo run -p cipherstash-proxy-burn-in -- conformance +cargo run -p cipherstash-proxy-burn-in -- soak --duration-seconds 300 +``` + +`soak` always runs `cargo build --locked --release --package cipherstash-proxy` and starts that +exact release binary. It samples the Proxy process RSS once per second and writes the full series +to `target/burn-in/soak-report.json`. Use `--max-rss-growth-mib` to turn retained growth into a +hard failure, and `--concurrency` to adjust load. The workload creates public fixture tables with +EQL domain columns and verifies CRUD through Proxy, so its memory measurements include encryption +and decryption work. + +Override connection URLs with `--proxy-database-url` / `--direct-database-url` or the +`BURN_IN_PROXY_DATABASE_URL` / `BURN_IN_DIRECT_DATABASE_URL` environment variables. + +CI runs a bounded PostgreSQL 17 soak with `mise run test:burn-in` and uploads +`target/burn-in/soak-report.json`. Override its defaults locally with +`BURN_IN_DURATION_SECONDS` and `BURN_IN_CONCURRENCY`. diff --git a/packages/cipherstash-proxy-burn-in/migrations/0001_schema.sql b/packages/cipherstash-proxy-burn-in/migrations/0001_schema.sql new file mode 100644 index 000000000..dbb9381bc --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/migrations/0001_schema.sql @@ -0,0 +1,54 @@ +-- EQL is installed from cipherstash-encrypt.sql before this migration runs. +-- Keep these tables in public: Proxy loads only schemas on its search path and +-- EQL Mapper resolves tables in a single, unqualified namespace. +DROP SCHEMA IF EXISTS burnin_type_lab CASCADE; +DROP SCHEMA IF EXISTS burnin_commerce CASCADE; + +DROP TABLE IF EXISTS public.burnin_commerce_order_lines; +DROP TABLE IF EXISTS public.burnin_commerce_orders; +DROP TABLE IF EXISTS public.burnin_commerce_products; +DROP TABLE IF EXISTS public.burnin_commerce_customers; +DROP TABLE IF EXISTS public.burnin_type_lab_bulk_values; +DROP TABLE IF EXISTS public.burnin_type_lab_samples; + +CREATE TABLE public.burnin_type_lab_samples ( + id integer PRIMARY KEY, + scalar eql_v3_integer_ord NOT NULL, + nullable_text eql_v3_text, + binary_value bytea NOT NULL, + tags text[] NOT NULL, + document eql_v3_json NOT NULL, + wide_text eql_v3_text NOT NULL +); + +CREATE TABLE public.burnin_type_lab_bulk_values ( + id integer PRIMARY KEY, + nullable_text eql_v3_text, + binary_value bytea NOT NULL, + wide_text eql_v3_text NOT NULL +); + +CREATE TABLE public.burnin_commerce_customers ( + id integer PRIMARY KEY, + name eql_v3_text NOT NULL +); + +CREATE TABLE public.burnin_commerce_products ( + id integer PRIMARY KEY, + sku eql_v3_text NOT NULL, + price_cents integer NOT NULL CHECK (price_cents > 0) +); + +CREATE TABLE public.burnin_commerce_orders ( + id integer PRIMARY KEY, + customer_id integer NOT NULL REFERENCES public.burnin_commerce_customers(id), + status eql_v3_text NOT NULL +); + +CREATE TABLE public.burnin_commerce_order_lines ( + order_id integer NOT NULL REFERENCES public.burnin_commerce_orders(id), + line_number integer NOT NULL, + product_id integer NOT NULL REFERENCES public.burnin_commerce_products(id), + quantity integer NOT NULL CHECK (quantity > 0), + PRIMARY KEY (order_id, line_number) +); diff --git a/packages/cipherstash-proxy-burn-in/migrations/0002_seed.sql b/packages/cipherstash-proxy-burn-in/migrations/0002_seed.sql new file mode 100644 index 000000000..c48febd80 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/migrations/0002_seed.sql @@ -0,0 +1,5 @@ +-- This migration must run through Proxy so values assigned to EQL domains are +-- encrypted before PostgreSQL stores them. +TRUNCATE burnin_commerce_order_lines, burnin_commerce_orders, + burnin_commerce_products, burnin_commerce_customers, + burnin_type_lab_bulk_values, burnin_type_lab_samples; diff --git a/packages/cipherstash-proxy-burn-in/src/conformance.rs b/packages/cipherstash-proxy-burn-in/src/conformance.rs new file mode 100644 index 000000000..479034b75 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/src/conformance.rs @@ -0,0 +1,168 @@ +//! Deterministic protocol and encrypted-data conformance checks. +//! +//! Setup first proves representative fixture values are ciphertext at rest. +//! The checks below then read those values through Proxy and cover typed +//! decryption, transactional encrypted CRUD, rollback, SQL error recovery, and +//! concurrent connections. + +use std::path::Path; + +use anyhow::{Context, Result}; +use serde_json::Value; +use tokio::time::timeout; + +use crate::database::{self, DatabaseTarget}; + +pub async fn run( + proxy_database: &DatabaseTarget, + direct_database: &DatabaseTarget, + eql_path: &Path, +) -> Result<()> { + let _run_lock = database::acquire_run_lock(direct_database).await?; + database::ensure_eql_installed(direct_database, eql_path).await?; + timeout( + database::MIGRATION_TIMEOUT, + database::migrate(proxy_database, direct_database), + ) + .await + .context("fixture migration timed out")??; + let mut client = database::connect(proxy_database).await?; + + client + .simple_query("SELECT current_database(), current_user") + .await + .context("simple-query startup conformance")?; + + let sample = client + .query_one( + "SELECT scalar, nullable_text, binary_value, tags, document, wide_text \ + FROM burnin_type_lab_samples WHERE id = $1", + &[&1_i32], + ) + .await + .context("extended-query type conformance")?; + anyhow::ensure!(sample.get::<_, i32>(0) == 10, "scalar value was corrupted"); + anyhow::ensure!( + sample.get::<_, Option>(1).is_none(), + "NULL was corrupted" + ); + anyhow::ensure!( + sample.get::<_, Vec>(2) == [0, 1, 2, 255], + "bytea was corrupted" + ); + anyhow::ensure!( + sample.get::<_, Vec>(3) == ["alpha", "one"], + "array was corrupted" + ); + anyhow::ensure!( + sample.get::<_, Value>(4)["kind"] == "alpha", + "jsonb was corrupted" + ); + anyhow::ensure!( + sample.get::<_, String>(5) == "wide-alpha-".repeat(40), + "wide text was corrupted" + ); + + let transaction = client + .transaction() + .await + .context("starting CRUD transaction")?; + let fixture_id = 900_001_i32; + transaction + .execute( + "INSERT INTO burnin_commerce_customers (id, name) VALUES ($1, $2)", + &[&fixture_id, &"conformance-customer"], + ) + .await?; + transaction + .execute( + "INSERT INTO burnin_commerce_products (id, sku, price_cents) VALUES ($1, $2, $3)", + &[&fixture_id, &"CONF-900001", &2_499_i32], + ) + .await?; + transaction + .execute( + "INSERT INTO burnin_commerce_orders (id, customer_id, status) VALUES ($1, $2, $3)", + &[&fixture_id, &fixture_id, &"open"], + ) + .await?; + transaction.execute( + "INSERT INTO burnin_commerce_order_lines (order_id, line_number, product_id, quantity) VALUES ($1, 1, $2, 2)", + &[&fixture_id, &fixture_id], + ).await?; + let total: Option = transaction + .query_one( + "SELECT sum(p.price_cents * l.quantity)::bigint \ + FROM burnin_commerce_orders o \ + JOIN burnin_commerce_order_lines l ON l.order_id = o.id \ + JOIN burnin_commerce_products p ON p.id = l.product_id \ + WHERE o.id = $1", + &[&fixture_id], + ) + .await? + .try_get(0) + .context("decoding joined CRUD total")?; + anyhow::ensure!(total == Some(4_998), "joined CRUD result was corrupted"); + transaction + .execute( + "UPDATE burnin_commerce_orders SET status = 'paid' WHERE id = $1", + &[&fixture_id], + ) + .await?; + transaction + .rollback() + .await + .context("rolling back CRUD transaction")?; + let rolled_back: i64 = client + .query_one( + "SELECT count(*) FROM burnin_commerce_orders WHERE id = $1", + &[&fixture_id], + ) + .await? + .get(0); + anyhow::ensure!(rolled_back == 0, "transaction rollback leaked a row"); + + let error = client + .execute( + "INSERT INTO burnin_commerce_products (id, sku, price_cents) VALUES ($1, $2, $3)", + &[&900_002_i32, &"INVALID-PRICE", &0_i32], + ) + .await + .expect_err("check constraint should reject a zero price"); + anyhow::ensure!( + error.code().is_some_and(|code| code.code() == "23514"), + "unexpected SQLSTATE: {error}" + ); + let recovered: i32 = client + .query_one("SELECT $1::integer", &[&42_i32]) + .await? + .get(0); + anyhow::ensure!( + recovered == 42, + "connection did not recover after SQL error" + ); + + let mut tasks = Vec::new(); + for worker in 0..8_i32 { + let target = proxy_database.clone(); + tasks.push(tokio::spawn(async move { + let client = database::connect(&target).await?; + for iteration in 0..25_i32 { + let value: i32 = client + .query_one("SELECT $1::integer + $2::integer", &[&worker, &iteration]) + .await? + .get(0); + anyhow::ensure!(value == worker + iteration, "concurrent result mismatch"); + } + Result::<()>::Ok(()) + })); + } + for task in tasks { + task.await.context("conformance worker panicked")??; + } + + println!( + "conformance passed: simple, extended, types, CRUD, rollback, error recovery, concurrency" + ); + Ok(()) +} diff --git a/packages/cipherstash-proxy-burn-in/src/database.rs b/packages/cipherstash-proxy-burn-in/src/database.rs new file mode 100644 index 000000000..e68c723e7 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/src/database.rs @@ -0,0 +1,328 @@ +//! Database lifecycle for encrypted burn-in fixtures. +//! +//! EQL itself is installed directly because Proxy cannot map statements until +//! its domains exist. Fixture DDL and seed writes then go through Proxy so DDL +//! triggers schema/encrypt-config reloads and seed values are encrypted. DDL +//! and seed use different Proxy connections because each connection snapshots +//! those configurations when it opens. + +use std::{fmt, path::Path, str::FromStr}; + +use anyhow::{Context, Result}; +use serde_json::{json, Value}; +use tokio_postgres::{config::Host, types::Type, Client, NoTls}; + +use crate::{SCHEMA_MIGRATION, SEED_MIGRATION}; + +const RUN_LOCK_ID: i64 = 0x4353_4255_524e_494e; +pub const MIGRATION_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); +const SEED_SAMPLE_SQL: &str = "INSERT INTO burnin_type_lab_samples \ + (id, scalar, nullable_text, binary_value, tags, document, wide_text) \ + VALUES ($1, $2, $3, $4, $5, $6, $7)"; + +fn seed_parameter_types() -> [Type; 7] { + [ + Type::INT4, + Type::INT4, + Type::TEXT, + Type::BYTEA, + Type::TEXT_ARRAY, + Type::JSONB, + Type::TEXT, + ] +} + +#[derive(Clone)] +pub struct DatabaseTarget { + config: tokio_postgres::Config, + identity: String, +} + +impl DatabaseTarget { + pub fn hostname(&self) -> Result<&str> { + match self.config.get_hosts() { + [Host::Tcp(host)] => Ok(host), + _ => anyhow::bail!("burn-in requires exactly one TCP database host"), + } + } + + pub fn port(&self) -> Result { + match self.config.get_ports() { + [] => Ok(5432), + [port] => Ok(*port), + _ => anyhow::bail!("burn-in requires exactly one database port"), + } + } + + pub fn configure_proxy_upstream(&self, command: &mut tokio::process::Command) -> Result<()> { + command + .env("CS_DATABASE__HOST", self.hostname()?) + .env("CS_DATABASE__PORT", self.port()?.to_string()) + .env( + "CS_DATABASE__NAME", + self.config.get_dbname().unwrap_or("postgres"), + ) + .env( + "CS_DATABASE__USERNAME", + self.config.get_user().unwrap_or("postgres"), + ); + if let Some(password) = self.config.get_password() { + command.env( + "CS_DATABASE__PASSWORD", + std::str::from_utf8(password).context("database password is not UTF-8")?, + ); + } else { + command.env_remove("CS_DATABASE__PASSWORD"); + } + Ok(()) + } +} + +impl FromStr for DatabaseTarget { + type Err = String; + + fn from_str(value: &str) -> std::result::Result { + let config = value + .parse::() + .map_err(|_| "invalid PostgreSQL connection configuration".to_string())?; + let host = match config.get_hosts() { + [Host::Tcp(host)] => host.clone(), + [Host::Unix(path)] => path.to_str().unwrap_or("unix-socket").to_string(), + _ => "multiple-hosts".to_string(), + }; + let port = config.get_ports().first().copied().unwrap_or(5432); + let user = config.get_user().unwrap_or("postgres").to_string(); + let database = config.get_dbname().unwrap_or("postgres").to_string(); + Ok(Self { + config, + identity: format!("postgresql://{user}@{host}:{port}/{database}"), + }) + } +} + +impl fmt::Debug for DatabaseTarget { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_tuple("DatabaseTarget") + .field(&self.identity) + .finish() + } +} + +impl fmt::Display for DatabaseTarget { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(&self.identity) + } +} + +pub async fn connect(target: &DatabaseTarget) -> Result { + let (client, connection) = target + .config + .connect(NoTls) + .await + .with_context(|| format!("connecting to {target}"))?; + tokio::spawn(async move { + if let Err(error) = connection.await { + eprintln!("database connection failed: {error}"); + } + }); + Ok(client) +} + +pub async fn acquire_run_lock(target: &DatabaseTarget) -> Result { + let client = connect(target).await?; + let acquired: bool = client + .query_one("SELECT pg_try_advisory_lock($1)", &[&RUN_LOCK_ID]) + .await + .context("acquiring the burn-in database lock")? + .get(0); + anyhow::ensure!( + acquired, + "another burn-in or conformance run already owns the database fixtures" + ); + Ok(client) +} + +/// Ensure the EQL domains exist before fixture DDL is sent through Proxy. +pub async fn ensure_eql_installed(direct_database: &DatabaseTarget, eql_path: &Path) -> Result<()> { + let client = connect(direct_database).await?; + let installed: bool = client + .query_one( + "SELECT EXISTS (\ + SELECT 1 FROM pg_type t \ + JOIN pg_namespace n ON n.oid = t.typnamespace \ + WHERE n.nspname = 'public' AND t.typname = 'eql_v3_text'\ + )", + &[], + ) + .await + .context("checking whether EQL is installed")? + .get(0); + if installed { + return Ok(()); + } + + let eql = tokio::fs::read_to_string(eql_path).await.with_context(|| { + format!( + "EQL is not installed and its migration could not be read from {}; run `mise run eql:download` or pass --eql-path", + eql_path.display() + ) + })?; + client + .batch_execute(&eql) + .await + .with_context(|| format!("installing EQL from {}", eql_path.display()))?; + Ok(()) +} + +/// Create and seed fixtures through Proxy so DDL reloads its schema and every +/// value assigned to an EQL domain traverses the encryption path. +pub async fn migrate( + proxy_database: &DatabaseTarget, + direct_database: &DatabaseTarget, +) -> Result<()> { + // A connection snapshots Proxy's schema and encrypt config when it opens. + // Apply DDL on one connection, let Proxy reload, then open a fresh + // connection whose snapshot includes the new encrypted fixture columns. + let ddl_client = connect(proxy_database).await?; + ddl_client + .batch_execute(SCHEMA_MIGRATION) + .await + .context("applying burn-in schema migration through Proxy")?; + drop(ddl_client); + + let client = connect(proxy_database).await?; + client + .batch_execute(SEED_MIGRATION) + .await + .context("clearing burn-in fixtures through Proxy")?; + + seed_sample( + &client, + 1, + 10, + None, + vec![0, 1, 2, 255], + vec!["alpha".into(), "one".into()], + json!({"kind": "alpha", "enabled": true}), + "wide-alpha-".repeat(40), + ) + .await?; + seed_sample( + &client, + 2, + 20, + Some("second".into()), + vec![0x10, 0x20, 0x30, 0x40], + vec![], + json!({"kind": "beta", "count": 2}), + "wide-beta-".repeat(40), + ) + .await?; + seed_sample( + &client, + 3, + 30, + None, + vec![0xde, 0xad, 0xbe, 0xef], + vec!["nullable".into()], + json!({"kind": "gamma", "values": [1, 2, 3]}), + "wide-gamma-".repeat(40), + ) + .await?; + seed_sample( + &client, + 4, + 40, + Some("fourth".into()), + vec![0xca, 0xfe, 0xba, 0xbe], + vec!["delta".into(), "four".into()], + json!({"kind": "delta", "value": null}), + "wide-delta-".repeat(40), + ) + .await?; + assert_seed_is_encrypted(direct_database).await?; + Ok(()) +} + +async fn assert_seed_is_encrypted(direct_database: &DatabaseTarget) -> Result<()> { + // Query around Proxy and inspect the JSON-backed domains themselves. A + // successful round trip through Proxy is insufficient proof: an unmappable + // statement can be passed through and appear correct while storing plaintext. + let client = connect(direct_database).await?; + let row = client + .query_one( + "SELECT scalar::jsonb, document::jsonb, wide_text::jsonb \ + FROM burnin_type_lab_samples WHERE id = 1", + &[], + ) + .await + .context("reading seeded ciphertext directly from PostgreSQL")?; + + for (index, column) in ["scalar", "document", "wide_text"].into_iter().enumerate() { + let ciphertext: Value = row.get(index); + anyhow::ensure!( + ciphertext.get("c").is_some() && ciphertext.get("v").is_some(), + "{column} was stored as plaintext instead of EQL ciphertext: {ciphertext}" + ); + } + Ok(()) +} + +#[allow(clippy::too_many_arguments)] +async fn seed_sample( + client: &Client, + id: i32, + scalar: i32, + nullable_text: Option, + binary_value: Vec, + tags: Vec, + document: Value, + wide_text: String, +) -> Result<()> { + let values: [&(dyn tokio_postgres::types::ToSql + Sync); 7] = [ + &id, + &scalar, + &nullable_text, + &binary_value, + &tags, + &document, + &wide_text, + ]; + let parameters: Vec<_> = values.into_iter().zip(seed_parameter_types()).collect(); + client + .query_typed(SEED_SAMPLE_SQL, ¶meters) + .await + .with_context(|| format!("seeding burn-in sample {id} through Proxy"))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn database_target_redacts_password() { + let target: DatabaseTarget = "postgresql://alice:super-secret@localhost:5544/app" + .parse() + .unwrap(); + assert_eq!(target.to_string(), "postgresql://alice@localhost:5544/app"); + assert!(!format!("{target:?}").contains("super-secret")); + } + + #[test] + fn seed_parameters_use_native_types_before_proxy_encryption() { + assert_eq!( + seed_parameter_types(), + [ + Type::INT4, + Type::INT4, + Type::TEXT, + Type::BYTEA, + Type::TEXT_ARRAY, + Type::JSONB, + Type::TEXT, + ] + ); + } +} diff --git a/packages/cipherstash-proxy-burn-in/src/lib.rs b/packages/cipherstash-proxy-burn-in/src/lib.rs new file mode 100644 index 000000000..f72690f67 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/src/lib.rs @@ -0,0 +1,56 @@ +//! End-to-end correctness and soak workloads for CipherStash Proxy. +//! +//! The fixtures deliberately use EQL domains on uniquely named tables in +//! `public`, and workload SQL deliberately leaves those table names +//! unqualified. Proxy only loads schemas on its search path and EQL Mapper +//! resolves a flat table namespace; changing either invariant can silently +//! turn this into a passthrough workload that never exercises encryption. + +pub mod conformance; +pub mod database; +pub mod resource; +pub mod soak; + +pub const SCHEMA_MIGRATION: &str = include_str!("../migrations/0001_schema.sql"); +pub const SEED_MIGRATION: &str = include_str!("../migrations/0002_seed.sql"); + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn fixtures_are_public_and_include_encrypted_domains() { + assert!(SCHEMA_MIGRATION.contains("CREATE TABLE public.burnin_")); + assert!(SCHEMA_MIGRATION.contains("eql_v3_integer_ord")); + assert!(SCHEMA_MIGRATION.contains("eql_v3_text")); + assert!(SCHEMA_MIGRATION.contains("eql_v3_json")); + } + + #[test] + fn workload_queries_do_not_use_schema_qualified_fixture_names() { + let workload = concat!(include_str!("conformance.rs"), include_str!("soak.rs")); + assert!(!workload.contains("burnin_type_lab.")); + assert!(!workload.contains("burnin_commerce.")); + } + + #[test] + fn ci_starts_proxy_against_a_database_without_encrypted_columns() { + let workflow = include_str!("../../../.github/workflows/test.yml"); + let burn_in_job = workflow.split(" burn-in:").nth(1).expect("burn-in CI job"); + + assert!(burn_in_job.contains("mise run eql:download")); + assert!(burn_in_job.contains("mise run postgres:eql:teardown")); + assert!(!burn_in_job.contains("mise run postgres:setup")); + } + + #[test] + fn eql_teardown_stops_on_sql_errors() { + let tasks = include_str!("../../../mise.toml"); + let teardown = tasks + .split("[tasks.\"postgres:eql:teardown\"]") + .nth(1) + .expect("EQL teardown task"); + + assert!(teardown.contains("ON_ERROR_STOP=1")); + } +} diff --git a/packages/cipherstash-proxy-burn-in/src/main.rs b/packages/cipherstash-proxy-burn-in/src/main.rs new file mode 100644 index 000000000..733d8db28 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/src/main.rs @@ -0,0 +1,106 @@ +//! CLI entry point for deterministic conformance and timed burn-in runs. +//! +//! Connection and EQL paths are explicit options with environment-variable +//! equivalents so the same binary works in local mise environments and CI. + +use std::{path::PathBuf, time::Duration}; + +use anyhow::Result; +use cipherstash_proxy_burn_in::database::DatabaseTarget; +use clap::{Args, Parser, Subcommand}; + +#[derive(Debug, Parser)] +#[command( + name = "cipherstash-proxy-burn-in", + about = "Conformance and release-mode soak testing for CipherStash Proxy" +)] +struct Cli { + #[command(subcommand)] + command: Command, +} + +#[derive(Debug, Subcommand)] +enum Command { + /// Run deterministic correctness and PostgreSQL-protocol scenarios. + Conformance(DatabaseArgs), + /// Build and start the proxy in release mode, then run a timed stress workload. + Soak(SoakArgs), +} + +#[derive(Debug, Args)] +struct DatabaseArgs { + /// Connection URL through CipherStash Proxy. + #[arg( + long, + env = "BURN_IN_PROXY_DATABASE_URL", + hide_env_values = true, + default_value = "postgresql://cipherstash:p%40ssword@localhost:6432/cipherstash" + )] + proxy_database_url: DatabaseTarget, + /// Direct PostgreSQL URL used only to install and seed the fixture schema. + #[arg( + long, + env = "BURN_IN_DIRECT_DATABASE_URL", + hide_env_values = true, + default_value = "postgresql://cipherstash:p%40ssword@localhost:5532/cipherstash" + )] + direct_database_url: DatabaseTarget, + /// EQL installation SQL used when the target database has no EQL domains. + #[arg( + long, + env = "BURN_IN_EQL_PATH", + default_value = "cipherstash-encrypt.sql" + )] + eql_path: PathBuf, +} + +#[derive(Debug, Args)] +struct SoakArgs { + #[command(flatten)] + database: DatabaseArgs, + /// Wall-clock duration of the stress workload. + #[arg(long)] + duration_seconds: u64, + /// Number of concurrent long-lived database connections. + #[arg(long, default_value_t = 8)] + concurrency: usize, + /// JSON report containing operation counts and one-second RSS samples. + #[arg(long, default_value = "target/burn-in/soak-report.json")] + output: PathBuf, + /// Optional hard gate for end-to-end proxy RSS growth, in MiB. + #[arg(long)] + max_rss_growth_mib: Option, +} + +#[tokio::main] +async fn main() -> Result<()> { + match Cli::parse().command { + Command::Conformance(args) => { + cipherstash_proxy_burn_in::conformance::run( + &args.proxy_database_url, + &args.direct_database_url, + &args.eql_path, + ) + .await + } + Command::Soak(args) => { + let max_rss_growth_bytes = args + .max_rss_growth_mib + .map(|mib| { + mib.checked_mul(1_048_576) + .ok_or_else(|| anyhow::anyhow!("--max-rss-growth-mib is too large")) + }) + .transpose()?; + cipherstash_proxy_burn_in::soak::run(cipherstash_proxy_burn_in::soak::Config { + duration: Duration::from_secs(args.duration_seconds), + concurrency: args.concurrency, + proxy_database_url: args.database.proxy_database_url, + direct_database_url: args.database.direct_database_url, + eql_path: args.database.eql_path, + output: args.output, + max_rss_growth_bytes, + }) + .await + } + } +} diff --git a/packages/cipherstash-proxy-burn-in/src/resource.rs b/packages/cipherstash-proxy-burn-in/src/resource.rs new file mode 100644 index 000000000..8af895767 --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/src/resource.rs @@ -0,0 +1,95 @@ +//! Cross-platform resident-memory sampling for the spawned Proxy process. +//! +//! Linux reads `/proc` for CI while other platforms use `ps`, keeping report +//! semantics identical for local and automated soak runs. + +use std::time::Instant; + +#[cfg(target_os = "linux")] +use std::fs; +#[cfg(not(target_os = "linux"))] +use std::process::Command; + +use anyhow::{Context, Result}; +use serde::Serialize; + +#[derive(Debug, Clone, Serialize)] +pub struct MemorySample { + pub elapsed_millis: u128, + pub rss_bytes: u64, +} + +pub fn sample(pid: u32, started_at: Instant) -> Result { + Ok(MemorySample { + elapsed_millis: started_at.elapsed().as_millis(), + rss_bytes: resident_bytes(pid)?, + }) +} + +#[cfg(target_os = "linux")] +fn resident_bytes(pid: u32) -> Result { + let status = fs::read_to_string(format!("/proc/{pid}/status")) + .with_context(|| format!("reading memory for proxy PID {pid}"))?; + let value = status + .lines() + .find_map(|line| line.strip_prefix("VmRSS:")) + .and_then(|line| line.split_whitespace().next()) + .context("VmRSS was absent from proc status")?; + Ok(value.parse::()? * 1024) +} + +#[cfg(not(target_os = "linux"))] +fn resident_bytes(pid: u32) -> Result { + let output = Command::new("ps") + .args(["-o", "rss=", "-p", &pid.to_string()]) + .output() + .with_context(|| format!("running ps for proxy PID {pid}"))?; + anyhow::ensure!( + output.status.success(), + "ps could not inspect proxy PID {pid}" + ); + let rss_kib = String::from_utf8(output.stdout)?.trim().parse::()?; + Ok(rss_kib * 1024) +} + +pub fn growth_bytes(samples: &[MemorySample]) -> u64 { + let Some(first) = samples.first() else { + return 0; + }; + samples + .last() + .map_or(0, |last| last.rss_bytes.saturating_sub(first.rss_bytes)) +} + +pub fn peak_bytes(samples: &[MemorySample]) -> u64 { + samples + .iter() + .map(|sample| sample.rss_bytes) + .max() + .unwrap_or(0) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn calculates_growth_and_peak() { + let samples = [ + MemorySample { + elapsed_millis: 0, + rss_bytes: 10, + }, + MemorySample { + elapsed_millis: 1, + rss_bytes: 25, + }, + MemorySample { + elapsed_millis: 2, + rss_bytes: 20, + }, + ]; + assert_eq!(growth_bytes(&samples), 10); + assert_eq!(peak_bytes(&samples), 25); + } +} diff --git a/packages/cipherstash-proxy-burn-in/src/soak.rs b/packages/cipherstash-proxy-burn-in/src/soak.rs new file mode 100644 index 000000000..44517c50f --- /dev/null +++ b/packages/cipherstash-proxy-burn-in/src/soak.rs @@ -0,0 +1,545 @@ +//! Timed encrypted CRUD workload with release-Proxy RSS sampling. +//! +//! This module builds and owns the exact Proxy process being measured. Every +//! CRUD cycle writes and reads EQL-domain columns; fixture setup also verifies +//! ciphertext directly in PostgreSQL before timing begins. Memory growth +//! therefore includes the encryption/decryption path rather than passthrough +//! SQL alone. + +use std::{ + net::TcpListener, + path::{Path, PathBuf}, + process::Stdio, + sync::{ + atomic::{AtomicU64, Ordering}, + Arc, + }, + time::{Duration, Instant}, +}; + +use anyhow::{Context, Result}; +use serde::{Deserialize, Serialize}; +use tokio::{ + process::{Child, Command}, + task::JoinSet, + time::timeout, +}; + +use crate::{ + database::{self, DatabaseTarget}, + resource::{self, MemorySample}, +}; + +#[derive(Debug)] +pub struct Config { + pub duration: Duration, + pub concurrency: usize, + pub proxy_database_url: DatabaseTarget, + pub direct_database_url: DatabaseTarget, + pub eql_path: PathBuf, + pub output: PathBuf, + pub max_rss_growth_bytes: Option, +} + +#[derive(Debug, Serialize)] +struct Report { + status: &'static str, + terminal_error: Option, + requested_duration_seconds: u64, + actual_elapsed_millis: u128, + concurrency: usize, + started_at_unix_seconds: u64, + artifact: String, + source_commit: String, + database: String, + operations: u64, + errors: u64, + proxy_pid: u32, + initial_rss_bytes: u64, + final_rss_bytes: u64, + peak_rss_bytes: u64, + rss_growth_bytes: u64, + memory_samples: Vec, +} + +const READY_TIMEOUT: Duration = Duration::from_secs(30); +const OPERATION_TIMEOUT: Duration = Duration::from_secs(10); +const WORKER_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(15); + +pub async fn run(config: Config) -> Result<()> { + anyhow::ensure!( + !config.duration.is_zero(), + "--duration-seconds must be positive" + ); + anyhow::ensure!(config.concurrency > 0, "--concurrency must be positive"); + preflight_output(&config.output).await?; + let _run_lock = database::acquire_run_lock(&config.direct_database_url).await?; + database::ensure_eql_installed(&config.direct_database_url, &config.eql_path).await?; + + let artifact = build_release_proxy().await?; + preflight_listener(&config.proxy_database_url)?; + let mut proxy = spawn_release_proxy( + &artifact, + &config.proxy_database_url, + &config.direct_database_url, + )?; + let proxy_pid = proxy.id().context("release proxy did not expose a PID")?; + let started_at = Instant::now(); + let mut report = Report { + status: "failed", + terminal_error: None, + requested_duration_seconds: config.duration.as_secs(), + actual_elapsed_millis: 0, + concurrency: config.concurrency, + started_at_unix_seconds: std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_secs(), + artifact: artifact.display().to_string(), + source_commit: source_commit(), + database: config.direct_database_url.to_string(), + operations: 0, + errors: 0, + proxy_pid, + initial_rss_bytes: 0, + final_rss_bytes: 0, + peak_rss_bytes: 0, + rss_growth_bytes: 0, + memory_samples: Vec::new(), + }; + let mut result = run_with_proxy(&config, &mut proxy, &mut report).await; + let _ = proxy.kill().await; + let _ = proxy.wait().await; + report.actual_elapsed_millis = started_at.elapsed().as_millis(); + refresh_rss_summary(&mut report); + + if result.is_ok() { + result = validate_report(&report, config.max_rss_growth_bytes); + } + match &result { + Ok(()) => report.status = "passed", + Err(error) => report.terminal_error = Some(format!("{error:#}")), + } + write_report_atomic(&config.output, &report).await?; + result?; + println!( + "soak passed: {} CRUD cycles, peak RSS {} MiB, RSS growth {} MiB; report: {}", + report.operations, + report.peak_rss_bytes / 1_048_576, + report.rss_growth_bytes / 1_048_576, + config.output.display() + ); + Ok(()) +} + +async fn run_with_proxy(config: &Config, proxy: &mut Child, report: &mut Report) -> Result<()> { + wait_until_ready(&config.proxy_database_url, proxy).await?; + ensure_child_running(proxy, "fixture migration")?; + timeout( + database::MIGRATION_TIMEOUT, + database::migrate(&config.proxy_database_url, &config.direct_database_url), + ) + .await + .context("fixture migration timed out")??; + ensure_child_running(proxy, "workload warm-up")?; + tokio::time::sleep(Duration::from_secs(1)).await; + ensure_child_running(proxy, "workload start")?; + let proxy_pid = proxy.id().context("release proxy did not expose a PID")?; + let started_at = Instant::now(); + let deadline = tokio::time::Instant::now() + config.duration; + let operations = Arc::new(AtomicU64::new(0)); + let errors = Arc::new(AtomicU64::new(0)); + let ids = Arc::new(AtomicU64::new(1_000_000)); + let mut workers = JoinSet::new(); + + // Workers keep connections open so the soak stresses repeated statement + // mapping and cipher use rather than connection establishment throughput. + for _ in 0..config.concurrency { + let url = config.proxy_database_url.clone(); + let operations = Arc::clone(&operations); + let errors = Arc::clone(&errors); + let ids = Arc::clone(&ids); + workers.spawn(async move { + let mut client = timeout(OPERATION_TIMEOUT, database::connect(&url)) + .await + .context("worker database connection timed out")??; + while tokio::time::Instant::now() < deadline { + let id = i32::try_from(ids.fetch_add(1, Ordering::Relaxed))?; + match timeout(OPERATION_TIMEOUT, crud_cycle(&mut client, id)).await { + Err(_) => { + errors.fetch_add(1, Ordering::Relaxed); + return Err(anyhow::anyhow!("CRUD cycle {id} timed out")); + } + Ok(Ok(())) => { + operations.fetch_add(1, Ordering::Relaxed); + } + Ok(Err(error)) => { + errors.fetch_add(1, Ordering::Relaxed); + return Err(error.context(format!("CRUD cycle {id}"))); + } + } + } + Result::<()>::Ok(()) + }); + } + + let mut ticker = tokio::time::interval_at( + tokio::time::Instant::now() + Duration::from_secs(1), + Duration::from_secs(1), + ); + while tokio::time::Instant::now() < deadline { + tokio::select! { + _ = ticker.tick() => {} + signal = tokio::signal::ctrl_c() => { + signal.context("installing interrupt handler")?; + anyhow::bail!("burn-in interrupted"); + } + } + ensure_child_running(proxy, "RSS sampling")?; + let sample = resource::sample(proxy_pid, started_at)?; + anyhow::ensure!(sample.rss_bytes > 0, "Proxy RSS sample was zero"); + report.memory_samples.push(sample); + refresh_counters(report, &operations, &errors); + } + ensure_child_running(proxy, "worker shutdown")?; + let worker_result = timeout(WORKER_SHUTDOWN_TIMEOUT, async { + while let Some(result) = workers.join_next().await { + result.context("soak worker panicked")??; + } + Result::<()>::Ok(()) + }) + .await; + if worker_result.is_err() { + workers.shutdown().await; + } + // Snapshot after aborting timed-out workers and before propagating the + // timeout, so failure reports do not retain the previous ticker's counts. + refresh_counters(report, &operations, &errors); + let worker_result = worker_result.context("workers did not stop within 15 seconds")?; + let final_sample = resource::sample(proxy_pid, started_at)?; + anyhow::ensure!( + final_sample.rss_bytes > 0, + "final Proxy RSS sample was zero" + ); + report.memory_samples.push(final_sample); + + refresh_counters(report, &operations, &errors); + worker_result?; + Ok(()) +} + +async fn crud_cycle(client: &mut tokio_postgres::Client, id: i32) -> Result<()> { + let transaction = client.transaction().await?; + let name = format!("soak-customer-{id}"); + let sku = format!("SOAK-{id}"); + transaction + .execute( + "INSERT INTO burnin_commerce_customers (id, name) VALUES ($1, $2)", + &[&id, &name], + ) + .await?; + transaction + .execute( + "INSERT INTO burnin_commerce_products (id, sku, price_cents) VALUES ($1, $2, $3)", + &[&id, &sku, &(100 + id % 10_000)], + ) + .await?; + transaction + .execute( + "INSERT INTO burnin_commerce_orders (id, customer_id, status) VALUES ($1, $1, 'open')", + &[&id], + ) + .await?; + transaction.execute( + "INSERT INTO burnin_commerce_order_lines (order_id, line_number, product_id, quantity) VALUES ($1, 1, $1, 2)", &[&id] + ).await?; + let row = transaction + .query_one( + "SELECT c.name, p.sku, p.price_cents * l.quantity \ + FROM burnin_commerce_orders o \ + JOIN burnin_commerce_customers c ON c.id = o.customer_id \ + JOIN burnin_commerce_order_lines l ON l.order_id = o.id \ + JOIN burnin_commerce_products p ON p.id = l.product_id WHERE o.id = $1", + &[&id], + ) + .await?; + anyhow::ensure!( + row.get::<_, String>(0) == name && row.get::<_, String>(1) == sku, + "read-after-write mismatch" + ); + transaction + .execute( + "UPDATE burnin_commerce_orders SET status = 'fulfilled' WHERE id = $1", + &[&id], + ) + .await?; + transaction + .execute( + "DELETE FROM burnin_commerce_order_lines WHERE order_id = $1", + &[&id], + ) + .await?; + transaction + .execute("DELETE FROM burnin_commerce_orders WHERE id = $1", &[&id]) + .await?; + transaction + .execute("DELETE FROM burnin_commerce_products WHERE id = $1", &[&id]) + .await?; + transaction + .execute( + "DELETE FROM burnin_commerce_customers WHERE id = $1", + &[&id], + ) + .await?; + transaction.commit().await?; + Ok(()) +} + +#[derive(Deserialize)] +struct CargoArtifact { + reason: String, + target: Option, + executable: Option, +} + +#[derive(Deserialize)] +struct CargoTarget { + name: String, +} + +async fn build_release_proxy() -> Result { + let output = Command::new("cargo") + .args([ + "build", + "--locked", + "--release", + "--package", + "cipherstash-proxy", + "--message-format=json-render-diagnostics", + ]) + .current_dir(workspace_root()) + .output() + .await + .context("building release proxy")?; + anyhow::ensure!( + output.status.success(), + "release proxy build failed: {}", + String::from_utf8_lossy(&output.stderr).trim() + ); + find_proxy_artifact(&output.stdout) +} + +fn find_proxy_artifact(messages: &[u8]) -> Result { + let mut artifact = None; + for line in messages.split(|byte| *byte == b'\n') { + let Ok(message) = serde_json::from_slice::(line) else { + continue; + }; + if message.reason == "compiler-artifact" + && message + .target + .is_some_and(|target| target.name == "cipherstash-proxy") + && message.executable.is_some() + { + artifact = message.executable; + } + } + artifact.context("Cargo did not report the release Proxy executable") +} + +fn spawn_release_proxy( + binary: &Path, + proxy_database: &DatabaseTarget, + direct_database: &DatabaseTarget, +) -> Result { + anyhow::ensure!( + binary.is_file(), + "release proxy binary is missing at {}", + binary.display() + ); + let host = proxy_bind_host(proxy_database)?; + let mut command = Command::new(binary); + direct_database.configure_proxy_upstream(&mut command)?; + command + .env("CS_SERVER__HOST", host) + .env("CS_SERVER__PORT", proxy_database.port()?.to_string()) + .kill_on_drop(true) + .current_dir(workspace_root()) + .stdin(Stdio::null()) + .stdout(Stdio::inherit()) + .stderr(Stdio::inherit()); + command.spawn().context("starting release proxy") +} + +fn preflight_listener(proxy_database: &DatabaseTarget) -> Result<()> { + let listener = TcpListener::bind((proxy_bind_host(proxy_database)?, proxy_database.port()?)) + .context("Proxy listen address is already in use")?; + drop(listener); + Ok(()) +} + +fn proxy_bind_host(proxy_database: &DatabaseTarget) -> Result<&'static str> { + match proxy_database.hostname()? { + "localhost" | "127.0.0.1" => Ok("127.0.0.1"), + "::1" => Ok("::1"), + _ => anyhow::bail!("spawned Proxy must use a loopback listener"), + } +} + +fn ensure_child_running(child: &mut Child, phase: &str) -> Result<()> { + if let Some(status) = child.try_wait().context("checking release Proxy status")? { + anyhow::bail!("release Proxy exited during {phase} with {status}"); + } + Ok(()) +} + +async fn wait_until_ready(target: &DatabaseTarget, child: &mut Child) -> Result<()> { + let deadline = tokio::time::Instant::now() + READY_TIMEOUT; + loop { + ensure_child_running(child, "startup")?; + let ready = timeout(Duration::from_secs(2), async { + let client = database::connect(target).await?; + client.simple_query("SELECT 1").await?; + Result::<()>::Ok(()) + }) + .await; + if matches!(ready, Ok(Ok(()))) { + return Ok(()); + } + if tokio::time::Instant::now() >= deadline { + anyhow::bail!("release Proxy did not become ready within 30 seconds"); + } + tokio::time::sleep(Duration::from_millis(250)).await; + } +} + +fn validate_report(report: &Report, max_rss_growth_bytes: Option) -> Result<()> { + anyhow::ensure!(report.operations > 0, "soak completed zero CRUD cycles"); + anyhow::ensure!(report.errors == 0, "soak observed {} errors", report.errors); + anyhow::ensure!( + !report.memory_samples.is_empty() && report.final_rss_bytes > 0, + "soak did not capture live Proxy RSS" + ); + if let Some(limit) = max_rss_growth_bytes { + anyhow::ensure!( + report.rss_growth_bytes <= limit, + "proxy RSS grew by {} bytes, above the {} byte limit", + report.rss_growth_bytes, + limit + ); + } + Ok(()) +} + +/// Recomputes summary fields even when the workload terminates early, so a +/// failed report retains all useful memory evidence collected before failure. +fn refresh_rss_summary(report: &mut Report) { + report.initial_rss_bytes = report + .memory_samples + .first() + .map_or(0, |sample| sample.rss_bytes); + report.final_rss_bytes = report + .memory_samples + .last() + .map_or(0, |sample| sample.rss_bytes); + report.peak_rss_bytes = resource::peak_bytes(&report.memory_samples); + report.rss_growth_bytes = resource::growth_bytes(&report.memory_samples); +} + +fn refresh_counters(report: &mut Report, operations: &AtomicU64, errors: &AtomicU64) { + report.operations = operations.load(Ordering::Relaxed); + report.errors = errors.load(Ordering::Relaxed); +} + +async fn preflight_output(path: &Path) -> Result<()> { + if let Some(parent) = path.parent().filter(|path| !path.as_os_str().is_empty()) { + tokio::fs::create_dir_all(parent).await?; + } + let temporary = temporary_report_path(path); + tokio::fs::write(&temporary, b"") + .await + .context("preflighting burn-in report output")?; + tokio::fs::remove_file(temporary).await?; + Ok(()) +} + +async fn write_report_atomic(path: &Path, report: &Report) -> Result<()> { + let temporary = temporary_report_path(path); + tokio::fs::write(&temporary, serde_json::to_vec_pretty(report)?) + .await + .context("writing temporary burn-in report")?; + tokio::fs::rename(&temporary, path) + .await + .context("publishing burn-in report atomically")?; + Ok(()) +} + +fn temporary_report_path(path: &Path) -> PathBuf { + let name = path.file_name().unwrap_or_default().to_string_lossy(); + path.with_file_name(format!(".{name}.{}.tmp", std::process::id())) +} + +fn source_commit() -> String { + std::process::Command::new("git") + .args(["rev-parse", "HEAD"]) + .current_dir(workspace_root()) + .output() + .ok() + .filter(|output| output.status.success()) + .and_then(|output| String::from_utf8(output.stdout).ok()) + .map(|commit| commit.trim().to_string()) + .unwrap_or_else(|| "unknown".to_string()) +} + +fn workspace_root() -> &'static Path { + Path::new(env!("CARGO_MANIFEST_DIR")) + .parent() + .and_then(Path::parent) + .expect("burn-in package must be under workspace packages") +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn discovers_executable_from_cargo_json() { + let messages = br#"{"reason":"compiler-artifact","target":{"name":"other"},"executable":"/tmp/other"} +{"reason":"compiler-artifact","target":{"name":"cipherstash-proxy"},"executable":"/custom/target/release/cipherstash-proxy"} +"#; + assert_eq!( + find_proxy_artifact(messages).unwrap(), + PathBuf::from("/custom/target/release/cipherstash-proxy") + ); + } + + #[test] + fn report_requires_work_and_live_rss() { + let mut report = Report { + status: "failed", + terminal_error: None, + requested_duration_seconds: 1, + actual_elapsed_millis: 1, + concurrency: 1, + started_at_unix_seconds: 0, + artifact: "proxy".into(), + source_commit: "commit".into(), + database: "postgresql://user@localhost:5432/db".into(), + operations: 0, + errors: 0, + proxy_pid: 1, + initial_rss_bytes: 0, + final_rss_bytes: 0, + peak_rss_bytes: 0, + rss_growth_bytes: 0, + memory_samples: vec![], + }; + assert!(validate_report(&report, None).is_err()); + + let operations = AtomicU64::new(7); + let errors = AtomicU64::new(2); + refresh_counters(&mut report, &operations, &errors); + assert_eq!(report.operations, 7); + assert_eq!(report.errors, 2); + } +} diff --git a/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_id.rs b/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_id.rs index f76f221ce..b3e023298 100644 --- a/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_id.rs +++ b/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_id.rs @@ -355,8 +355,8 @@ mod tests { // Test cases that should potentially fail or be handled gracefully let invalid = vec![ format!("SET CIPHERSTASH.KEYSET_ID = {tenant_keyset_id_1}"), // unquoted string - format!("SET CIPHERSTASH.KEYSET_ID = NULL"), - format!("SET CIPHERSTASH.KEYSET_ID = 123"), + "SET CIPHERSTASH.KEYSET_ID = NULL".to_string(), + "SET CIPHERSTASH.KEYSET_ID = 123".to_string(), ]; for sql in invalid { diff --git a/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_name.rs b/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_name.rs index 0d9af1757..82dbd19a1 100644 --- a/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_name.rs +++ b/packages/cipherstash-proxy-integration/src/multitenant/set_keyset_name.rs @@ -378,8 +378,8 @@ mod tests { // Test cases that should potentially fail or be handled gracefully let invalid_cases = vec![ - format!("SET CIPHERSTASH.KEYSET_NAME = test-1"), // unquoted string that is NOT a valid pg Identifier - format!("SET CIPHERSTASH.KEYSET_NAME = NULL"), // null value + "SET CIPHERSTASH.KEYSET_NAME = test-1".to_string(), // unquoted string that is NOT a valid pg Identifier + "SET CIPHERSTASH.KEYSET_NAME = NULL".to_string(), // null value ]; for invalid_sql in invalid_cases { diff --git a/packages/cipherstash-proxy/src/postgresql/backend.rs b/packages/cipherstash-proxy/src/postgresql/backend.rs index d22730092..1e2a24f11 100644 --- a/packages/cipherstash-proxy/src/postgresql/backend.rs +++ b/packages/cipherstash-proxy/src/postgresql/backend.rs @@ -177,6 +177,17 @@ where client_id = self.context.client_id, msg = "Passthrough enabled" ); + + // A Proxy started against a database with no encrypted columns is + // initially in passthrough mode. DDL on that connection is how an + // encrypted schema can first appear, so publish the reload before + // forwarding ReadyForQuery. This ordering also guarantees that a + // client opening its next connection after ReadyForQuery observes + // the newly loaded schema and encrypt configuration. + if matches!(code.into(), BackendCode::ReadyForQuery) { + self.context.reload_schema_if_changed().await; + } + self.write_with_flush(bytes).await?; // The frontend starts a session and enqueues an execute for every @@ -280,9 +291,7 @@ where client_id = self.context.client_id, msg = "ReadyForQuery" ); - if self.context.schema_changed() { - self.context.reload_schema().await; - } + self.context.reload_schema_if_changed().await; } code => { @@ -832,6 +841,49 @@ mod tests { backend_message(b'E', b"SERROR\0CXX000\0Mboom\0\0") } + /// `'Z'` ReadyForQuery, with an idle transaction status. + fn ready_for_query_bytes() -> BytesMut { + backend_message(b'Z', b"I") + } + + #[tokio::test] + async fn passthrough_reloads_changed_schema_before_ready_for_query() { + let config = Arc::new(TandemConfig::for_testing()); + let encrypt_config = Arc::new(EncryptConfig::default()); + let schema = Arc::new(Schema::new("public")); + let (reload_sender, mut reload_receiver) = mpsc::unbounded_channel(); + let context = Context::new( + 1, + config, + encrypt_config, + schema, + TestService {}, + reload_sender, + ); + context.set_schema_changed(); + + let reload_task = tokio::spawn(async move { + let Some(crate::proxy::ReloadCommand::DatabaseSchema(responder)) = + reload_receiver.recv().await + else { + panic!("expected a database schema reload command"); + }; + responder.send(true).expect("reload receiver must be open"); + }); + + let (client_sender, mut client_receiver) = mpsc::unbounded_channel(); + let reader = Cursor::new(ready_for_query_bytes().to_vec()); + let mut backend = Backend::new(client_sender, reader, context); + + backend.rewrite().await.unwrap(); + reload_task.await.unwrap(); + assert_eq!( + client_receiver.recv().await.unwrap(), + ready_for_query_bytes() + ); + assert!(!backend.context.take_schema_changed()); + } + /// Regression test for BUG-300 (passthrough memory leak). /// /// The frontend enqueues a session + execute for *every* statement. Those diff --git a/packages/cipherstash-proxy/src/postgresql/context/mod.rs b/packages/cipherstash-proxy/src/postgresql/context/mod.rs index d42e015cf..3ea0cbd54 100644 --- a/packages/cipherstash-proxy/src/postgresql/context/mod.rs +++ b/packages/cipherstash-proxy/src/postgresql/context/mod.rs @@ -28,7 +28,7 @@ pub use statement_metadata::StatementMetadata; use std::{ collections::{HashMap, VecDeque}, sync::{ - atomic::{AtomicU64, Ordering}, + atomic::{AtomicBool, AtomicU64, Ordering}, Arc, LazyLock, RwLock, }, time::{Duration, Instant}, @@ -70,7 +70,7 @@ where portals: Arc>>, describe: Arc>, execute: Arc>, - schema_changed: Arc>, + schema_changed: Arc, session_metrics: Arc>, table_resolver: Arc, unsafe_disable_mapping: bool, @@ -185,7 +185,7 @@ where portals: Arc::new(RwLock::new(HashMap::new())), describe: Arc::new(RwLock::from(Queue::new())), execute: Arc::new(RwLock::from(Queue::new())), - schema_changed: Arc::new(RwLock::from(false)), + schema_changed: Arc::new(AtomicBool::new(false)), session_metrics: Arc::new(RwLock::from(Queue::new())), table_resolver: Arc::new(TableResolver::new_editable(schema)), client_id, @@ -567,11 +567,11 @@ where client_id = self.client_id, msg = "Schema changed" ); - let _ = self.schema_changed.write().map(|mut guard| *guard = true); + self.schema_changed.store(true, Ordering::Release); } - pub fn schema_changed(&self) -> bool { - self.schema_changed.read().ok().is_some_and(|s| *s) + pub fn take_schema_changed(&self) -> bool { + self.schema_changed.swap(false, Ordering::AcqRel) } pub fn get_table_resolver(&self) -> Arc { @@ -768,7 +768,7 @@ where self.encryption.decrypt(keyset_id, ciphertexts).await } - pub async fn reload_schema(&self) { + pub async fn reload_schema(&self) -> bool { let (responder, receiver) = oneshot::channel(); match self .reload_sender @@ -780,18 +780,22 @@ where msg = "Database schema could not be reloaded", error = err.to_string() ); + return false; } } debug!(target: CONTEXT, msg = "Waiting for schema reload"); let response = receiver.await; debug!(target: CONTEXT, msg = "Database schema reloaded", ?response); + matches!(response, Ok(true)) } /// Reload schema if it has changed since last check. pub async fn reload_schema_if_changed(&self) { - if self.schema_changed() { - self.reload_schema().await; + if self.take_schema_changed() && !self.reload_schema().await { + // Preserve the dirty state when the reload task is unavailable so + // a later statement can retry instead of silently losing the DDL. + self.set_schema_changed(); } } @@ -1065,7 +1069,7 @@ mod tests { messages::{Name, Target}, Column, }, - proxy::{EncryptConfig, EncryptionService}, + proxy::{EncryptConfig, EncryptionService, ReloadCommand}, TandemConfig, }; use cipherstash_client::IdentifiedBy; @@ -1117,6 +1121,68 @@ mod tests { ) } + #[tokio::test] + async fn successful_schema_reload_consumes_change_flag_once() { + let config = Arc::new(TandemConfig::for_testing()); + let encrypt_config = Arc::new(EncryptConfig::default()); + let schema = Arc::new(Schema::new("public")); + let (reload_sender, mut reload_receiver) = mpsc::unbounded_channel(); + let context = Context::new( + 1, + config, + encrypt_config, + schema, + TestService {}, + reload_sender, + ); + let reload_task = tokio::spawn(async move { + let Some(ReloadCommand::DatabaseSchema(responder)) = reload_receiver.recv().await + else { + panic!("expected database schema reload"); + }; + responder.send(true).expect("reload receiver is alive"); + tokio::time::timeout(std::time::Duration::from_millis(20), reload_receiver.recv()) + .await + .is_err() + }); + + context.set_schema_changed(); + context.reload_schema_if_changed().await; + context.reload_schema_if_changed().await; + + assert!(!context.take_schema_changed()); + assert!(reload_task.await.expect("reload task did not panic")); + } + + #[tokio::test] + async fn failed_schema_reload_keeps_change_flag_for_retry() { + let config = Arc::new(TandemConfig::for_testing()); + let encrypt_config = Arc::new(EncryptConfig::default()); + let schema = Arc::new(Schema::new("public")); + let (reload_sender, mut reload_receiver) = mpsc::unbounded_channel(); + let context = Context::new( + 1, + config, + encrypt_config, + schema, + TestService {}, + reload_sender, + ); + let reload_task = tokio::spawn(async move { + let Some(ReloadCommand::DatabaseSchema(responder)) = reload_receiver.recv().await + else { + panic!("expected database schema reload"); + }; + responder.send(false).expect("reload receiver is alive"); + }); + + context.set_schema_changed(); + context.reload_schema_if_changed().await; + + reload_task.await.expect("reload task did not panic"); + assert!(context.take_schema_changed()); + } + fn statement() -> Statement { Statement { param_columns: vec![], diff --git a/packages/cipherstash-proxy/src/postgresql/message_buffer.rs b/packages/cipherstash-proxy/src/postgresql/message_buffer.rs index 750b79b2c..15fcd98f7 100644 --- a/packages/cipherstash-proxy/src/postgresql/message_buffer.rs +++ b/packages/cipherstash-proxy/src/postgresql/message_buffer.rs @@ -21,7 +21,7 @@ impl MessageBuffer { } pub fn drain(&mut self) -> Vec { - self.buffer.drain(..).collect() + std::mem::take(&mut self.buffer) } pub fn clear(&mut self) { diff --git a/packages/cipherstash-proxy/src/proxy/encrypt_config/manager.rs b/packages/cipherstash-proxy/src/proxy/encrypt_config/manager.rs index 759074549..6b923fc57 100644 --- a/packages/cipherstash-proxy/src/proxy/encrypt_config/manager.rs +++ b/packages/cipherstash-proxy/src/proxy/encrypt_config/manager.rs @@ -67,19 +67,21 @@ impl EncryptConfigManager { self.encrypt_config.load().is_empty() } - pub async fn reload(&self) { + pub async fn reload(&self) -> bool { match load_encrypt_config_with_retry(&self.config).await { Ok(reloaded) => { debug!(target: ENCRYPT_CONFIG, msg = "Reloaded encrypt configuration"); self.encrypt_config.swap(Arc::new(reloaded)); + true } Err(err) => { warn!( msg = "Error reloading encrypt configuration", error = err.to_string() ); + false } - }; + } } } diff --git a/packages/cipherstash-proxy/src/proxy/mod.rs b/packages/cipherstash-proxy/src/proxy/mod.rs index be45ed321..091fcba18 100644 --- a/packages/cipherstash-proxy/src/proxy/mod.rs +++ b/packages/cipherstash-proxy/src/proxy/mod.rs @@ -23,7 +23,7 @@ pub type ReloadSender = UnboundedSender; type ReloadReceiver = UnboundedReceiver; -pub type ReloadResponder = Sender<()>; +pub type ReloadResponder = Sender; /// SQL Statement for loading database schema. /// @@ -114,13 +114,13 @@ impl Proxy { debug!(msg = "ReloadCommand received", ?command); match command { ReloadCommand::DatabaseSchema(responder) => { - schema_manager.reload().await; - encrypt_config_manager.reload().await; - let _ = responder.send(()); + let schema_reloaded = schema_manager.reload().await; + let encrypt_config_reloaded = encrypt_config_manager.reload().await; + let _ = responder.send(schema_reloaded && encrypt_config_reloaded); } ReloadCommand::EncryptSchema(responder) => { - encrypt_config_manager.reload().await; - let _ = responder.send(()); + let reloaded = encrypt_config_manager.reload().await; + let _ = responder.send(reloaded); } } } diff --git a/packages/cipherstash-proxy/src/proxy/schema/manager.rs b/packages/cipherstash-proxy/src/proxy/schema/manager.rs index b8e6d8fcc..94eea3c59 100644 --- a/packages/cipherstash-proxy/src/proxy/schema/manager.rs +++ b/packages/cipherstash-proxy/src/proxy/schema/manager.rs @@ -28,19 +28,21 @@ impl SchemaManager { self.schema.load().clone() } - pub async fn reload(&self) { + pub async fn reload(&self) -> bool { match load_schema_with_retry(&self.config).await { Ok(reloaded) => { debug!(target: SCHEMA, msg = "Reloaded database schema"); self.schema.swap(Arc::new(reloaded)); + true } Err(err) => { warn!( msg = "Error reloading database schema", error = err.to_string() ); + false } - }; + } } }