diff --git a/CHANGELOG.md b/CHANGELOG.md index 2b7186a..d8cdc6e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- Knot event firehose: durable `events` table plus a `GET /events` WebSocket + that backfills from `?cursor=` then live-tails via Postgres `LISTEN/NOTIFY`. + Push handlers emit `sh.tangled.git.refUpdate` records (with `meta.isDefaultRef`, + `meta.commitCount.byEmail`, and `meta.langBreakdown` on cache hit) inside + the same push transaction so events are atomic with the ref edits they + describe. Cross-replica ordering is enforced by a transaction-scoped + advisory lock; the `count_commits_by_email` walk implements proper + `git rev-list old..new` semantics (proptested against synthetic DAGs). + A `stress_events` test (gated by `STRESS_LEN`, run via `just stress`) checks + correctness and reports throughput / catch-up margin under load, and a + criterion bench (`just bench`) characterizes the lock-hold path's + contention curve with saveable baselines for regression detection. - Postgres-backed Git storage: async `ObjectStore` and `RefStore` traits with an sqlx `PgBackend` and a gitoxide `FileBackend` for parity testing. - Smart-HTTP clone: `info_refs` advertisement and `git-upload-pack` over HTTP, diff --git a/Cargo.lock b/Cargo.lock index f4bf6e2..2300c51 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -108,6 +108,12 @@ dependencies = [ "libc", ] +[[package]] +name = "anes" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4b46cbb362ab8752921c97e041f5e366ee6297bd428a31275b9fcf1e380f7299" + [[package]] name = "anstream" version = "1.0.0" @@ -312,6 +318,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" dependencies = [ "axum-core", + "base64", "bytes", "form_urlencoded", "futures-util", @@ -330,8 +337,10 @@ dependencies = [ "serde_json", "serde_path_to_error", "serde_urlencoded", + "sha1 0.10.6", "sync_wrapper", "tokio", + "tokio-tungstenite 0.29.0", "tower", "tower-layer", "tower-service", @@ -628,6 +637,12 @@ dependencies = [ "serde", ] +[[package]] +name = "cast" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37b2a672a2cb129a2e41c10b1224bb368f9f37a2b16b612598138befd7b37eb5" + [[package]] name = "cbc" version = "0.1.2" @@ -983,6 +998,44 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "criterion" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2b12d017a929603d80db1831cd3a24082f8137ce19c69e6447f54f5fc8d692f" +dependencies = [ + "anes", + "cast", + "ciborium", + "clap", + "criterion-plot", + "futures", + "is-terminal", + "itertools 0.10.5", + "num-traits", + "once_cell", + "oorandom", + "plotters", + "rayon", + "regex", + "serde", + "serde_derive", + "serde_json", + "tinytemplate", + "tokio", + "walkdir", +] + +[[package]] +name = "criterion-plot" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b50826342786a51a89e2da3a28f1c32b06e387201bc2d19791f622c673706b1" +dependencies = [ + "cast", + "itertools 0.10.5", +] + [[package]] name = "critical-section" version = "1.2.0" @@ -3560,6 +3613,12 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hermit-abi" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c" + [[package]] name = "hex" version = "0.4.3" @@ -4049,12 +4108,32 @@ version = "2.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d98f6fed1fde3f8c21bc40a1abb88dd75e67924f9cffc3ef95607bad8017f8e2" +[[package]] +name = "is-terminal" +version = "0.4.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" +dependencies = [ + "hermit-abi", + "libc", + "windows-sys 0.61.2", +] + [[package]] name = "is_terminal_polyfill" version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" +[[package]] +name = "itertools" +version = "0.10.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0fd2260e829bddf4cb6ea802289de2f86d6a7a690192fbe91b3f46e0f2c8473" +dependencies = [ + "either", +] + [[package]] name = "itertools" version = "0.14.0" @@ -4549,6 +4628,7 @@ dependencies = [ "async-trait", "bstr", "bytes", + "criterion", "ctor", "gix-hash 0.25.0", "gix-object 0.60.0", @@ -4579,9 +4659,11 @@ dependencies = [ "gix-packetline", "gix-ref 0.63.0", "http-body-util", + "jacquard-common", "mnemosyne_git", "mnemosyne_postgres", "serde", + "serde_json", "tempfile", "thiserror 2.0.18", "tokio", @@ -4618,8 +4700,10 @@ dependencies = [ "axum", "base64", "bstr", + "bytes", "ctor", "flate2", + "futures", "gix 0.83.0", "gix-actor 0.41.0", "gix-date", @@ -4646,6 +4730,7 @@ dependencies = [ "tempfile", "testcontainers-modules", "tokio", + "tokio-tungstenite 0.29.0", "tower", ] @@ -4871,6 +4956,12 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe" +[[package]] +name = "oorandom" +version = "11.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e" + [[package]] name = "opaque-debug" version = "0.3.1" @@ -5225,6 +5316,34 @@ version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b4596b6d070b27117e987119b4dac604f3c58cfb0b191112e24771b2faeac1a6" +[[package]] +name = "plotters" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5aeb6f403d7a4911efb1e33402027fc44f29b5bf6def3effcc22d7bb75f2b747" +dependencies = [ + "num-traits", + "plotters-backend", + "plotters-svg", + "wasm-bindgen", + "web-sys", +] + +[[package]] +name = "plotters-backend" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df42e13c12958a16b3f7f4386b9ab1f3e7933914ecea48da7139435263a4172a" + +[[package]] +name = "plotters-svg" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "51bae2ac328883f7acdfea3d66a7c35751187f870bc81f94563733a154d7a670" +dependencies = [ + "plotters-backend", +] + [[package]] name = "poly1305" version = "0.8.0" @@ -5411,7 +5530,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "27c6023962132f4b30eb4c172c91ce92d933da334c59c23cddee82358ddafb0b" dependencies = [ "anyhow", - "itertools", + "itertools 0.14.0", "proc-macro2", "quote", "syn", @@ -6915,7 +7034,7 @@ dependencies = [ "ferroid", "futures", "http", - "itertools", + "itertools 0.14.0", "log", "memchr", "parse-display", @@ -7029,6 +7148,16 @@ dependencies = [ "zerovec", ] +[[package]] +name = "tinytemplate" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "be4d6b5f19ff7664e8c98d03e2139cb510db9b0a60b55f8e8709b689d939b6bc" +dependencies = [ + "serde", + "serde_json", +] + [[package]] name = "tinyvec" version = "1.11.0" diff --git a/Cargo.toml b/Cargo.toml index 2ed476f..8db159a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -55,8 +55,9 @@ gix-protocol = "0.61" # OwnShared resolves to Arc (not Rc), making gix_ref::file::Store Send+Sync. gix-features = { version = "0.48", features = ["parallel"] } -# HTTP server -axum = "0.8" +# HTTP server. `ws` adds the WebSocket extractor used by the +# `/events` firehose endpoint in the main binary. +axum = { version = "0.8", features = ["ws"] } tower = { version = "0.5", features = ["util"] } http-body-util = "0.1" clap = { version = "4", features = ["derive", "env"] } @@ -106,6 +107,13 @@ tempfile = "3" tokio-test = "0.4" testcontainers-modules = { version = "0.15", features = ["postgres"] } proptest = "1.11" + +# Microbenchmarks. Used by `mnemosyne_postgres/benches/*` (and +# more in the future) to characterize hot-path performance with +# baselines that `cargo bench --save-baseline` / `--baseline` can +# compare against in CI. `async_tokio` lets bench iterations +# `.await` directly without juggling a runtime per case. +criterion = { version = "0.5", features = ["async_tokio", "html_reports"] } # gix porcelain — used in the e2e clone proptest only (dev-dep) gix = { version = "0.83", default-features = false, features = ["blocking-network-client", "blocking-http-transport-reqwest"] } diff --git a/justfile b/justfile index 08f179d..512c4c7 100644 --- a/justfile +++ b/justfile @@ -61,6 +61,47 @@ stress depth="100000": ' EXIT INT TERM STRESS_LEN={{depth}} cargo nextest run -p mnemosyne_tests -E 'test(/^stress(::|_)/)' +# Run criterion benchmarks against a one-shot Postgres container. +# +# Pairs with `just stress`: stress runs the design's correctness +# claims under load (binary outcome), this runs the perf +# characterization (continuous signal). +# +# just bench # run every bench +# just bench event_insert # filter to a bench name +# just bench event_insert/uncontended +# +# Criterion saves baselines under `target/criterion/`. To compare +# against a known-good main, run on main first: +# +# git checkout main && just bench --save-baseline main +# git checkout my-branch && just bench --baseline main +bench *args: + #!/usr/bin/env bash + set -euo pipefail + fixture_dir="${CARGO_TARGET_DIR:-target}/bench-fixtures" + mkdir -p "$fixture_dir" + container_name="mnemosyne-bench-pg-$$" + trap ' + docker rm -f "$container_name" > /dev/null 2>&1 || true + ' EXIT INT TERM + docker run -d \ + --name "$container_name" \ + -e POSTGRES_PASSWORD=postgres \ + -p 5432 \ + postgres:17 > /dev/null + # Wait for readiness — docker reports started before pg actually + # listens. ~1 s on a warm host, longer on cold. + for _ in $(seq 1 60); do + if docker exec "$container_name" pg_isready -U postgres > /dev/null 2>&1; then + break + fi + sleep 0.5 + done + port="$(docker port "$container_name" 5432/tcp | head -1 | awk -F: '{print $NF}')" + export MNEMOSYNE_BENCH_PG_URL="postgres://postgres:postgres@127.0.0.1:${port}/postgres" + cargo bench -p mnemosyne_postgres --bench event_emission -- {{args}} + # Run the property tests with a custom case count (default: 64). # # just proptest # 64 cases per property diff --git a/mnemosyne/src/main.rs b/mnemosyne/src/main.rs index 6ebf73d..6234b8f 100644 --- a/mnemosyne/src/main.rs +++ b/mnemosyne/src/main.rs @@ -11,22 +11,23 @@ use std::sync::Arc; use anyhow::Context; +use axum::Router; use axum::body::Bytes; use axum::extract::{Path, Query, State}; use axum::http::StatusCode; use axum::response::IntoResponse; use axum::routing::{get, post}; -use axum::Router; use clap::Parser; use gix_hash::Kind as HashKind; -use mnemosyne_postgres::{ensure_repo, lookup_repo_by_did, migrate, resolve_repo_did, PgBackend}; +use mnemosyne_postgres::{PgBackend, ensure_repo, lookup_repo_by_did, migrate, resolve_repo_did}; +use mnemosyne_protocol::events_stream::{self as events_stream, EventsStreamState}; use mnemosyne_protocol::info_refs::InfoRefsQuery; -use mnemosyne_protocol::{info_refs, receive_pack, upload_pack, RepoState}; +use mnemosyne_protocol::{RepoState, info_refs, receive_pack, upload_pack}; use mnemosyne_ssh::keys::{Algorithm, PrivateKey}; use mnemosyne_ssh::{Server as SshServer, ServerConfig as SshServerConfig}; -use mnemosyne_xrpc::{router as xrpc_router, XrpcState}; -use sqlx::postgres::PgPoolOptions; +use mnemosyne_xrpc::{XrpcState, router as xrpc_router}; use sqlx::PgPool; +use sqlx::postgres::PgPoolOptions; use tokio::net::TcpListener; mod seed; @@ -48,7 +49,7 @@ mod seed; fn init_tracing() { use tracing_subscriber::fmt::format::FmtSpan; use tracing_subscriber::prelude::*; - use tracing_subscriber::{fmt, EnvFilter}; + use tracing_subscriber::{EnvFilter, fmt}; let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| { EnvFilter::new( @@ -70,7 +71,10 @@ fn init_tracing() { } #[derive(Debug, Parser)] -#[command(name = "mnemosyne", about = "Postgres-backed Git HTTP server (Tangled knot)")] +#[command( + name = "mnemosyne", + about = "Postgres-backed Git HTTP server (Tangled knot)" +)] struct Config { /// Postgres connection URL. #[arg(long, env = "MNEMOSYNE_DATABASE_URL")] @@ -95,11 +99,7 @@ struct Config { repo_did: String, /// AT Protocol DID of the repository owner. - #[arg( - long, - default_value = "did:web:localhost", - env = "MNEMOSYNE_OWNER_DID" - )] + #[arg(long, default_value = "did:web:localhost", env = "MNEMOSYNE_OWNER_DID")] owner_did: String, /// Human-readable name of the seeded repository. @@ -165,7 +165,11 @@ async fn git_info_refs( Path(repo_did): Path, query: Query, ) -> impl IntoResponse { - let Some(repo_did) = lookup_repo_by_did(&app.pool, &repo_did).await.ok().flatten() else { + let Some(repo_did) = lookup_repo_by_did(&app.pool, &repo_did) + .await + .ok() + .flatten() + else { return repo_not_found(); }; info_refs::handler(State(app.repo_state(repo_did)), query) @@ -178,7 +182,11 @@ async fn git_upload_pack( Path(repo_did): Path, body: Bytes, ) -> impl IntoResponse { - let Some(repo_did) = lookup_repo_by_did(&app.pool, &repo_did).await.ok().flatten() else { + let Some(repo_did) = lookup_repo_by_did(&app.pool, &repo_did) + .await + .ok() + .flatten() + else { return repo_not_found(); }; upload_pack::handler(State(app.repo_state(repo_did)), body) @@ -191,7 +199,11 @@ async fn git_receive_pack( Path(repo_did): Path, body: Bytes, ) -> impl IntoResponse { - let Some(repo_did) = lookup_repo_by_did(&app.pool, &repo_did).await.ok().flatten() else { + let Some(repo_did) = lookup_repo_by_did(&app.pool, &repo_did) + .await + .ok() + .flatten() + else { return repo_not_found(); }; receive_pack::handler(State(app.repo_state(repo_did)), body) @@ -310,6 +322,15 @@ async fn main() -> anyhow::Result<()> { owner_did: Some(cfg.owner_did.clone()), }); + // /events WebSocket — drains the durable `events` table to + // subscribers via LISTEN/NOTIFY. Mounted on the root router so + // it isn't subject to /xrpc-style routing. + let events_router: Router = Router::new() + .route("/events", get(events_stream::handler)) + .with_state(EventsStreamState { + pool: state.pool.clone(), + }); + let app = Router::new() // /:repo_did/... .route("/{repo_did}/info/refs", get(git_info_refs)) @@ -329,6 +350,7 @@ async fn main() -> anyhow::Result<()> { post(git_receive_pack_named), ) .with_state(state) + .merge(events_router) .nest("/xrpc", xrpc); // HTTP listener. diff --git a/mnemosyne_postgres/Cargo.toml b/mnemosyne_postgres/Cargo.toml index 99e238d..18a3958 100644 --- a/mnemosyne_postgres/Cargo.toml +++ b/mnemosyne_postgres/Cargo.toml @@ -30,9 +30,15 @@ mnemosyne_harness.workspace = true # in the integration tests' `OnceCell` static. See # `tests/common/mod.rs` for the rationale. ctor = "0.4" +criterion.workspace = true +serde_json = "1" tempfile.workspace = true testcontainers-modules.workspace = true tokio = { workspace = true, features = ["rt-multi-thread", "macros"] } +[[bench]] +name = "event_emission" +harness = false + [lints] workspace = true diff --git a/mnemosyne_postgres/benches/event_emission.rs b/mnemosyne_postgres/benches/event_emission.rs new file mode 100644 index 0000000..506ddbc --- /dev/null +++ b/mnemosyne_postgres/benches/event_emission.rs @@ -0,0 +1,228 @@ +//! Microbenchmarks for the event-emission hot path. +//! +//! Pairs with the `stress_events` integration test (in +//! `mnemosyne_tests`) — same code path, different signals: +//! +//! - **`stress_events`** runs the design correctness claim under load: +//! "did anything break?" Binary outcome, gated by `STRESS_LEN`. +//! - **`event_emission`** bench characterizes performance over time: +//! "is this getting slower?" Continuous signal via criterion's +//! baseline-comparison feature. +//! +//! Some overlap (both call [`insert_event_in_tx`]) is intentional — +//! they detect different failure modes (correctness regressions vs +//! performance regressions) and you want both signals. +//! +//! ## Cases +//! +//! - `event_insert/uncontended` — one writer in a loop. Floor cost +//! of `pg_advisory_xact_lock` + INSERT + `pg_notify` + commit. +//! Lock acquisition is uncontended so this is the lower bound; any +//! real contention only adds wait time on top. +//! - `event_insert/contended_{2,8,32}` — N concurrent writers +//! sharing the lock. The shape of throughput / per-event latency +//! vs N tells us whether the advisory-lock design degrades +//! gracefully (linear contention, expected) or pathologically +//! (super-linear / deadlock-y, would force the outbox upgrade). +//! +//! ## Environment +//! +//! `MNEMOSYNE_BENCH_PG_URL` must point at a writable Postgres. The +//! `just bench` recipe boots a fresh container, runs migrations, +//! exports the URL, and tears down on exit. Running `cargo bench` +//! directly requires the operator to provide the URL. +//! +//! ## Baselines (M-series laptop, release build, fresh container) +//! +//! | case | per-iter time | per-event | throughput | +//! |-----------------|---------------|-----------|-------------| +//! | uncontended | ~3.5 ms | 3.5 ms | ~280 ev/s | +//! | contended / 2 | ~3.1 ms | 1.6 ms | ~640 ev/s | +//! | contended / 8 | ~10.9 ms | 1.4 ms | ~730 ev/s | +//! | contended / 32 | ~44.4 ms | 1.4 ms | ~720 ev/s | +//! +//! Things to notice: +//! +//! - Per-event latency *drops* from uncontended → contended because +//! concurrent transactions pipeline through the lock — while one +//! commit drains, the next is already queued. Round-trip latency +//! amortizes. +//! - Aggregate throughput plateaus at ~700 ev/s by N=8 and stays +//! flat at N=32. This is the expected signature of a clean +//! serialization point: ceiling = `1 / lock_hold_time`. +//! **Super-linear degradation past N=8 would be the warning sign** +//! that the advisory-lock design is misbehaving and the outbox +//! upgrade (issue #1) is worth doing. +//! - Bench throughput (~700 ev/s) is lower than the integration +//! stress test (~2 000 ev/s on the same hardware) because the +//! bench uses a fresh container and isolates per-iter overhead; +//! the stress test runs against a warmed pool. Both are valid +//! regressions surface — track them independently. +//! +//! Save a baseline with `--save-baseline ` (the `just bench` +//! recipe passes args through); compare with `--baseline `. + +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{Duration, Instant}; + +use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; +use mnemosyne_postgres::{ensure_repo, insert_event_in_tx, migrate}; +use serde_json::json; +use sqlx::PgPool; +use sqlx::postgres::PgPoolOptions; +use tokio::runtime::Runtime; + +const BENCH_REPO_OWNER: &str = "did:test:bench-owner"; +const BENCH_REPO_NAME: &str = "bench-events"; +const BENCH_REPO_RKEY: &str = "bench-rkey"; + +/// Pull the Postgres URL from the env. We deliberately do NOT fall +/// back to a sentinel default — silently pointing at a wrong DB +/// would be a very confusing failure mode for a perf bench. +fn pg_url() -> String { + std::env::var("MNEMOSYNE_BENCH_PG_URL").unwrap_or_else(|_| { + panic!( + "MNEMOSYNE_BENCH_PG_URL is required to run event_emission \ + benches — see the file header or run via `just bench`.", + ) + }) +} + +/// Lazily-constructed tokio runtime + pool + scratch repo. Built +/// once per bench process so the per-iteration overhead is just the +/// work being measured. +struct Harness { + rt: Runtime, + pool: PgPool, + repo_did: String, +} + +impl Harness { + fn boot() -> Self { + let rt = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build() + .expect("tokio runtime"); + let (pool, repo_did) = rt.block_on(async { + let pool = PgPoolOptions::new() + // Generous so high-N contended cases never starve + // for connections — the bench measures lock-wait, + // not pool-wait. + .max_connections(64) + .connect(&pg_url()) + .await + .expect("connect bench pg"); + migrate(&pool).await.expect("migrate bench pg"); + // Per-process repo so re-running benches in the same + // database doesn't trip the unique constraint. PID + // gives uniqueness across parallel `cargo bench` runs + // pointed at the same DB. + let repo_did = format!("did:test:bench-{}", std::process::id()); + ensure_repo( + &pool, + &repo_did, + BENCH_REPO_OWNER, + BENCH_REPO_NAME, + BENCH_REPO_RKEY, + ) + .await + .expect("ensure bench repo"); + (pool, repo_did) + }); + Self { rt, pool, repo_did } + } +} + +/// Atomically-incremented rkey counter — each emit needs a distinct +/// rkey for realism (the lexicon keys records by TID, all distinct +/// in production) and to keep the bench from accidentally hammering +/// any uniqueness index that might one day land on `events.rkey`. +static RKEY_COUNTER: AtomicU64 = AtomicU64::new(0); + +fn next_rkey() -> String { + let n = RKEY_COUNTER.fetch_add(1, Ordering::Relaxed); + let pid = std::process::id(); + format!("bench-{pid}-{n:020}") +} + +/// One synthetic event emission inside its own transaction. The +/// payload is small so we're measuring lock + INSERT + commit cost, +/// not JSON serialization. Realistic refUpdate payloads are ~300 B; +/// adjust if a future bench wants to measure payload-size sensitivity. +async fn emit_one(pool: &PgPool, repo_did: &str) { + let mut tx = pool.begin().await.expect("begin"); + let rkey = next_rkey(); + insert_event_in_tx( + &mut tx, + &rkey, + "sh.tangled.git.refUpdate", + Some(repo_did), + &json!({"bench": true}), + ) + .await + .expect("insert"); + tx.commit().await.expect("commit"); +} + +fn bench_uncontended(c: &mut Criterion) { + let h = Harness::boot(); + let mut group = c.benchmark_group("event_insert"); + group.throughput(Throughput::Elements(1)); + group.bench_function("uncontended", |b| { + b.to_async(&h.rt) + .iter(|| async { emit_one(&h.pool, &h.repo_did).await }); + }); + group.finish(); +} + +fn bench_contended(c: &mut Criterion) { + let h = Harness::boot(); + let pool = Arc::new(h.pool.clone()); + let repo = Arc::new(h.repo_did.clone()); + + let mut group = c.benchmark_group("event_insert"); + for workers in [2usize, 8, 32] { + group.throughput(Throughput::Elements(workers as u64)); + group.bench_with_input( + BenchmarkId::new("contended", workers), + &workers, + |b, &workers| { + let rt = &h.rt; + let pool = pool.clone(); + let repo = repo.clone(); + // `iter_custom` so we control the per-iter setup — + // we time the "spawn N tasks, await all" envelope, + // not the closure-call overhead criterion's default + // iter would impose per task. + b.to_async(rt).iter_custom(move |iters| { + let pool = pool.clone(); + let repo = repo.clone(); + async move { + let mut total = Duration::ZERO; + for _ in 0..iters { + let start = Instant::now(); + let mut handles = Vec::with_capacity(workers); + for _ in 0..workers { + let p = pool.clone(); + let r = repo.clone(); + handles.push(tokio::spawn(async move { + emit_one(&p, &r).await; + })); + } + for h in handles { + h.await.expect("worker join"); + } + total += start.elapsed(); + } + total + } + }); + }, + ); + } + group.finish(); +} + +criterion_group!(benches, bench_uncontended, bench_contended); +criterion_main!(benches); diff --git a/mnemosyne_postgres/migrations/0011_create_events.sql b/mnemosyne_postgres/migrations/0011_create_events.sql new file mode 100644 index 0000000..6acd91b --- /dev/null +++ b/mnemosyne_postgres/migrations/0011_create_events.sql @@ -0,0 +1,60 @@ +-- Knot event stream — durable backing for `GET /events`. +-- +-- One row per emitted event (currently: `sh.tangled.git.refUpdate` +-- records produced by `git-receive-pack`; future: pipeline events, +-- repo lifecycle events, etc.). The appview tails this table over +-- WebSocket to learn about activity on the knot. +-- +-- ## Ordering contract +-- +-- `seq` is a monotonic `BIGSERIAL` cursor. Subscribers query +-- `WHERE seq > $cursor ORDER BY seq` and advance their stored +-- cursor to the latest seq they have processed. +-- +-- Under concurrent writes, Postgres can allocate sequence values to +-- transactions that commit out of order — a subscriber would see +-- seq=N+1 first and then never observe seq=N if the older tx commits +-- afterwards. We sidestep that by requiring every writer to hold a +-- session-level **advisory lock** (`EVENTS_ADVISORY_LOCK_ID` in the +-- Rust side) immediately before INSERT, until commit. The lock is +-- only held for the INSERT + commit window — pack ingestion and ref +-- updates earlier in the push transaction remain concurrent. +-- +-- This is a deliberately simple ordering model. When we add a second +-- emission source (ingester, XRPC mutations, …) that doesn't want +-- to serialize against pushes, the upgrade path is a `stream` column +-- + per-stream advisory lock. See issue #1 architectural overview. +-- +-- ## NOTIFY +-- +-- Writers fire `pg_notify('mnemosyne_events', seq::text)` inside +-- the same transaction. Postgres queues notifications until commit, +-- so subscribers never see a notification for an uncommitted seq. +-- +-- ## Schema +-- +-- - `seq` — global cursor, monotonic under the advisory-lock contract. +-- - `rkey` — TID-format record key for the event (matches lexicon). +-- - `nsid` — record collection NSID, e.g. `sh.tangled.git.refUpdate`. +-- - `repo_did` — repo the event is about, when applicable. Indexed for +-- future per-repo subscription support. +-- - `payload` — full event record as JSONB, exactly what gets sent +-- to subscribers as the `event` field. +-- - `created_at` — wall-clock time at INSERT, populated by Postgres. +-- Wire format converts this to nanoseconds-since-epoch +-- for the appview, matching knotserver's `created` field. +CREATE TABLE events ( + seq BIGSERIAL PRIMARY KEY, + rkey TEXT NOT NULL, + nsid TEXT NOT NULL, + repo_did TEXT, + payload JSONB NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +-- Subscribers ORDER BY seq; the PK index covers this. Per-repo +-- subscription is on the roadmap; index ahead so the cost is paid +-- at write-time (cheap) not at first query under load (visible +-- latency spike). +CREATE INDEX events_repo_did_seq_idx ON events (repo_did, seq) + WHERE repo_did IS NOT NULL; diff --git a/mnemosyne_postgres/src/events.rs b/mnemosyne_postgres/src/events.rs new file mode 100644 index 0000000..669fa86 --- /dev/null +++ b/mnemosyne_postgres/src/events.rs @@ -0,0 +1,162 @@ +//! Knot event-stream persistence — backing store for `GET /events`. +//! +//! Writers (currently: the receive-pack handler emitting +//! `sh.tangled.git.refUpdate`) call [`insert_event_in_tx`] inside +//! their existing transaction. Readers (the WebSocket handler in the +//! main binary) poll [`fetch_events_since`] and tail +//! [`EVENTS_NOTIFY_CHANNEL`] via [`sqlx::postgres::PgListener`]. +//! +//! ## Ordering +//! +//! Sequence values must be monotonic against commit order so that +//! subscribers can use `seq` as a resumable cursor without missing +//! events. Postgres allocates `BIGSERIAL` values at INSERT time, +//! independent of commit order — without coordination, two concurrent +//! transactions can commit `seq=N+1` before `seq=N`, and a +//! subscriber that advances past `N+1` will never observe `N`. +//! +//! [`insert_event_in_tx`] acquires [`EVENTS_ADVISORY_LOCK_ID`] as a +//! transaction-scoped advisory lock immediately before the INSERT. +//! The lock auto-releases on transaction end (commit or rollback), +//! so the serialized window is just `INSERT + commit` — sub-ms in +//! practice. Earlier work in the push transaction (object inserts, +//! ref edits) remains concurrent. +//! +//! This is the simplest design that produces a correct totally-ordered +//! stream. If we later add a second emission source that shouldn't +//! serialize against pushes, the upgrade path is a `stream` column +//! plus a per-stream lock id. See issue #1. + +use sqlx::{PgPool, Postgres, Transaction}; + +use crate::PgError; + +/// Postgres `pg_notify` channel name. Subscribers `LISTEN` on this +/// to wake up when new events have been committed. +pub const EVENTS_NOTIFY_CHANNEL: &str = "mnemosyne_events"; + +/// Advisory lock id used to serialize event INSERTs. +/// +/// Bytes spell `"mnem_evt"` in ASCII so it's recognizable in +/// `pg_locks`, but the actual value is arbitrary — any constant +/// unique to this codebase works. +pub const EVENTS_ADVISORY_LOCK_ID: i64 = 0x6D6E_656D_5F65_7674; + +/// One row of the `events` table, denormalized for the wire format. +/// +/// `created_nanos` is the wall-clock time of the row's INSERT, +/// rendered as nanoseconds since the Unix epoch to match the field +/// shape knotserver emits on its `/events` WebSocket. +#[derive(Debug, Clone)] +pub struct EventRow { + pub seq: i64, + pub rkey: String, + pub nsid: String, + pub repo_did: Option, + pub payload: serde_json::Value, + pub created_nanos: i64, +} + +/// Persist one event inside the caller's open transaction. +/// +/// Takes the advisory lock so commit order matches `seq` order across +/// concurrent writers, INSERTs the row, and queues a `pg_notify` +/// carrying the allocated seq as its payload (subscribers can use +/// this to bound their follow-up `SELECT`). +/// +/// Returns the freshly-allocated seq. The lock is released +/// automatically when the caller commits or rolls back. +pub async fn insert_event_in_tx( + tx: &mut Transaction<'_, Postgres>, + rkey: &str, + nsid: &str, + repo_did: Option<&str>, + payload: &serde_json::Value, +) -> Result { + // Acquire BEFORE the INSERT so any concurrent writer is blocked + // from allocating its seq until we commit. Holding-window = + // INSERT + commit only; pack ingest + ref edits ran earlier. + sqlx::query("SELECT pg_advisory_xact_lock($1)") + .bind(EVENTS_ADVISORY_LOCK_ID) + .execute(&mut **tx) + .await?; + + let (seq,): (i64,) = sqlx::query_as( + "INSERT INTO events (rkey, nsid, repo_did, payload) \ + VALUES ($1, $2, $3, $4) \ + RETURNING seq", + ) + .bind(rkey) + .bind(nsid) + .bind(repo_did) + .bind(sqlx::types::Json(payload)) + .fetch_one(&mut **tx) + .await?; + + // `pg_notify` is deferred to commit by Postgres, so subscribers + // can't see this notification before the row is visible. Payload + // is the seq as a decimal string so the listener can do a bounded + // `WHERE seq <= $payload` if it wants to. + sqlx::query("SELECT pg_notify($1, $2)") + .bind(EVENTS_NOTIFY_CHANNEL) + .bind(seq.to_string()) + .execute(&mut **tx) + .await?; + + Ok(seq) +} + +/// Row shape used by [`fetch_events_since`] for the sqlx tuple +/// decode. Named `pub(crate)` so the public API stays the +/// denormalized [`EventRow`]. +type EventRowTuple = ( + i64, + String, + String, + Option, + sqlx::types::Json, + i64, +); + +/// Drain up to `limit` rows from `events` with `seq > cursor`, +/// ordered by `seq` ascending. +/// +/// Subscribers call this from their cursor on connect (backfill) and +/// again after every `LISTEN` notification (live tail). They advance +/// their stored cursor to the largest seq returned. +/// +/// `created_at` is rendered as nanoseconds-since-epoch in-query to +/// avoid pulling a `chrono` / `time` feature into sqlx — the wire +/// format is `nanos` anyway, matching knotserver's `created` field. +pub async fn fetch_events_since( + pool: &PgPool, + cursor: i64, + limit: i64, +) -> Result, PgError> { + let rows: Vec = sqlx::query_as( + "SELECT seq, rkey, nsid, repo_did, payload, \ + (EXTRACT(EPOCH FROM created_at) * 1000000000)::BIGINT \ + FROM events \ + WHERE seq > $1 \ + ORDER BY seq \ + LIMIT $2", + ) + .bind(cursor) + .bind(limit) + .fetch_all(pool) + .await?; + + Ok(rows + .into_iter() + .map( + |(seq, rkey, nsid, repo_did, payload, created_nanos)| EventRow { + seq, + rkey, + nsid, + repo_did, + payload: payload.0, + created_nanos, + }, + ) + .collect()) +} diff --git a/mnemosyne_postgres/src/lib.rs b/mnemosyne_postgres/src/lib.rs index 50118f0..bdf73c0 100644 --- a/mnemosyne_postgres/src/lib.rs +++ b/mnemosyne_postgres/src/lib.rs @@ -1,9 +1,16 @@ //! Postgres-backed implementation of [`mnemosyne_git`] storage traits. +pub mod events; + +pub use events::{ + EVENTS_ADVISORY_LOCK_ID, EVENTS_NOTIFY_CHANNEL, EventRow, fetch_events_since, + insert_event_in_tx, +}; + use async_trait::async_trait; use bstr::BString; use bytes::Bytes; -use gix_hash::{oid, Kind as HashKind, ObjectId}; +use gix_hash::{Kind as HashKind, ObjectId, oid}; use gix_object::{Header, Kind}; use gix_ref::transaction::{Change, PreviousValue, RefEdit}; use gix_ref::{FullName, FullNameRef, PartialNameRef, Reference, Target}; @@ -50,7 +57,9 @@ pub enum PgError { source: gix_object::decode::Error, }, - #[error("connectivity check failed: object {missing} (referenced from {referenced_from}) is not present")] + #[error( + "connectivity check failed: object {missing} (referenced from {referenced_from}) is not present" + )] MissingObject { missing: ObjectId, referenced_from: ObjectId, @@ -131,15 +140,11 @@ pub async fn ensure_repo( /// `None` if this ingester has never written a cursor yet (i.e., /// first start — caller should subscribe from the firehose's live /// tail). -pub async fn read_ingest_cursor( - pool: &PgPool, - id: &str, -) -> Result, PgError> { - let row: Option<(i64,)> = - sqlx::query_as("SELECT cursor FROM ingest_cursor WHERE id = $1") - .bind(id) - .fetch_optional(pool) - .await?; +pub async fn read_ingest_cursor(pool: &PgPool, id: &str) -> Result, PgError> { + let row: Option<(i64,)> = sqlx::query_as("SELECT cursor FROM ingest_cursor WHERE id = $1") + .bind(id) + .fetch_optional(pool) + .await?; Ok(row.map(|(c,)| c)) } @@ -148,11 +153,7 @@ pub async fn read_ingest_cursor( /// as the event's effect on the rest of the schema so that a crash /// can't lose progress without rolling back the change too — see the /// `apply_event` flow in `mnemosyne_ingest`. -pub async fn advance_ingest_cursor<'e, E>( - executor: E, - id: &str, - cursor: i64, -) -> Result<(), PgError> +pub async fn advance_ingest_cursor<'e, E>(executor: E, id: &str, cursor: i64) -> Result<(), PgError> where E: sqlx::Executor<'e, Database = Postgres>, { @@ -182,11 +183,10 @@ pub async fn lookup_repo_owner<'e, E>( where E: sqlx::Executor<'e, Database = Postgres>, { - let row: Option<(String,)> = - sqlx::query_as("SELECT owner_did FROM repos WHERE did = $1") - .bind(repo_did) - .fetch_optional(executor) - .await?; + let row: Option<(String,)> = sqlx::query_as("SELECT owner_did FROM repos WHERE did = $1") + .bind(repo_did) + .fetch_optional(executor) + .await?; Ok(row.map(|(owner_did,)| owner_did)) } @@ -194,15 +194,11 @@ where /// /// Returns `Some(repo_did)` if found, `None` otherwise. Useful for the /// `/:repo_did` routing form where the DID is supplied directly. -pub async fn lookup_repo_by_did( - pool: &PgPool, - repo_did: &str, -) -> Result, PgError> { - let row: Option<(String,)> = - sqlx::query_as("SELECT did FROM repos WHERE did = $1") - .bind(repo_did) - .fetch_optional(pool) - .await?; +pub async fn lookup_repo_by_did(pool: &PgPool, repo_did: &str) -> Result, PgError> { + let row: Option<(String,)> = sqlx::query_as("SELECT did FROM repos WHERE did = $1") + .bind(repo_did) + .fetch_optional(pool) + .await?; Ok(row.map(|(did,)| did)) } @@ -222,12 +218,11 @@ pub async fn describe_repo( pool: &PgPool, repo_did: &str, ) -> Result, PgError> { - let row: Option<(String, String, String, String)> = sqlx::query_as( - "SELECT did, owner_did, name, rkey FROM repos WHERE did = $1", - ) - .bind(repo_did) - .fetch_optional(pool) - .await?; + let row: Option<(String, String, String, String)> = + sqlx::query_as("SELECT did, owner_did, name, rkey FROM repos WHERE did = $1") + .bind(repo_did) + .fetch_optional(pool) + .await?; Ok(row.map(|(did, owner_did, name, rkey)| RepoDescription { repo_did: did, owner_did, @@ -289,12 +284,11 @@ pub async fn find_public_key_did( pool: &PgPool, key_openssh: &str, ) -> Result, PgError> { - let row: Option<(String,)> = sqlx::query_as( - "SELECT did FROM public_keys WHERE key = $1 LIMIT 1", - ) - .bind(key_openssh) - .fetch_optional(pool) - .await?; + let row: Option<(String,)> = + sqlx::query_as("SELECT did FROM public_keys WHERE key = $1 LIMIT 1") + .bind(key_openssh) + .fetch_optional(pool) + .await?; Ok(row.map(|(did,)| did)) } @@ -363,11 +357,7 @@ where /// Remove the `public_keys` row identified by `(did, rkey)`. Silent /// no-op if the row doesn't exist — Delete events for unknown rkeys /// are valid (e.g. cursor replay) and must not error. -pub async fn delete_public_key<'e, E>( - executor: E, - did: &str, - rkey: &str, -) -> Result<(), PgError> +pub async fn delete_public_key<'e, E>(executor: E, did: &str, rkey: &str) -> Result<(), PgError> where E: sqlx::Executor<'e, Database = Postgres>, { @@ -589,13 +579,12 @@ where { // FK from repo_objects.oid → objects.oid enforces that inclusion // implies existence, so we only need the inclusion check. - let row: Option<(i32,)> = sqlx::query_as( - "SELECT 1 FROM repo_objects WHERE repo_did = $1 AND oid = $2 LIMIT 1", - ) - .bind(repo_did) - .bind(id.as_bytes()) - .fetch_optional(executor) - .await?; + let row: Option<(i32,)> = + sqlx::query_as("SELECT 1 FROM repo_objects WHERE repo_did = $1 AND oid = $2 LIMIT 1") + .bind(repo_did) + .bind(id.as_bytes()) + .fetch_optional(executor) + .await?; Ok(row.is_some()) } @@ -887,19 +876,15 @@ async fn apply_ref_edits_in_tx( impl RefStore for PgBackend { type Error = PgError; - async fn try_find_ref( - &self, - name: &FullNameRef, - ) -> Result, Self::Error> { - let row: Option<(i16, Option>, Option, Option>)> = - sqlx::query_as( - "SELECT kind, target_oid, target_name, peeled FROM refs \ + async fn try_find_ref(&self, name: &FullNameRef) -> Result, Self::Error> { + let row: Option<(i16, Option>, Option, Option>)> = sqlx::query_as( + "SELECT kind, target_oid, target_name, peeled FROM refs \ WHERE repo_did = $1 AND name = $2", - ) - .bind(&self.repo_did) - .bind(name_to_string(name)) - .fetch_optional(&self.pool) - .await?; + ) + .bind(&self.repo_did) + .bind(name_to_string(name)) + .fetch_optional(&self.pool) + .await?; match row { None => Ok(None), @@ -922,14 +907,19 @@ impl RefStore for PgBackend { &self, prefix: Option<&PartialNameRef>, ) -> Result, Self::Error> { - let rows: Vec<(String, i16, Option>, Option, Option>)> = - sqlx::query_as( - "SELECT name, kind, target_oid, target_name, peeled FROM refs \ + let rows: Vec<( + String, + i16, + Option>, + Option, + Option>, + )> = sqlx::query_as( + "SELECT name, kind, target_oid, target_name, peeled FROM refs \ WHERE repo_did = $1", - ) - .bind(&self.repo_did) - .fetch_all(&self.pool) - .await?; + ) + .bind(&self.repo_did) + .fetch_all(&self.pool) + .await?; let prefix_bytes = prefix.map(|p| p.as_bstr().to_vec()); let mut out = Vec::with_capacity(rows.len()); @@ -1006,13 +996,57 @@ impl PushTx { /// Apply a batch of ref edits with the same CAS / lock-ordering /// semantics as [`RefStore::commit_ref_transaction`], but inside the /// caller-owned tx (no commit here). - pub async fn apply_ref_edits( - &mut self, - edits: Vec, - ) -> Result, PgError> { + pub async fn apply_ref_edits(&mut self, edits: Vec) -> Result, PgError> { apply_ref_edits_in_tx(&mut self.tx, &self.repo_did, edits).await } + /// Look up a ref by name inside the push tx. Reads see prior edits + /// in the same push (read-your-writes), so meta computation that + /// runs after `apply_ref_edits` observes the post-update state. + /// + /// Used by the post-receive event emitter to read HEAD when + /// deciding `meta.isDefaultRef`. Returns `None` if the ref doesn't + /// exist. + pub async fn find_ref_target(&mut self, name: &FullNameRef) -> Result, PgError> { + let row: Option<(i16, Option>, Option)> = sqlx::query_as( + "SELECT kind, target_oid, target_name FROM refs \ + WHERE repo_did = $1 AND name = $2", + ) + .bind(&self.repo_did) + .bind(name_to_string(name)) + .fetch_optional(&mut *self.tx) + .await?; + row.map(|(kind, target_oid, target_name)| { + target_from_row(RefRow { + kind, + target_oid, + target_name, + }) + }) + .transpose() + } + + /// Resolve HEAD to the branch short-name it points at, e.g. + /// `"main"` for `HEAD -> refs/heads/main`. Returns `None` if HEAD + /// is unset, detached (points at an OID directly), or points at + /// something other than `refs/heads/*`. + /// + /// The post-receive emitter uses this to compute + /// `meta.isDefaultRef` for each pushed ref. + pub async fn find_default_branch(&mut self) -> Result, PgError> { + let head: FullName = BString::from(b"HEAD".to_vec()) + .try_into() + .expect("HEAD is a valid ref name"); + let Some(target) = self.find_ref_target(head.as_ref()).await? else { + return Ok(None); + }; + let Target::Symbolic(full) = target else { + return Ok(None); + }; + let name = full.as_bstr().to_string(); + Ok(name.strip_prefix("refs/heads/").map(|s| s.to_string())) + } + /// Walk the reachable object graph from each tip and confirm every /// referenced oid is present in `objects`. Returns the first missing /// oid as a `MissingObject` error; otherwise `Ok(())`. @@ -1024,8 +1058,7 @@ impl PushTx { let mut visited: HashSet = HashSet::new(); // (oid_to_check, oid_that_referenced_it). For tips the two are equal. - let mut worklist: Vec<(ObjectId, ObjectId)> = - tips.iter().map(|t| (*t, *t)).collect(); + let mut worklist: Vec<(ObjectId, ObjectId)> = tips.iter().map(|t| (*t, *t)).collect(); while let Some((oid, referenced_from)) = worklist.pop() { if !visited.insert(oid) { @@ -1058,8 +1091,7 @@ impl PushTx { } } Kind::Tree => { - for entry in - gix_object::TreeRefIter::from_bytes(&object.data, self.object_hash) + for entry in gix_object::TreeRefIter::from_bytes(&object.data, self.object_hash) { let entry = entry.map_err(|source| PgError::DecodeObject { kind: Kind::Tree, @@ -1070,8 +1102,7 @@ impl PushTx { } } Kind::Tag => { - let iter = - gix_object::TagRefIter::from_bytes(&object.data, self.object_hash); + let iter = gix_object::TagRefIter::from_bytes(&object.data, self.object_hash); let target_id = iter.target_id().map_err(|source| PgError::DecodeObject { kind: Kind::Tag, oid, diff --git a/mnemosyne_protocol/Cargo.toml b/mnemosyne_protocol/Cargo.toml index 9d130e5..2b06a8d 100644 --- a/mnemosyne_protocol/Cargo.toml +++ b/mnemosyne_protocol/Cargo.toml @@ -32,8 +32,12 @@ http-body-util.workspace = true anyhow.workspace = true serde = { version = "1", features = ["derive"] } +serde_json = "1" tempfile.workspace = true +# TID rkey generation for emitted records. +jacquard-common.workspace = true + [dev-dependencies] tokio = { workspace = true, features = ["rt-multi-thread", "macros"] } diff --git a/mnemosyne_protocol/src/events_stream.rs b/mnemosyne_protocol/src/events_stream.rs new file mode 100644 index 0000000..3eb2390 --- /dev/null +++ b/mnemosyne_protocol/src/events_stream.rs @@ -0,0 +1,260 @@ +//! `GET /events` — knot firehose WebSocket. +//! +//! Streams every committed row of the `events` table (currently +//! `sh.tangled.git.refUpdate` records) to subscribers in `seq` order. +//! Used by the Tangled appview to learn about activity on this knot +//! without polling per-repo XRPCs. +//! +//! ## Protocol +//! +//! - **Query param `cursor` (optional)** — last seen seq. The server +//! replays every event with `seq > cursor` in ascending order +//! before going live. +//! - **Message format** — one JSON object per WebSocket text frame: +//! +//! ```json +//! { "seq": 42, +//! "rkey": "3kx...", +//! "nsid": "sh.tangled.git.refUpdate", +//! "repo": "did:plc:...", // optional +//! "event": { /* lexicon record body */ }, +//! "created": 1716916873000000000 } +//! ``` +//! +//! `seq` is the cursor token; clients persist it and pass on +//! reconnect. `created` is nanoseconds-since-epoch to match +//! knotserver's `/events` shape. +//! +//! - **Backfill** — drained in pages of [`PAGE_SIZE`] until exhausted. +//! No artificial cap; if a subscriber resumes from a stale cursor +//! on a busy knot we send everything they missed before going live. +//! - **Live tail** — server subscribes to `LISTEN mnemosyne_events` +//! and re-polls on every notification, advancing the local cursor. +//! 30-second pings keep idle connections alive. +//! - **Backpressure** — if the websocket send fails (slow client, +//! dropped connection), the loop exits and the connection closes. +//! Subscribers reconnect with their last persisted cursor; no +//! server-side queueing. + +use std::time::Duration; + +use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade}; +use axum::extract::{Query, State}; +use axum::response::Response; +use mnemosyne_postgres::sqlx::PgPool; +use mnemosyne_postgres::sqlx::postgres::PgListener; +use mnemosyne_postgres::{EVENTS_NOTIFY_CHANNEL, EventRow, fetch_events_since}; +use serde::Deserialize; +use serde_json::json; + +/// Max rows fetched per `fetch_events_since` call. Tunable; small +/// enough that each frame batch fits comfortably in memory, large +/// enough that a fresh subscriber draining a long backlog doesn't +/// pay 100k round-trips. +pub const PAGE_SIZE: i64 = 100; + +/// Idle keep-alive interval. If no events have flowed for this long, +/// the server sends a WebSocket ping so middleboxes don't reap the +/// connection. +pub const KEEPALIVE_INTERVAL: Duration = Duration::from_secs(30); + +/// Query string for `GET /events`. Only one field — but defined as a +/// struct so axum's `Query` extractor gives us clean failure on bad +/// input (e.g. `?cursor=abc`). +#[derive(Debug, Deserialize)] +pub struct EventsQuery { + /// Last seen seq. Resume the stream from immediately after this + /// value. Defaults to 0 (full history) when absent. + #[serde(default)] + pub cursor: i64, +} + +/// State injected into [`handler`]. Holds the pool the WS task will +/// use both for backfill `SELECT`s and for opening the dedicated +/// `LISTEN` connection. +#[derive(Clone)] +pub struct EventsStreamState { + pub pool: PgPool, +} + +#[tracing::instrument( + name = "events.subscribe", + skip(ws, state), + fields( + events.cursor_in = query.cursor, + events.delivered = tracing::field::Empty, + events.exit_reason = tracing::field::Empty, + ), +)] +pub async fn handler( + State(state): State, + Query(query): Query, + ws: WebSocketUpgrade, +) -> Response { + let cursor = query.cursor; + ws.on_upgrade(move |socket| stream_loop(state.pool, socket, cursor)) +} + +/// Drive the WS stream: +/// 1. Open `LISTEN` connection first, so any commit landing during +/// backfill produces a notification we'll observe on the live +/// branch (no events lost in the gap between backfill drain and +/// `LISTEN`). +/// 2. Drain backfill from `cursor` until exhausted. +/// 3. Loop: wait for notification or keep-alive timeout; drain +/// again; ping on timeout. +async fn stream_loop(pool: PgPool, mut socket: WebSocket, mut cursor: i64) { + let span = tracing::Span::current(); + let mut delivered: u64 = 0; + + // Listener first — see step 1 in the function doc. + let mut listener = match PgListener::connect_with(&pool).await { + Ok(l) => l, + Err(e) => { + span.record("events.exit_reason", "listener_connect_failed"); + tracing::warn!(error = %e, "events: failed to open LISTEN connection"); + let _ = socket + .send(Message::Close(Some(axum::extract::ws::CloseFrame { + code: 1011, // internal error + reason: "listener connect failed".into(), + }))) + .await; + return; + } + }; + if let Err(e) = listener.listen(EVENTS_NOTIFY_CHANNEL).await { + span.record("events.exit_reason", "listen_failed"); + tracing::warn!(error = %e, "events: failed to issue LISTEN"); + let _ = socket + .send(Message::Close(Some(axum::extract::ws::CloseFrame { + code: 1011, + reason: "LISTEN failed".into(), + }))) + .await; + return; + } + + // Step 2: backfill drain. + if let Err(reason) = drain(&pool, &mut socket, &mut cursor, &mut delivered).await { + span.record("events.delivered", delivered); + span.record("events.exit_reason", reason); + return; + } + + // Step 3: live tail. + loop { + let recv = tokio::time::timeout(KEEPALIVE_INTERVAL, listener.recv()); + let next = tokio::select! { + // Client-initiated close / ping arrives here. + client = socket.recv() => { + match client { + Some(Ok(Message::Close(_))) | None => { + span.record("events.delivered", delivered); + span.record("events.exit_reason", "client_closed"); + return; + } + Some(Ok(Message::Ping(p))) => { + // Axum auto-pongs by default but explicit is fine. + if socket.send(Message::Pong(p)).await.is_err() { + span.record("events.delivered", delivered); + span.record("events.exit_reason", "send_failed"); + return; + } + continue; + } + Some(Ok(_)) => continue, + Some(Err(e)) => { + span.record("events.delivered", delivered); + span.record("events.exit_reason", "recv_failed"); + tracing::debug!(error = %e, "events: ws recv error"); + return; + } + } + } + notify = recv => notify, + }; + + match next { + Ok(Ok(_notification)) => { + if let Err(reason) = drain(&pool, &mut socket, &mut cursor, &mut delivered).await { + span.record("events.delivered", delivered); + span.record("events.exit_reason", reason); + return; + } + } + Ok(Err(e)) => { + span.record("events.delivered", delivered); + span.record("events.exit_reason", "listener_recv_failed"); + tracing::warn!(error = %e, "events: LISTEN connection died"); + return; + } + Err(_timeout) => { + if socket.send(Message::Ping(Vec::new().into())).await.is_err() { + span.record("events.delivered", delivered); + span.record("events.exit_reason", "keepalive_send_failed"); + return; + } + } + } + } +} + +/// Drain `events` rows newer than `cursor` until exhausted, sending +/// each as a single text frame. Advances `cursor` to the last +/// seq sent. Returns the exit reason on send failure. +async fn drain( + pool: &PgPool, + socket: &mut WebSocket, + cursor: &mut i64, + delivered: &mut u64, +) -> Result<(), &'static str> { + loop { + let batch = match fetch_events_since(pool, *cursor, PAGE_SIZE).await { + Ok(b) => b, + Err(e) => { + tracing::warn!(error = %e, "events: fetch_events_since failed"); + return Err("fetch_failed"); + } + }; + if batch.is_empty() { + return Ok(()); + } + let batch_len = batch.len(); + let last_seq = batch.last().map_or(*cursor, |e| e.seq); + for row in batch { + let payload = wire_format(row); + let frame = match serde_json::to_string(&payload) { + Ok(s) => s, + Err(e) => { + tracing::warn!(error = %e, "events: serialize failed"); + return Err("serialize_failed"); + } + }; + if socket.send(Message::Text(frame.into())).await.is_err() { + return Err("send_failed"); + } + *delivered += 1; + } + *cursor = last_seq; + // Cast is safe: batch_len <= PAGE_SIZE (an i64), well under i64::MAX. + if i64::try_from(batch_len).unwrap_or(i64::MAX) < PAGE_SIZE { + return Ok(()); + } + } +} + +/// Wire shape of one event row. Matches knotserver's `/events` +/// frame fields plus our own `seq` cursor token. `repo` is included +/// (vs knotserver's omission) because mnemosyne always carries +/// `repo_did` on a `refUpdate` row and surfacing it lets the +/// appview skip a JSON decode round-trip just to learn the repo. +fn wire_format(row: EventRow) -> serde_json::Value { + json!({ + "seq": row.seq, + "rkey": row.rkey, + "nsid": row.nsid, + "repo": row.repo_did, + "event": row.payload, + "created": row.created_nanos, + }) +} diff --git a/mnemosyne_protocol/src/lib.rs b/mnemosyne_protocol/src/lib.rs index 7039ee5..4a9eb18 100644 --- a/mnemosyne_protocol/src/lib.rs +++ b/mnemosyne_protocol/src/lib.rs @@ -9,11 +9,13 @@ pub mod error; pub mod events; +pub mod events_stream; pub mod find_adapter; pub mod info_refs; pub mod pack; pub mod pkt; pub mod receive_pack; +pub mod ref_update; pub mod upload_pack; pub use error::ProtocolError; diff --git a/mnemosyne_protocol/src/receive_pack.rs b/mnemosyne_protocol/src/receive_pack.rs index a145c16..f3380e8 100644 --- a/mnemosyne_protocol/src/receive_pack.rs +++ b/mnemosyne_protocol/src/receive_pack.rs @@ -40,7 +40,7 @@ use std::sync::{Arc, Mutex}; use axum::body::{Body, Bytes}; use axum::extract::State; -use axum::http::{header, StatusCode}; +use axum::http::{StatusCode, header}; use axum::response::Response; use bstr::ByteSlice; use gix_hash::ObjectId; @@ -51,7 +51,7 @@ use mnemosyne_postgres::{PgError, PushTx}; use tokio::runtime::Handle; use crate::find_adapter::PushTxFind; -use crate::{pkt, ProtocolError, RepoState}; +use crate::{ProtocolError, RepoState, pkt}; // Spare 5 bytes for the 4-hex length prefix + 1-byte channel marker. const SIDE_BAND_MAX_DATA: usize = pkt::MAX_DATA_LEN - 1; @@ -131,6 +131,7 @@ struct ReceivePackRequest { request.command_count = tracing::field::Empty, request.pack_bytes = tracing::field::Empty, request.delete_only = tracing::field::Empty, + emit.seq_max = tracing::field::Empty, ), )] pub async fn handler( @@ -163,18 +164,16 @@ pub async fn handler( // PDS and ingested into the `collaborators` table. Unauthenticated // callers (e.g. raw HTTP that bypassed our SSH layer) get denied // because `authenticated_did` is `None`. - let owner_did = mnemosyne_postgres::lookup_repo_owner( - state.backend.pool(), - state.backend.repo_did(), - ) - .await - .map_err(ProtocolError::Storage)? - .ok_or_else(|| { - ProtocolError::BadRequest(format!( - "repo {} not registered", - state.backend.repo_did() - )) - })?; + let owner_did = + mnemosyne_postgres::lookup_repo_owner(state.backend.pool(), state.backend.repo_did()) + .await + .map_err(ProtocolError::Storage)? + .ok_or_else(|| { + ProtocolError::BadRequest(format!( + "repo {} not registered", + state.backend.repo_did() + )) + })?; match state.authenticated_did.as_deref() { Some(did) if did == owner_did => { span.record("auth.role", "owner"); @@ -246,14 +245,44 @@ pub async fn handler( }; // Step 3: commit if everything passed, otherwise drop the tx. - let everything_ok = unpack_status.is_ok() - && command_outcomes.iter().all(|(_, r)| r.is_ok()); + let everything_ok = unpack_status.is_ok() && command_outcomes.iter().all(|(_, r)| r.is_ok()); if everything_ok { + // Persist `sh.tangled.git.refUpdate` rows inside the same + // push tx so the events table is atomic with the ref edits. + // `committer_did` is the authenticated principal (verified + // owner-or-collaborator above); the `authenticated_did: + // None` branch was rejected earlier. The unwrap is therefore + // safe — if we reach here we have a DID. + let committer_did = state + .authenticated_did + .as_deref() + .expect("authenticated_did populated post-authz check"); + let updates: Vec = req + .commands + .iter() + .map(|c| crate::ref_update::RefUpdateInput { + ref_name: c.ref_name.clone(), + old_oid: c.old, + new_oid: c.new, + }) + .collect(); + let emitted = crate::ref_update::emit_ref_updates( + &mut push_tx, + state.backend.pool(), + state.backend.repo_did(), + &owner_did, + committer_did, + &updates, + ) + .await?; + span.record("emit.seq_max", emitted.last().map(|e| e.seq).unwrap_or(0)); + push_tx.commit().await.map_err(ProtocolError::Storage)?; - // Fire the post-receive event after the commit lands, so - // subscribers can react to a known-durable state. Drops silently - // if no broadcaster is configured on RepoState. + // Fire the in-process post-receive event after the commit + // lands. Same-replica subscribers (cache warmers, etc.) react + // here; cross-replica subscribers use the durable `events` + // table via `pg_notify` from `emit_ref_updates`. if let Some(events) = &state.events { let refs_updated = req .commands @@ -283,7 +312,10 @@ pub async fn handler( Response::builder() .status(StatusCode::OK) - .header(header::CONTENT_TYPE, "application/x-git-receive-pack-result") + .header( + header::CONTENT_TYPE, + "application/x-git-receive-pack-result", + ) .header(header::CACHE_CONTROL, "no-cache") .body(Body::from(body_bytes)) .map_err(|e| ProtocolError::Internal(e.to_string())) @@ -357,20 +389,14 @@ fn parse_request( .map_err(|e| ProtocolError::BadRequest(format!("old oid: {e}")))?; let new = ObjectId::from_hex(parts[1]) .map_err(|e| ProtocolError::BadRequest(format!("new oid: {e}")))?; - let ref_name: FullName = bstr::BString::from(parts[2].to_vec()) - .try_into() - .map_err(|_| { - ProtocolError::BadRequest(format!( - "invalid ref name: {:?}", - parts[2].as_bstr() - )) - })?; + let ref_name: FullName = + bstr::BString::from(parts[2].to_vec()) + .try_into() + .map_err(|_| { + ProtocolError::BadRequest(format!("invalid ref name: {:?}", parts[2].as_bstr())) + })?; - commands.push(UpdateCommand { - old, - new, - ref_name, - }); + commands.push(UpdateCommand { old, new, ref_name }); if is_first { if let Some(caps) = caps_part { @@ -487,16 +513,16 @@ fn unpack_into_tx( for idx in 0..bundle.index.num_objects() { buf.clear(); - let (kind, bytes) = match bundle.get_object_by_index(idx, &mut buf, &mut inflate, &mut cache) - { - Ok((data, _location)) => (data.kind, data.data.to_vec()), - Err(e) => { - return ( - push_tx, - Err(ProtocolError::PackDecode(format!("decode entry: {e}"))), - ); - } - }; + let (kind, bytes) = + match bundle.get_object_by_index(idx, &mut buf, &mut inflate, &mut cache) { + Ok((data, _location)) => (data.kind, data.data.to_vec()), + Err(e) => { + return ( + push_tx, + Err(ProtocolError::PackDecode(format!("decode entry: {e}"))), + ); + } + }; match handle.block_on(push_tx.insert_object(kind, &bytes)) { Ok(_oid) => {} diff --git a/mnemosyne_protocol/src/ref_update.rs b/mnemosyne_protocol/src/ref_update.rs new file mode 100644 index 0000000..79a8ebd --- /dev/null +++ b/mnemosyne_protocol/src/ref_update.rs @@ -0,0 +1,372 @@ +//! `sh.tangled.git.refUpdate` event construction + emission. +//! +//! Called from [`crate::receive_pack`] after pack ingestion and ref +//! edits succeed, but before the push transaction commits — so the +//! refUpdate row lands atomically with the ref state it describes. +//! +//! ## Wire shape +//! +//! Each event is the JSON-serialized lexicon record: +//! +//! ```json +//! { "ref": "refs/heads/main", +//! "committerDid": "did:plc:abc...", +//! "ownerDid": "did:plc:owner...", +//! "repo": "did:plc:repo...", +//! "oldSha": "<40 hex>", +//! "newSha": "<40 hex>", +//! "meta": { +//! "isDefaultRef": true, +//! "commitCount": { "byEmail": [{ "email": "...", "count": 1 }] }, +//! "langBreakdown": { "inputs": [...] } // optional, omitted on cache miss +//! } } +//! ``` +//! +//! Rkey is a freshly-minted TID, courtesy of `jacquard_common::types::tid::Tid`. +//! +//! ## Meta computation +//! +//! - `isDefaultRef` reads HEAD via the push tx's read-your-writes +//! view, so it observes any HEAD update that landed in the same push. +//! - `commitCount.byEmail` walks the new ref's history bounded at +//! [`COMMIT_WALK_LIMIT`] commits, stopping at `oldSha` when that +//! exists. Matches knotserver's `--max-count=100` behavior. +//! - `langBreakdown` is populated only on `language_stats_cache` hit +//! for the new commit's tree. We never block push commit to compute +//! it; the receive-pack path spawns a cache warmer post-commit so +//! subsequent pushes / queries see it ready. +//! +//! ## Atomicity +//! +//! Event INSERT happens inside the push tx and takes the +//! `EVENTS_ADVISORY_LOCK_ID` lock immediately before INSERTing, then +//! commits. A crash between unpacking and commit aborts both the ref +//! edits and the event row — no half-states reach subscribers. + +use std::collections::{HashMap, HashSet}; + +use bytes::Bytes; +use gix_hash::{Kind as HashKind, ObjectId}; +use gix_object::{CommitRef, Kind as ObjectKind}; +use gix_ref::FullName; +use mnemosyne_postgres::sqlx::PgPool; +use mnemosyne_postgres::{PgError, PushTx, fetch_language_stats, insert_event_in_tx}; +use serde_json::{Value, json}; + +use crate::ProtocolError; + +/// Minimal async object-source abstraction over which the commit +/// walker is generic. `PushTx` and any mock test source implement +/// it; the walker doesn't care. +/// +/// Kept narrow on purpose — `ObjectStore` from `mnemosyne_git` is +/// the right home long-term, but it's `async_trait`-based and brings +/// a heavier interface than the walker needs. We can collapse this +/// into `ObjectStore` later if a second caller emerges. +#[async_trait::async_trait] +pub trait CommitSource { + /// Look up a commit by id. Returns `None` for missing OIDs (the + /// walker treats those as walked-off-the-end leaves, not errors). + /// Non-commit objects also return `None` — the caller has + /// already enforced that walked OIDs are commits. + async fn try_find_commit(&mut self, oid: ObjectId) -> Result, PgError>; +} + +#[async_trait::async_trait] +impl CommitSource for &mut PushTx { + async fn try_find_commit(&mut self, oid: ObjectId) -> Result, PgError> { + let Some(obj) = self.find_object(oid.as_ref()).await? else { + return Ok(None); + }; + if obj.kind != ObjectKind::Commit { + return Ok(None); + } + Ok(Some(obj.data)) + } +} + +/// Hard cap on the *new-side* commit walk per refUpdate — i.e. how +/// many commits we'll attribute to author emails in the resulting +/// meta. Matches knotserver's `git rev-list --max-count=100`. +pub const COMMIT_WALK_LIMIT: usize = 100; + +/// Hard cap on the *old-side* ancestor walk used to mark +/// uninteresting commits for proper `new..old` set-difference +/// semantics. Set generously (10× new-side) so very deep histories +/// still produce a well-formed `commitCount` when the divergence +/// itself is small. If old's ancestor chain exceeds this, the walker +/// may over-count by including commits that *are* reachable from old +/// — acceptable degradation; the alternative (silent error) is +/// worse for an event stream. +pub const OLD_ANCESTOR_WALK_LIMIT: usize = COMMIT_WALK_LIMIT * 10; + +/// NSID of the record we emit. Exposed as a constant so the WS +/// handler / tests can filter without restringing. +pub const REF_UPDATE_NSID: &str = "sh.tangled.git.refUpdate"; + +/// One emitted event with its allocated sequence number. Returned by +/// [`emit_ref_updates`] so the caller can log / instrument. +#[derive(Debug)] +pub struct EmittedRefUpdate { + pub seq: i64, + pub rkey: String, + pub ref_name: String, +} + +/// Input to [`emit_ref_updates`]: one update that just landed in the +/// push tx. The receive-pack handler builds this from its +/// `UpdateCommand`s after they've been validated against connectivity +/// and applied. +#[derive(Debug, Clone)] +pub struct RefUpdateInput { + pub ref_name: FullName, + pub old_oid: ObjectId, + pub new_oid: ObjectId, +} + +/// Build and persist a `sh.tangled.git.refUpdate` row for every input. +/// +/// Runs inside `tx` (the push transaction). Each event lands atomically +/// with the rest of the push — drop the tx and the events disappear +/// along with the ref edits. +/// +/// `owner_did` is the repo owner DID (looked up by the caller from the +/// `repos` table); `committer_did` is the authenticated principal that +/// pushed (already validated as owner-or-collaborator). +/// +/// Errors are wrapped in [`ProtocolError::Storage`] / `Internal` so +/// they unwind through the receive_pack handler the same way pack / +/// ref-edit errors do. +#[tracing::instrument( + name = "ref_update.emit", + skip(tx, pool, updates), + fields( + repo.did = %repo_did, + emit.count = updates.len(), + ), +)] +pub async fn emit_ref_updates( + tx: &mut PushTx, + pool: &PgPool, + repo_did: &str, + owner_did: &str, + committer_did: &str, + updates: &[RefUpdateInput], +) -> Result, ProtocolError> { + let object_hash = tx.object_hash(); + let default_branch = tx + .find_default_branch() + .await + .map_err(ProtocolError::Storage)?; + + let mut emitted = Vec::with_capacity(updates.len()); + for input in updates { + let payload = build_payload( + tx, + pool, + repo_did, + owner_did, + committer_did, + default_branch.as_deref(), + input, + object_hash, + ) + .await?; + + let rkey = jacquard_common::types::tid::Tid::now_0() + .as_str() + .to_string(); + let seq = insert_event_in_tx(tx.tx(), &rkey, REF_UPDATE_NSID, Some(repo_did), &payload) + .await + .map_err(ProtocolError::Storage)?; + + emitted.push(EmittedRefUpdate { + seq, + rkey, + ref_name: input.ref_name.as_bstr().to_string(), + }); + } + + Ok(emitted) +} + +async fn build_payload( + tx: &mut PushTx, + pool: &PgPool, + repo_did: &str, + owner_did: &str, + committer_did: &str, + default_branch: Option<&str>, + input: &RefUpdateInput, + object_hash: HashKind, +) -> Result { + let ref_name = input.ref_name.as_bstr().to_string(); + let is_default_ref = default_branch.is_some_and(|db| ref_name == format!("refs/heads/{db}")); + + let counts = count_commits_by_email(&mut &mut *tx, input.old_oid, input.new_oid, object_hash) + .await + .map_err(ProtocolError::Storage)?; + + let by_email: Vec = counts + .into_iter() + .map(|(email, count)| json!({ "email": email, "count": count })) + .collect(); + + let mut meta = json!({ + "isDefaultRef": is_default_ref, + "commitCount": { "byEmail": by_email }, + }); + + // Optional langBreakdown — read from the language_stats_cache + // outside the push tx. The cache is keyed by `(repo_did, tree_oid)` + // so a stale-read is impossible (any prior cached value is by + // definition for a different commit's tree). On miss we omit the + // field; the post-commit warmer task fills it for next time. + if !input.new_oid.is_null() { + if let Some(tree_oid) = head_tree_oid(tx, input.new_oid, object_hash).await? { + if let Some(stats) = fetch_language_stats(pool, repo_did, tree_oid.as_bytes()) + .await + .map_err(ProtocolError::Storage)? + { + let inputs: Vec = stats + .into_iter() + .map(|(lang, size)| json!({ "lang": lang, "size": size })) + .collect(); + meta["langBreakdown"] = json!({ "inputs": inputs }); + } + } + } + + Ok(json!({ + "ref": ref_name, + "committerDid": committer_did, + "ownerDid": owner_did, + "repo": repo_did, + "oldSha": input.old_oid.to_string(), + "newSha": input.new_oid.to_string(), + "meta": meta, + })) +} + +/// Resolve the tree OID for a commit. Returns `None` for non-commit +/// objects (shouldn't happen on a valid ref update but we don't panic +/// for malformed input). +async fn head_tree_oid( + tx: &mut PushTx, + commit_oid: ObjectId, + object_hash: HashKind, +) -> Result, ProtocolError> { + let Some(obj) = tx + .find_object(commit_oid.as_ref()) + .await + .map_err(ProtocolError::Storage)? + else { + return Ok(None); + }; + if obj.kind != ObjectKind::Commit { + return Ok(None); + } + let commit = CommitRef::from_bytes(&obj.data, object_hash) + .map_err(|e| ProtocolError::Internal(format!("decode commit {commit_oid}: {e}")))?; + Ok(Some(commit.tree())) +} + +/// Compute the per-author-email count of commits reachable from +/// `new_oid` but not from `old_oid` — i.e. the post-push equivalent +/// of `git rev-list old..new`. Returns at most [`COMMIT_WALK_LIMIT`] +/// attributions. +/// +/// - For pure creates (`old_oid` is zero), walks up to the limit +/// from `new_oid` without any exclusion. Note this over-counts +/// versus knotserver, which excludes commits already reachable +/// from other branches; v1 deliberately punts that subtlety. +/// - For deletes (`new_oid` is zero), returns an empty map — no new +/// history to attribute. +/// - For updates, performs a two-phase walk: +/// 1. **Uninteresting marking** — BFS from `old_oid` populating a +/// set of "already-knew-about-this" oids, bounded by +/// [`OLD_ANCESTOR_WALK_LIMIT`]. +/// 2. **New-side count** — BFS from `new_oid`, skipping any oid in +/// the uninteresting set, counting up to [`COMMIT_WALK_LIMIT`] +/// attributions. Because uninteresting oids short-circuit +/// descent into their ancestors, this correctly produces the +/// set-difference (matches `git rev-list old..new`). +/// +/// Generic over [`CommitSource`] so proptests can drive it against +/// an in-memory DAG without spinning up Postgres. In production the +/// source is `&mut &mut PushTx`, providing read-your-writes +/// visibility into objects unpacked earlier in the same push. +pub async fn count_commits_by_email( + source: &mut S, + old_oid: ObjectId, + new_oid: ObjectId, + object_hash: HashKind, +) -> Result, PgError> { + let mut counts: HashMap = HashMap::new(); + if new_oid.is_null() { + return Ok(counts); + } + + // Phase 1: mark `old_oid`'s ancestor closure as uninteresting. + // Bounded — see OLD_ANCESTOR_WALK_LIMIT for the degradation rule. + let mut uninteresting: HashSet = HashSet::new(); + if !old_oid.is_null() { + let mut stack: Vec = vec![old_oid]; + while let Some(oid) = stack.pop() { + if uninteresting.len() >= OLD_ANCESTOR_WALK_LIMIT { + break; + } + if !uninteresting.insert(oid) { + continue; + } + let Some(bytes) = source.try_find_commit(oid).await? else { + continue; + }; + let Ok(commit_ref) = CommitRef::from_bytes(&bytes, object_hash) else { + continue; + }; + for parent in commit_ref.parents() { + stack.push(parent); + } + } + } + + // Phase 2: count commits on the new side, skipping the + // uninteresting closure. Iteration order matches knotserver's + // pop-from-back behavior — properties don't depend on it, but + // matching makes debugging easier. + let mut visited: HashSet = HashSet::new(); + let mut worklist: Vec = vec![new_oid]; + let mut walked = 0usize; + + while let Some(oid) = worklist.pop() { + if walked >= COMMIT_WALK_LIMIT { + break; + } + if uninteresting.contains(&oid) { + // Don't descend into the uninteresting closure — its + // ancestors are by definition also uninteresting. + continue; + } + if !visited.insert(oid) { + continue; + } + let Some(bytes) = source.try_find_commit(oid).await? else { + continue; + }; + let Ok(commit_ref) = CommitRef::from_bytes(&bytes, object_hash) else { + continue; + }; + if let Ok(author) = commit_ref.author() { + let email = author.email.to_string(); + *counts.entry(email).or_insert(0) += 1; + } + walked += 1; + + for parent in commit_ref.parents() { + worklist.push(parent); + } + } + + Ok(counts) +} diff --git a/mnemosyne_tests/Cargo.toml b/mnemosyne_tests/Cargo.toml index 7358ab3..fa7381b 100644 --- a/mnemosyne_tests/Cargo.toml +++ b/mnemosyne_tests/Cargo.toml @@ -47,12 +47,14 @@ anyhow.workspace = true axum.workspace = true base64.workspace = true bstr.workspace = true +bytes.workspace = true # Process-exit destructors so we can `stop()` shared testcontainers # held in `OnceCell` statics — those statics don't run Drop at exit, # so the underlying testcontainers-rs cleanup never fires without # this hook. ctor = "0.4" flate2.workspace = true +futures.workspace = true http-body-util.workspace = true reqwest = { version = "0.12", default-features = false, features = ["rustls-tls"] } proptest.workspace = true @@ -65,6 +67,10 @@ sqlx.workspace = true tempfile.workspace = true testcontainers-modules.workspace = true tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net"] } +# WebSocket client for `/events` round-trip tests. Uses the same +# `tokio-tungstenite` line that powers the Jetstream ingester (see +# workspace deps), just on the receiving side. +tokio-tungstenite.workspace = true tower.workspace = true diff --git a/mnemosyne_tests/tests/integration/events_websocket.rs b/mnemosyne_tests/tests/integration/events_websocket.rs new file mode 100644 index 0000000..a02ddb9 --- /dev/null +++ b/mnemosyne_tests/tests/integration/events_websocket.rs @@ -0,0 +1,325 @@ +//! `/events` WebSocket round-trip integration tests. +//! +//! Boots the events router against a fresh Postgres tx, inserts +//! events via [`mnemosyne_postgres::insert_event_in_tx`], and asserts +//! a tokio-tungstenite subscriber receives them in the expected wire +//! shape. Verifies the LISTEN/NOTIFY path end-to-end and the +//! `?cursor=` resume semantics — both of which are too easy to break +//! and too hard to catch with unit-level coverage. +//! +//! Properties exercised: +//! +//! - **Backfill** — events inserted *before* the subscriber connects +//! are replayed in `seq` order. +//! - **Live tail** — events inserted *after* connection are pushed +//! via NOTIFY without polling. +//! - **Cursor resume** — `?cursor=K` skips everything with `seq <= K` +//! on backfill and (because new events have `seq > K`) sees only +//! strictly newer rows on the live path. +//! - **Wire shape** — each frame is JSON with the documented fields. + +#![allow(unreachable_pub)] + +use std::time::Duration; + +use axum::{Router, routing::get}; +use futures::StreamExt; +use mnemosyne_postgres::insert_event_in_tx; +use mnemosyne_protocol::events_stream::{self, EventsStreamState}; +use serde_json::{Value, json}; +use tokio::net::TcpListener; +use tokio_tungstenite::tungstenite::Message; + +async fn spawn_server(pool: sqlx::PgPool) -> (std::net::SocketAddr, tokio::task::AbortHandle) { + let app: Router = Router::new() + .route("/events", get(events_stream::handler)) + .with_state(EventsStreamState { pool }); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let handle = tokio::spawn(async move { + let _ = axum::serve(listener, app).await; + }); + (addr, handle.abort_handle()) +} + +async fn connect_ws( + addr: std::net::SocketAddr, + cursor: Option, +) -> tokio_tungstenite::WebSocketStream> { + let url = match cursor { + Some(c) => format!("ws://{addr}/events?cursor={c}"), + None => format!("ws://{addr}/events"), + }; + let (stream, _resp) = tokio_tungstenite::connect_async(&url) + .await + .expect("ws connect"); + stream +} + +/// Insert one synthetic event in a fresh tx and commit. Returns the +/// allocated seq. +async fn insert_one(pool: &sqlx::PgPool, repo_did: &str, payload: Value) -> i64 { + let mut tx = pool.begin().await.expect("begin"); + let rkey = format!( + "rk-{}", + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() + ); + let seq = insert_event_in_tx( + &mut tx, + &rkey, + "sh.tangled.git.refUpdate", + Some(repo_did), + &payload, + ) + .await + .expect("insert_event_in_tx"); + tx.commit().await.expect("commit"); + seq +} + +/// Pull the next text frame within `timeout`, parsing it as JSON. +/// Skips pings/pongs (they're never asserted in these tests). +async fn next_event( + stream: &mut tokio_tungstenite::WebSocketStream< + tokio_tungstenite::MaybeTlsStream, + >, + timeout: Duration, +) -> Option { + loop { + let msg = tokio::time::timeout(timeout, stream.next()).await.ok()??; + let frame = match msg { + Ok(Message::Text(t)) => t, + Ok(Message::Close(_)) | Err(_) => return None, + Ok(_) => continue, + }; + return serde_json::from_str(&frame).ok(); + } +} + +/// Like [`next_event`] but skips frames belonging to a different +/// `repo_did`. The shared `events` table interleaves rows from every +/// test in the suite; per-repo filtering is how we get isolation +/// without a per-test schema. +async fn next_event_for_repo( + stream: &mut tokio_tungstenite::WebSocketStream< + tokio_tungstenite::MaybeTlsStream, + >, + repo_did: &str, + timeout: Duration, +) -> Option { + let deadline = tokio::time::Instant::now() + timeout; + loop { + let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); + if remaining.is_zero() { + return None; + } + let frame = next_event(stream, remaining).await?; + if frame["repo"].as_str() == Some(repo_did) { + return Some(frame); + } + } +} + +/// Snapshot the current global `seq` ceiling. Tests pass this as +/// `cursor` so they skip events from earlier-suite test cases that +/// inserted into the shared table. +async fn starting_max(pool: &sqlx::PgPool) -> i64 { + sqlx::query_scalar("SELECT COALESCE(MAX(seq), 0) FROM events") + .fetch_one(pool) + .await + .unwrap() +} + +#[test] +fn backfill_drains_in_seq_order() { + crate::common::runtime().block_on(async { + let pool = crate::common::shared_pool().await; + let repo_did = crate::common::unique_repo_did(); + mnemosyne_postgres::ensure_repo( + &pool, + &repo_did, + "did:test:owner", + "events-backfill", + crate::common::PLACEHOLDER_RKEY, + ) + .await + .unwrap(); + let cursor_start = starting_max(&pool).await; + + // Pre-load 3 events before any subscriber connects. + let s1 = insert_one(&pool, &repo_did, json!({"ref": "refs/heads/a"})).await; + let s2 = insert_one(&pool, &repo_did, json!({"ref": "refs/heads/b"})).await; + let s3 = insert_one(&pool, &repo_did, json!({"ref": "refs/heads/c"})).await; + assert!( + s1 < s2 && s2 < s3, + "seqs must be monotonic: {s1} < {s2} < {s3}" + ); + + let (addr, abort) = spawn_server(pool.clone()).await; + let mut ws = connect_ws(addr, Some(cursor_start)).await; + + // Drain the three rows belonging to us, in seq order. + for expected_seq in [s1, s2, s3] { + let frame = next_event_for_repo(&mut ws, &repo_did, Duration::from_secs(5)) + .await + .expect("expected backfilled frame"); + assert_eq!( + frame["seq"].as_i64(), + Some(expected_seq), + "frame seq mismatch: {frame}", + ); + assert_eq!(frame["nsid"], "sh.tangled.git.refUpdate"); + assert_eq!(frame["repo"], repo_did); + assert!(frame["created"].is_i64(), "created should be i64 nanos"); + assert!(frame["event"].is_object()); + assert!(frame["rkey"].is_string()); + } + + let _ = ws.close(None).await; + abort.abort(); + }); +} + +#[test] +fn live_tail_pushes_new_events_via_notify() { + crate::common::runtime().block_on(async { + let pool = crate::common::shared_pool().await; + let repo_did = crate::common::unique_repo_did(); + mnemosyne_postgres::ensure_repo( + &pool, + &repo_did, + "did:test:owner", + "events-live", + crate::common::PLACEHOLDER_RKEY, + ) + .await + .unwrap(); + + let (addr, abort) = spawn_server(pool.clone()).await; + let cursor_start = starting_max(&pool).await; + let mut ws = connect_ws(addr, Some(cursor_start)).await; + + // Insert AFTER connect — the LISTEN/NOTIFY path must wake + // the handler and the new row must reach the client. + let s_live = insert_one(&pool, &repo_did, json!({"ref": "refs/heads/live"})).await; + assert!(s_live > cursor_start); + + let frame = next_event_for_repo(&mut ws, &repo_did, Duration::from_secs(5)) + .await + .expect("expected live-tail frame after NOTIFY"); + assert_eq!(frame["seq"].as_i64(), Some(s_live)); + assert_eq!(frame["event"]["ref"], "refs/heads/live"); + + let _ = ws.close(None).await; + abort.abort(); + }); +} + +#[test] +fn cursor_resume_skips_already_seen() { + crate::common::runtime().block_on(async { + let pool = crate::common::shared_pool().await; + let repo_did = crate::common::unique_repo_did(); + mnemosyne_postgres::ensure_repo( + &pool, + &repo_did, + "did:test:owner", + "events-resume", + crate::common::PLACEHOLDER_RKEY, + ) + .await + .unwrap(); + + let s1 = insert_one(&pool, &repo_did, json!({"step": 1})).await; + let s2 = insert_one(&pool, &repo_did, json!({"step": 2})).await; + + let (addr, abort) = spawn_server(pool.clone()).await; + // Resume from after s1 — we should see only s2 in backfill + // for OUR repo (other tests' events may sit between s1 and s2 + // but `next_event_for_repo` filters them out). + let mut ws = connect_ws(addr, Some(s1)).await; + let frame = next_event_for_repo(&mut ws, &repo_did, Duration::from_secs(5)) + .await + .expect("expected one frame after cursor"); + assert_eq!(frame["seq"].as_i64(), Some(s2)); + assert_eq!(frame["event"]["step"], 2); + + // No further frames for OUR repo before live tail begins. + let no_more = next_event_for_repo(&mut ws, &repo_did, Duration::from_millis(300)).await; + assert!(no_more.is_none(), "unexpected extra frame: {no_more:?}"); + + let _ = ws.close(None).await; + abort.abort(); + }); +} + +/// Verify the contract that concurrent writers don't break seq +/// monotonicity from a subscriber's perspective. Spawns N +/// concurrent inserters under the advisory-lock contract; the +/// subscribed client must see all of them in strictly increasing +/// seq order. +#[test] +fn concurrent_writers_yield_strictly_increasing_seq() { + crate::common::runtime().block_on(async { + let pool = crate::common::shared_pool().await; + let repo_did = crate::common::unique_repo_did(); + mnemosyne_postgres::ensure_repo( + &pool, + &repo_did, + "did:test:owner", + "events-concurrent", + crate::common::PLACEHOLDER_RKEY, + ) + .await + .unwrap(); + + let cursor_start = starting_max(&pool).await; + + let (addr, abort) = spawn_server(pool.clone()).await; + let mut ws = connect_ws(addr, Some(cursor_start)).await; + + const N: usize = 16; + // Fire N concurrent inserts. Each opens its own tx, so the + // advisory lock serializes their seq allocation against each + // other. + let mut handles = Vec::new(); + for i in 0..N { + let p = pool.clone(); + let rd = repo_did.clone(); + handles.push(tokio::spawn(async move { + insert_one(&p, &rd, json!({"writer": i})).await + })); + } + let mut expected_seqs: Vec = Vec::with_capacity(N); + for h in handles { + expected_seqs.push(h.await.unwrap()); + } + expected_seqs.sort(); + + let mut received: Vec = Vec::with_capacity(N); + while received.len() < N { + let frame = next_event_for_repo(&mut ws, &repo_did, Duration::from_secs(5)) + .await + .expect("expected concurrent frame"); + received.push(frame["seq"].as_i64().expect("seq is i64")); + } + + // The subscriber must see every committed seq. + assert_eq!( + received, expected_seqs, + "subscriber missed events or saw them out of order" + ); + // Strictly increasing — restating the property explicitly so + // a regression shows up as a more legible assertion than the + // set equality above. + for w in received.windows(2) { + assert!(w[0] < w[1], "non-monotonic seq pair: {} >= {}", w[0], w[1]); + } + + let _ = ws.close(None).await; + abort.abort(); + }); +} diff --git a/mnemosyne_tests/tests/integration/main.rs b/mnemosyne_tests/tests/integration/main.rs index 55d5d4a..df2ce2c 100644 --- a/mnemosyne_tests/tests/integration/main.rs +++ b/mnemosyne_tests/tests/integration/main.rs @@ -21,6 +21,7 @@ mod common; +mod events_websocket; mod ingest_dispatch; mod ingest_public_key; mod ingest_repo; @@ -33,12 +34,14 @@ mod protocol_pack; mod protocol_pkt; mod protocol_push_concurrent; mod protocol_push_roundtrip; +mod ref_update_walk; mod ssh_auth; mod ssh_push; mod ssh_push_parity; mod storage_dedup; mod stress; mod stress_concurrent_push; +mod stress_events; mod xrpc_browse; mod xrpc_endpoints; mod xrpc_history; diff --git a/mnemosyne_tests/tests/integration/ref_update_walk.rs b/mnemosyne_tests/tests/integration/ref_update_walk.rs new file mode 100644 index 0000000..913ea51 --- /dev/null +++ b/mnemosyne_tests/tests/integration/ref_update_walk.rs @@ -0,0 +1,329 @@ +//! Proptests for [`mnemosyne_protocol::ref_update::count_commits_by_email`]. +//! +//! The walker is the most subtle piece of slice-1's ref-update emission: +//! it traverses a synthetic commit DAG bounded by [`COMMIT_WALK_LIMIT`], +//! attributes each visited commit's author to a counts map, and stops at +//! a configurable `old_oid` boundary. The properties we care about: +//! +//! 1. **Termination** — the walk always ends; no cycles, no infinite parents. +//! 2. **Bounded** — the total visited count never exceeds the limit. +//! 3. **No leakage past `old_oid`** — when `old_oid` is reachable from +//! `new_oid`, no commit at-or-beyond `old_oid` contributes to the +//! counts. +//! 4. **Empty on delete** — when `new_oid` is null, the result is empty. +//! 5. **Full attribution** — the sum of counts equals the number of +//! walked commits (every visited commit's author was attributed). +//! +//! The walker is generic over [`CommitSource`], so we drive it against an +//! in-memory map of synthetic commit objects encoded with `gix_object`. +//! No Postgres in the loop. + +#![allow(unreachable_pub)] + +use std::collections::HashMap; + +use async_trait::async_trait; +use bytes::Bytes; +use gix_actor::Signature; +use gix_hash::{Kind as HashKind, ObjectId}; +use gix_object::{Commit, Kind, WriteTo}; +use mnemosyne_postgres::PgError; +use mnemosyne_protocol::ref_update::{COMMIT_WALK_LIMIT, CommitSource, count_commits_by_email}; +use proptest::collection::vec; +use proptest::prelude::*; + +/// In-memory backing for the walker's `CommitSource`. Maps oid -> raw +/// commit bytes; deliberately holds no non-commit objects so a +/// would-be commit reference to a tree returns `None` via the +/// not-a-commit branch. +struct MemSource { + commits: HashMap, +} + +#[async_trait] +impl CommitSource for &mut MemSource { + async fn try_find_commit(&mut self, oid: ObjectId) -> Result, PgError> { + Ok(self.commits.get(&oid).cloned()) + } +} + +/// One commit in the generated DAG. `parents` reference earlier +/// indices (i.e. the DAG is topologically sorted by construction — +/// no cycles). +#[derive(Debug, Clone)] +struct PlannedCommit { + parents: Vec, + email: String, +} + +/// A linearized DAG plus the encoded objects for each node, indexed +/// by the plan's position. +struct EncodedDag { + /// Plan position -> oid of the encoded commit. + oids: Vec, + /// Plan position -> parents (as plan positions). + parents: Vec>, + /// Plan position -> author email. + emails: Vec, + source: MemSource, +} + +/// Encode `plan` into commit objects so each commit's parents point at +/// the oids of earlier-in-plan commits. Returns the DAG. +fn encode_dag(plan: &[PlannedCommit]) -> EncodedDag { + let mut oids: Vec = Vec::with_capacity(plan.len()); + let mut commits: HashMap = HashMap::new(); + let mut emails: Vec = Vec::with_capacity(plan.len()); + let mut parents: Vec> = Vec::with_capacity(plan.len()); + + // Placeholder tree oid: all-zero is a valid hex pattern but doesn't + // exist in our `MemSource`. The walker only descends into parents, + // not trees, so the tree oid never gets resolved. + let tree_oid = ObjectId::null(HashKind::Sha1); + + for (idx, planned) in plan.iter().enumerate() { + let sig = Signature { + name: "Tester".into(), + email: planned.email.clone().into(), + time: gix_date::Time { + seconds: idx as i64, + offset: 0, + }, + }; + let parent_oids: Vec = planned.parents.iter().map(|p| oids[*p]).collect(); + + let commit = Commit { + tree: tree_oid, + parents: parent_oids.into(), + author: sig.clone(), + committer: sig, + encoding: None, + message: format!("planned commit {idx}").into(), + extra_headers: vec![], + }; + + let mut bytes = Vec::new(); + commit.write_to(&mut bytes).expect("encode commit"); + let oid = + gix_object::compute_hash(HashKind::Sha1, Kind::Commit, &bytes).expect("hash commit"); + commits.insert(oid, Bytes::from(bytes)); + oids.push(oid); + emails.push(planned.email.clone()); + parents.push(planned.parents.clone()); + } + + EncodedDag { + oids, + parents, + emails, + source: MemSource { commits }, + } +} + +/// Reachable set from a tip in plan-index space, optionally stopping +/// when the boundary index is encountered. Used as the oracle: the +/// walker must visit at most this set, never more. +fn reachable_from(parents: &[Vec], tip: usize, stop: Option) -> Vec { + let mut out: Vec = Vec::new(); + let mut seen = std::collections::HashSet::new(); + let mut stack = vec![tip]; + while let Some(i) = stack.pop() { + if Some(i) == stop { + continue; + } + if !seen.insert(i) { + continue; + } + out.push(i); + for p in &parents[i] { + stack.push(*p); + } + } + out +} + +// ---------- strategies ---------- + +/// A small set of distinct author emails the strategy chooses from. +/// Tightly bounded so we hit duplicate-author counting paths often. +const EMAILS: &[&str] = &["a@e", "b@e", "c@e"]; + +/// Generate a topologically sorted DAG of N commits (1..=N_MAX). Each +/// commit's parents are a sorted subset of earlier indices, capped at +/// 2 parents (covers merge commits and linear chains). +fn dag_strategy(n_max: usize) -> impl Strategy> { + (1usize..=n_max).prop_flat_map(|n| { + let mut nodes: Vec> = Vec::with_capacity(n); + for i in 0..n { + // Each commit i may pick up to 2 parents from {0..i}. + let parent_pool = if i == 0 { 0..1usize } else { 0..i }; + let email = proptest::sample::select(EMAILS).prop_map(|s| s.to_string()); + let parents_strat = if i == 0 { + Just(Vec::::new()).boxed() + } else { + vec(parent_pool, 0..=2usize.min(i)) + .prop_map(|mut v| { + v.sort_unstable(); + v.dedup(); + v + }) + .boxed() + }; + nodes.push( + (parents_strat, email) + .prop_map(|(parents, email)| PlannedCommit { parents, email }) + .boxed(), + ); + } + nodes + }) +} + +// ---------- properties ---------- + +proptest! { + #![proptest_config(ProptestConfig { + // Walker has no I/O; iterate generously. + cases: 256, + max_shrink_iters: 1024, + ..ProptestConfig::default() + })] + + /// Termination + the walk visits no more than COMMIT_WALK_LIMIT + /// commits. (Implicitly: termination — proptest would time out + /// if the walker hung.) + #[test] + fn walker_is_bounded(plan in dag_strategy(40)) { + let dag = run_walker(&plan, /*old=*/ None); + let total: i64 = dag.counts.values().sum(); + prop_assert!(total <= COMMIT_WALK_LIMIT as i64); + } + + /// When `new_oid` is reachable from `old_oid` (= `old` is an + /// ancestor of `new`), no commit at-or-beyond `old` is counted. + /// The walker should hit the stop boundary and abandon that + /// subtree. + #[test] + fn walker_stops_at_old_oid(plan in dag_strategy(40), tip_idx in 0usize..40) { + prop_assume!(tip_idx < plan.len()); + // Choose an ancestor of `tip_idx` as the boundary. If none, + // skip — there's nothing to test. + let parents_view: Vec> = plan.iter().map(|p| p.parents.clone()).collect(); + let ancestors = reachable_from(&parents_view, tip_idx, None); + let Some(&boundary) = ancestors + .iter() + .filter(|&&i| i != tip_idx) + .next() + else { + return Ok(()); + }; + + let dag = run_walker_at(&plan, tip_idx, Some(boundary)); + // Build the oracle: everything reachable from tip, MINUS + // everything at-or-beyond boundary. + let visited_oracle: std::collections::HashSet = reachable_from(&parents_view, tip_idx, None) + .into_iter() + .filter(|&i| !reachable_from(&parents_view, boundary, None).contains(&i)) + .collect(); + + // Every counted email must originate from a commit in the + // oracle visit set. + let mut oracle_emails: HashMap = HashMap::new(); + for i in &visited_oracle { + *oracle_emails.entry(dag.plan_emails[*i].clone()).or_insert(0) += 1; + } + // Cap by the limit. The walker uses pop-from-back order + // (LIFO), so the boundedness check applies to whichever + // commits get popped first; we just need: counts.sum() <= + // min(|oracle|, COMMIT_WALK_LIMIT). + let total: i64 = dag.counts.values().sum(); + let oracle_total: i64 = oracle_emails.values().sum(); + prop_assert!(total <= oracle_total.min(COMMIT_WALK_LIMIT as i64), + "counts {total} exceed oracle {oracle_total} or limit"); + // Every email that appears with positive count must appear + // in the oracle (no leakage past the boundary). + for (email, count) in &dag.counts { + prop_assert!(*count > 0); + let oracle_count = oracle_emails.get(email).copied().unwrap_or(0); + prop_assert!(oracle_count >= *count, + "email {email}: walker counted {count}, oracle has {oracle_count}"); + } + } + + /// `new_oid == 0...0` (a delete) yields an empty counts map. + #[test] + fn walker_empty_on_delete(plan in dag_strategy(10)) { + let dag = encode_dag(&plan); + let mut source = dag.source; + let result = crate::common::runtime().block_on(count_commits_by_email( + &mut &mut source, + ObjectId::null(HashKind::Sha1), + ObjectId::null(HashKind::Sha1), + HashKind::Sha1, + )).expect("walker"); + prop_assert!(result.is_empty()); + } + + /// The sum of returned counts equals the number of unique + /// commits the walker successfully attributed — there's no + /// double-counting and no off-by-one. + #[test] + fn walker_no_double_count(plan in dag_strategy(40)) { + let dag = run_walker(&plan, /*old=*/ None); + let total: i64 = dag.counts.values().sum(); + // The walker walks N <= min(reachable, limit) commits and + // each contributes exactly 1 to its author email. So total + // must equal the visited count. We don't have direct + // observability into visited, but bounding both sides: + // it must be <= limit AND <= reachable. + let parents_view: Vec> = plan.iter().map(|p| p.parents.clone()).collect(); + let reachable_total = reachable_from(&parents_view, plan.len() - 1, None).len() as i64; + prop_assert!(total <= reachable_total); + prop_assert!(total <= COMMIT_WALK_LIMIT as i64); + // Authors are drawn from a 3-element pool; if total > 0, + // at least one bucket must be present. + if reachable_total > 0 { + prop_assert!(total > 0); + prop_assert!(!dag.counts.is_empty()); + } + } +} + +// ---------- harness ---------- + +struct WalkOutput { + counts: HashMap, + plan_emails: Vec, +} + +fn run_walker(plan: &[PlannedCommit], _old: Option) -> WalkOutput { + let tip = plan.len() - 1; + run_walker_at(plan, tip, None) +} + +fn run_walker_at(plan: &[PlannedCommit], tip: usize, boundary: Option) -> WalkOutput { + let dag = encode_dag(plan); + let plan_emails = dag.emails.clone(); + let new_oid = dag.oids[tip]; + let old_oid = boundary + .map(|i| dag.oids[i]) + .unwrap_or_else(|| ObjectId::null(HashKind::Sha1)); + + // Silence the parents field — only used for oracle building in + // the call sites; here we hand `source` to the walker. + let _ = dag.parents; + + let mut source = dag.source; + let counts = crate::common::runtime() + .block_on(count_commits_by_email( + &mut &mut source, + old_oid, + new_oid, + HashKind::Sha1, + )) + .expect("walker should not error on synthetic in-memory source"); + WalkOutput { + counts, + plan_emails, + } +} diff --git a/mnemosyne_tests/tests/integration/stress_events.rs b/mnemosyne_tests/tests/integration/stress_events.rs new file mode 100644 index 0000000..bc53fdc --- /dev/null +++ b/mnemosyne_tests/tests/integration/stress_events.rs @@ -0,0 +1,294 @@ +//! Sustained event-emission stress test. +//! +//! Pumps `STRESS_EVENTS_TOTAL` events through [`insert_event_in_tx`] +//! from `STRESS_EVENTS_WORKERS` concurrent producers, while one +//! subscriber tails `/events` via WebSocket. Verifies the +//! late-advisory-lock design (slice 1, issue #1) actually scales: +//! +//! - **Correctness** — every committed event reaches the subscriber +//! exactly once, in strictly increasing seq order. +//! - **Throughput** — events committed / second across all +//! producers, reported as wall-clock metric. +//! - **Catch-up margin** — wall-clock seconds the subscriber takes +//! to drain everything *after* producers stop. A small margin +//! means the live-tail is keeping up; a large one is early warning +//! that the advisory lock is starting to bottleneck the path and +//! the outbox-pattern upgrade (see issue #1) is worth doing. +//! +//! ## Baseline (recorded on author's M-series laptop, release build) +//! +//! - 2 000 events / 16 workers → ~2 350 events/sec +//! - 20 000 events / 64 workers → ~2 120 events/sec +//! +//! Throughput **does not scale with worker count past ~16** — every +//! committer serializes on `EVENTS_ADVISORY_LOCK_ID` during the +//! INSERT + commit window. The ceiling is "how fast a single +//! Postgres connection can commit one INSERT," around 0.5 ms here. +//! That's ample for the immediate appview-firehose use case; if a +//! deployment ever observes sustained > 2 k ev/s of push activity +//! the outbox upgrade should land. Use this test under +//! `STRESS_LEN=20000 STRESS_EVENTS_WORKERS=64` to spot a regression +//! (e.g. accidentally enlarging the locked section). +//! +//! Catch-up margin should be sub-second for any healthy run — a +//! growing margin under load is the first sign the WS path is +//! lagging behind producers. +//! +//! ## Tunable environment variables +//! +//! - `STRESS_LEN` — required; gates the test out of the +//! default `just test` filter, matching +//! the existing stress convention. +//! Also acts as the default for +//! `STRESS_EVENTS_TOTAL`. +//! - `STRESS_EVENTS_TOTAL` — total events to commit (default +//! `STRESS_LEN` if set, else 10_000). +//! - `STRESS_EVENTS_WORKERS` — concurrent producers (default 16). +//! +//! The test reuses the suite's shared Postgres pool (sized in +//! `common/mod.rs`); raising worker counts past that pool's +//! `max_connections` just serializes producers at the connection +//! acquire — not a useful stress signal. Increase the shared pool +//! ceiling there if your stress profile needs it. +//! +//! ## Why this isn't a proptest +//! +//! The proptest `concurrent_writers_yield_strictly_increasing_seq` +//! covers correctness at N=16. This test is here to surface +//! *performance* characteristics (throughput, lock-wait, catch-up +//! margin) and crank N high enough that a regression in the +//! lock-hold window — e.g. accidentally moving expensive work into +//! the locked section — shows up as a noticeable throughput drop. + +#![allow(unreachable_pub)] + +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{Duration, Instant}; + +use axum::Router; +use axum::routing::get; +use futures::StreamExt; +use mnemosyne_postgres::insert_event_in_tx; +use mnemosyne_protocol::events_stream::{self, EventsStreamState}; +use serde_json::{Value, json}; +use sqlx::PgPool; +use tokio::net::TcpListener; +use tokio_tungstenite::tungstenite::Message; + +const DEFAULT_TOTAL_FALLBACK: u64 = 10_000; +const DEFAULT_WORKERS: u64 = 16; + +fn env_u64(name: &str, default: u64) -> u64 { + std::env::var(name) + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(default) +} + +struct StressConfig { + total: u64, + workers: u64, +} + +impl StressConfig { + fn from_env() -> Self { + let stress_len = std::env::var("STRESS_LEN") + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(DEFAULT_TOTAL_FALLBACK); + let total = env_u64("STRESS_EVENTS_TOTAL", stress_len); + let workers = env_u64("STRESS_EVENTS_WORKERS", DEFAULT_WORKERS).max(1); + Self { total, workers } + } +} + +async fn spawn_server(pool: PgPool) -> (std::net::SocketAddr, tokio::task::AbortHandle) { + let app: Router = Router::new() + .route("/events", get(events_stream::handler)) + .with_state(EventsStreamState { pool }); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let handle = tokio::spawn(async move { + let _ = axum::serve(listener, app).await; + }); + (addr, handle.abort_handle()) +} + +/// Drain one text frame, skipping pings/pongs. Returns `None` on +/// timeout or non-text frames. +async fn next_text( + stream: &mut tokio_tungstenite::WebSocketStream< + tokio_tungstenite::MaybeTlsStream, + >, + timeout: Duration, +) -> Option { + loop { + let msg = tokio::time::timeout(timeout, stream.next()).await.ok()??; + match msg { + Ok(Message::Text(t)) => return serde_json::from_str(&t).ok(), + Ok(Message::Close(_)) | Err(_) => return None, + Ok(_) => continue, + } + } +} + +#[test] +fn events_emit_and_drain_under_load() { + if std::env::var("STRESS_LEN").is_err() { + eprintln!( + "skipping stress_events: set STRESS_LEN to enable. \ + See file header for tunable env vars." + ); + return; + } + + let cfg = StressConfig::from_env(); + eprintln!("stress_events: total={} workers={}", cfg.total, cfg.workers,); + + crate::common::runtime().block_on(async { + let pool = crate::common::shared_pool().await; + let repo_did = crate::common::unique_repo_did(); + mnemosyne_postgres::ensure_repo( + &pool, + &repo_did, + "did:test:stress-owner", + "stress-events", + crate::common::PLACEHOLDER_RKEY, + ) + .await + .expect("ensure_repo"); + + let starting_max: i64 = sqlx::query_scalar("SELECT COALESCE(MAX(seq), 0) FROM events") + .fetch_one(&pool) + .await + .unwrap(); + + // Subscriber up first, then producers — exercises the + // backfill + live-tail handoff under load. + let (addr, abort) = spawn_server(pool.clone()).await; + let url = format!("ws://{addr}/events?cursor={starting_max}"); + let (mut ws, _resp) = tokio_tungstenite::connect_async(&url) + .await + .expect("ws connect"); + + // Producers: divide `total` evenly across `workers`. + let per_worker = cfg.total / cfg.workers; + let leftover = cfg.total % cfg.workers; + let emitted = Arc::new(AtomicU64::new(0)); + let producers_start = Instant::now(); + + let mut producer_handles = Vec::with_capacity(cfg.workers as usize); + for w in 0..cfg.workers { + let p = pool.clone(); + let rd = repo_did.clone(); + let n = per_worker + if w < leftover { 1 } else { 0 }; + let emitted_counter = emitted.clone(); + producer_handles.push(tokio::spawn(async move { + let mut local_seqs = Vec::with_capacity(n as usize); + for i in 0..n { + let mut tx = p.begin().await.expect("begin"); + let rkey = format!("w{w}-i{i}"); + let seq = insert_event_in_tx( + &mut tx, + &rkey, + "sh.tangled.git.refUpdate", + Some(&rd), + &json!({"worker": w, "i": i}), + ) + .await + .expect("insert"); + tx.commit().await.expect("commit"); + local_seqs.push(seq); + emitted_counter.fetch_add(1, Ordering::Relaxed); + } + local_seqs + })); + } + + // Collect every produced seq. + let mut all_emitted: Vec = Vec::with_capacity(cfg.total as usize); + for h in producer_handles { + let seqs = h.await.expect("producer join"); + all_emitted.extend(seqs); + } + all_emitted.sort(); + let producers_elapsed = producers_start.elapsed(); + + let throughput_eps = cfg.total as f64 / producers_elapsed.as_secs_f64(); + eprintln!( + "stress_events: emitted {} in {:.3}s = {:.0} events/sec", + cfg.total, + producers_elapsed.as_secs_f64(), + throughput_eps, + ); + + // Drain — the subscriber must have observed every emitted + // seq, in order. Filter to OUR repo_did to avoid picking up + // any concurrent test rows (though in practice stress runs + // alone under STRESS_LEN gating). + let drain_start = Instant::now(); + let mut received: Vec = Vec::with_capacity(cfg.total as usize); + // Per-frame timeout grows with `total` so a slow CI doesn't + // false-positive a hang on legitimate work. + let per_frame_timeout = Duration::from_secs(30); + while received.len() < all_emitted.len() { + let Some(frame) = next_text(&mut ws, per_frame_timeout).await else { + panic!( + "stress_events: timed out after receiving {}/{} events", + received.len(), + all_emitted.len(), + ); + }; + if frame["repo"].as_str() != Some(&repo_did) { + // Drop noise from other-repo rows (shouldn't happen + // when running gated, but guards against test + // pollution). + continue; + } + received.push(frame["seq"].as_i64().expect("seq i64")); + } + let drain_elapsed = drain_start.elapsed(); + let catch_up_margin = drain_elapsed.saturating_sub(producers_elapsed); + eprintln!( + "stress_events: drained {} in {:.3}s; catch-up margin = {:.3}s", + received.len(), + drain_elapsed.as_secs_f64(), + catch_up_margin.as_secs_f64(), + ); + + // ----- Correctness assertions ----- + + // 1. Every emitted seq was received. + assert_eq!( + received.len(), + all_emitted.len(), + "subscriber count mismatch: emitted {}, received {}", + all_emitted.len(), + received.len(), + ); + + // 2. Strictly increasing seq order on the wire. + for w in received.windows(2) { + assert!( + w[0] < w[1], + "non-monotonic seq on the wire: {} >= {}", + w[0], + w[1], + ); + } + + // 3. Set equality (no duplicates, no losses). Sorted-equal + // is the strongest form given (1) + (2). + let mut received_sorted = received.clone(); + received_sorted.sort(); + assert_eq!( + received_sorted, all_emitted, + "subscriber received set != producer emitted set (look for \ + missing or duplicated seqs)", + ); + + let _ = ws.close(None).await; + abort.abort(); + }); +} diff --git a/mnemosyne_tests/tests/proptest-regressions/protocol_push_roundtrip.txt b/mnemosyne_tests/tests/proptest-regressions/protocol_push_roundtrip.txt new file mode 100644 index 0000000..fa1a1f8 --- /dev/null +++ b/mnemosyne_tests/tests/proptest-regressions/protocol_push_roundtrip.txt @@ -0,0 +1,8 @@ +# Seeds for failure cases proptest has generated in the past. It is +# automatically read and these particular cases re-run before any +# novel cases are generated. +# +# It is recommended to check this file in to source control so that +# everyone who runs the test benefits from these saved cases. +cc 86ff421224fb9a22f14e5ed856ad27e0fd6be5518ea40a393b9b87d3bad27618 # shrinks to (graph, refs) = (ObjectGraph { objects: [(Sha1(b5f49ec5aa9a968f5c26024c745032c0647d2deb), Blob, [86, 121, 169, 244, 245, 62, 222, 25, 238, 3, 147, 158, 49, 146, 120, 249, 216, 45, 129, 52]), (Sha1(e9aeb710df97e7bbeb79c8803bffddae75a162b6), Tree, [49, 48, 48, 54, 52, 52, 32, 102, 105, 108, 101, 48, 0, 181, 244, 158, 197, 170, 154, 150, 143, 92, 38, 2, 76, 116, 80, 50, 192, 100, 125, 45, 235]), (Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01), Commit, [116, 114, 101, 101, 32, 101, 57, 97, 101, 98, 55, 49, 48, 100, 102, 57, 55, 101, 55, 98, 98, 101, 98, 55, 57, 99, 56, 56, 48, 51, 98, 102, 102, 100, 100, 97, 101, 55, 53, 97, 49, 54, 50, 98, 54, 10, 97, 117, 116, 104, 111, 114, 32, 77, 110, 101, 109, 111, 115, 121, 110, 101, 32, 84, 101, 115, 116, 101, 114, 32, 60, 116, 101, 115, 116, 64, 101, 120, 97, 109, 112, 108, 101, 46, 99, 111, 109, 62, 32, 48, 32, 43, 48, 48, 48, 48, 10, 99, 111, 109, 109, 105, 116, 116, 101, 114, 32, 77, 110, 101, 109, 111, 115, 121, 110, 101, 32, 84, 101, 115, 116, 101, 114, 32, 60, 116, 101, 115, 116, 64, 101, 120, 97, 109, 112, 108, 101, 46, 99, 111, 109, 62, 32, 48, 32, 43, 48, 48, 48, 48, 10, 10, 112, 114, 111, 112, 116, 101, 115, 116, 32, 99, 111, 109, 109, 105, 116])], head_commit: Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01) }, ProtocolRefSet { edits: [RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01)) }, name: FullName("refs/heads/main"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01)) }, name: FullName("refs/heads/gebc"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Symbolic(FullName("refs/heads/main")) }, name: FullName("HEAD"), deref: false }], expected_advertised: [(FullName("HEAD"), Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01)), (FullName("refs/heads/gebc"), Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01)), (FullName("refs/heads/main"), Sha1(c2d84697625d3cdebe0f66fede9108d4e17add01))], head_symref_target: Some(FullName("refs/heads/main")) }) +cc 223ed04c6a0706b07d4c90d88e05ddba9066e3fcc71a42cf09119af8defcfe46 # shrinks to (graph, refs) = (ObjectGraph { objects: [(Sha1(c544304fc0b250bc8995251406cdf3f6367886a7), Blob, [48, 190, 105, 224, 201, 211, 55, 158, 151, 218, 103, 85, 254, 101, 73, 132, 26, 232, 198, 230, 30, 25, 187, 10, 186, 33, 141, 155, 193, 103, 109, 74, 92, 7, 168, 167, 206]), (Sha1(a609d39ba4af722792913851051cd4a65bb0adc8), Blob, [117, 112, 65, 251, 137, 53, 157, 60, 241, 214, 206, 103, 1, 152, 200, 62, 106, 10, 233, 93, 82, 154, 193, 223, 56, 202, 34, 212, 149, 72, 53, 66, 94, 119, 9, 146, 98, 40, 43, 6, 221, 51, 194, 35, 191, 62, 102]), (Sha1(2102f5f1234bdd381584f610e703feb8b572cb8a), Blob, [217, 19, 220, 63, 204, 175, 65, 83, 86, 48, 30, 40, 202, 238, 144, 72, 135, 30, 177, 82, 32, 122, 41, 72, 16, 62, 30, 119, 23, 33, 117, 86, 234, 114, 128, 229, 192, 226, 150, 183, 26, 44, 109, 68, 59, 198, 126, 44, 60, 37, 192, 210, 65, 110, 131, 66, 84, 119, 153]), (Sha1(3e31490af0af9da97ca20491304ec5b7c940dbf0), Tree, [49, 48, 48, 54, 52, 52, 32, 102, 105, 108, 101, 48, 0, 197, 68, 48, 79, 192, 178, 80, 188, 137, 149, 37, 20, 6, 205, 243, 246, 54, 120, 134, 167, 49, 48, 48, 54, 52, 52, 32, 102, 105, 108, 101, 49, 0, 166, 9, 211, 155, 164, 175, 114, 39, 146, 145, 56, 81, 5, 28, 212, 166, 91, 176, 173, 200, 49, 48, 48, 54, 52, 52, 32, 102, 105, 108, 101, 50, 0, 33, 2, 245, 241, 35, 75, 221, 56, 21, 132, 246, 16, 231, 3, 254, 184, 181, 114, 203, 138]), (Sha1(8fa61c101a19843e38745446813d2db902998379), Commit, [116, 114, 101, 101, 32, 51, 101, 51, 49, 52, 57, 48, 97, 102, 48, 97, 102, 57, 100, 97, 57, 55, 99, 97, 50, 48, 52, 57, 49, 51, 48, 52, 101, 99, 53, 98, 55, 99, 57, 52, 48, 100, 98, 102, 48, 10, 97, 117, 116, 104, 111, 114, 32, 77, 110, 101, 109, 111, 115, 121, 110, 101, 32, 84, 101, 115, 116, 101, 114, 32, 60, 116, 101, 115, 116, 64, 101, 120, 97, 109, 112, 108, 101, 46, 99, 111, 109, 62, 32, 48, 32, 43, 48, 48, 48, 48, 10, 99, 111, 109, 109, 105, 116, 116, 101, 114, 32, 77, 110, 101, 109, 111, 115, 121, 110, 101, 32, 84, 101, 115, 116, 101, 114, 32, 60, 116, 101, 115, 116, 64, 101, 120, 97, 109, 112, 108, 101, 46, 99, 111, 109, 62, 32, 48, 32, 43, 48, 48, 48, 48, 10, 10, 112, 114, 111, 112, 116, 101, 115, 116, 32, 99, 111, 109, 109, 105, 116])], head_commit: Sha1(8fa61c101a19843e38745446813d2db902998379) }, ProtocolRefSet { edits: [RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(8fa61c101a19843e38745446813d2db902998379)) }, name: FullName("refs/heads/main"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(8fa61c101a19843e38745446813d2db902998379)) }, name: FullName("refs/heads/ah"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Symbolic(FullName("refs/heads/main")) }, name: FullName("HEAD"), deref: false }], expected_advertised: [(FullName("HEAD"), Sha1(8fa61c101a19843e38745446813d2db902998379)), (FullName("refs/heads/ah"), Sha1(8fa61c101a19843e38745446813d2db902998379)), (FullName("refs/heads/main"), Sha1(8fa61c101a19843e38745446813d2db902998379))], head_symref_target: Some(FullName("refs/heads/main")) }) diff --git a/mnemosyne_tests/tests/proptest-regressions/ref_update_walk.txt b/mnemosyne_tests/tests/proptest-regressions/ref_update_walk.txt new file mode 100644 index 0000000..0d199e0 --- /dev/null +++ b/mnemosyne_tests/tests/proptest-regressions/ref_update_walk.txt @@ -0,0 +1,7 @@ +# Seeds for failure cases proptest has generated in the past. It is +# automatically read and these particular cases re-run before any +# novel cases are generated. +# +# It is recommended to check this file in to source control so that +# everyone who runs the test benefits from these saved cases. +cc c8ef9de16e3410f159e69b6ed43648de9e4d329e06cc09dd1b0b79c6f6cd86b4 # shrinks to plan = [PlannedCommit { parents: [], email: "a@e" }, PlannedCommit { parents: [0], email: "a@e" }, PlannedCommit { parents: [0, 1], email: "a@e" }], tip_idx = 2