diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..5f5d669 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,46 @@ +name: CI + +on: + workflow_call: + pull_request: + push: + branches: [main] + +permissions: + contents: read + +concurrency: + group: queue-ci-${{ github.workflow }}-${{ github.ref }} + cancel-in-progress: true + +jobs: + verify: + runs-on: ubuntu-latest + services: + postgres: + image: postgres:16 + env: + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + POSTGRES_DB: animus_test + ports: + - 5432:5432 + options: >- + --health-cmd "pg_isready -U postgres -d animus_test" + --health-interval 5s + --health-timeout 5s + --health-retries 10 + env: + TEST_DATABASE_URL: postgres://postgres:postgres@localhost:5432/animus_test + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + with: + components: clippy, rustfmt + - uses: Swatinem/rust-cache@v2 + - name: Formatting + run: cargo fmt --all -- --check + - name: Clippy + run: cargo clippy --all-targets --all-features -- -D warnings + - name: Tests + run: cargo test --all-targets diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index b99d5ad..ea8bcb6 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -1,5 +1,12 @@ name: Release Binaries +permissions: + contents: read + +concurrency: + group: queue-release-${{ github.ref }} + cancel-in-progress: false + on: push: tags: @@ -12,8 +19,12 @@ on: type: string jobs: + verify: + uses: ./.github/workflows/ci.yml + build: name: Build (${{ matrix.target }}) + needs: verify runs-on: ${{ matrix.os }} strategy: fail-fast: false @@ -131,9 +142,37 @@ jobs: find dist -type f \( -name '*.tar.gz' -o -name '*.tar.gz.sha256' \) ! -path "${ASSETS_DIR}/*" -exec cp {} "${ASSETS_DIR}/" \; ls -la "${ASSETS_DIR}" + - name: Fail closed if the release already exists + env: + GH_TOKEN: ${{ github.token }} + shell: bash + run: | + set -euo pipefail + if gh release view "${GITHUB_REF_NAME}" >/dev/null 2>&1; then + echo "::error::release ${GITHUB_REF_NAME} already exists; release tags are write-once" >&2 + exit 1 + fi + - name: Publish release uses: softprops/action-gh-release@v2 with: files: dist/release-assets/* fail_on_unmatched_files: true generate_release_notes: true + + - name: Verify published assets against bound checksums + env: + GH_TOKEN: ${{ github.token }} + shell: bash + run: | + set -euo pipefail + mkdir -p verify + gh release download "${GITHUB_REF_NAME}" --dir verify + test "$(find verify -maxdepth 1 -name '*.tar.gz' | wc -l | tr -d ' ')" = "3" + test "$(find verify -maxdepth 1 -name '*.tar.gz.sha256' | wc -l | tr -d ' ')" = "3" + for checksum in verify/*.tar.gz.sha256; do + archive="${checksum%.sha256}" + expected="$(tr -d '[:space:]' < "${checksum}")" + actual="$(sha256sum "${archive}" | awk '{print $1}')" + test "${actual}" = "${expected}" + done diff --git a/Cargo.lock b/Cargo.lock index 4244440..b0c2057 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -39,10 +39,30 @@ dependencies = [ "libc", ] +[[package]] +name = "animus-actor" +version = "0.1.0" +source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.7.0-rc.14#63c60d573090a98c1a36c8f469e86e0dbcd6c712" +dependencies = [ + "schemars", + "serde", +] + +[[package]] +name = "animus-execution-protocol" +version = "0.1.0" +source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.7.0-rc.14#63c60d573090a98c1a36c8f469e86e0dbcd6c712" +dependencies = [ + "chrono", + "schemars", + "serde", + "serde_json", +] + [[package]] name = "animus-plugin-protocol" -version = "0.1.14" -source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.5.10#a3b39dcbe02e8c3701f85b89aca55ae010193620" +version = "0.1.18" +source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.7.0-rc.14#63c60d573090a98c1a36c8f469e86e0dbcd6c712" dependencies = [ "schemars", "serde", @@ -51,8 +71,9 @@ dependencies = [ [[package]] name = "animus-queue-postgres" -version = "0.1.0" +version = "0.2.0" dependencies = [ + "animus-execution-protocol", "animus-plugin-protocol", "animus-queue-protocol", "animus-subject-protocol", @@ -71,9 +92,10 @@ dependencies = [ [[package]] name = "animus-queue-protocol" -version = "0.3.2" -source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.5.10#a3b39dcbe02e8c3701f85b89aca55ae010193620" +version = "0.4.0" +source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.7.0-rc.14#63c60d573090a98c1a36c8f469e86e0dbcd6c712" dependencies = [ + "animus-execution-protocol", "animus-plugin-protocol", "animus-subject-protocol", "schemars", @@ -83,9 +105,10 @@ dependencies = [ [[package]] name = "animus-subject-protocol" -version = "0.1.15" -source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.5.10#a3b39dcbe02e8c3701f85b89aca55ae010193620" +version = "0.2.0" +source = "git+https://github.com/launchapp-dev/animus-protocol?tag=v0.7.0-rc.14#63c60d573090a98c1a36c8f469e86e0dbcd6c712" dependencies = [ + "animus-actor", "animus-plugin-protocol", "anyhow", "async-trait", diff --git a/Cargo.toml b/Cargo.toml index 87d10d3..dc344d3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,9 +1,9 @@ [package] name = "animus-queue-postgres" -version = "0.1.1" +version = "0.2.0" edition = "2021" license = "Elastic-2.0" -description = "Durable Postgres-backed queue plugin for Animus — survives daemon restarts/redeploys and reclaims crashed leases. Wire-compatible with animus-queue-default (the `queue` PluginKind)." +description = "Durable generation-fenced Postgres queue plugin for Animus." repository = "https://github.com/launchapp-dev/animus-queue-postgres" homepage = "https://github.com/launchapp-dev/animus-cli" default-run = "animus-queue-postgres" @@ -17,14 +17,12 @@ name = "animus-queue-postgres" path = "src/main.rs" [dependencies] -# Match animus-queue-default v0.3.3's protocol pin EXACTLY (tag v0.5.10 on -# launchapp-dev/animus-protocol) so the wire types are byte-identical to the -# reference queue plugin the v0.6.x daemon already talks to. The queue RPC -# surface is frozen here regardless of the portal's config/subject/chat plugins -# pinning v0.1.25 — those are a different protocol crate family. -animus-plugin-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.5.10" } -animus-queue-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.5.10" } -animus-subject-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.5.10" } +# Pin the released execution-fence contract exactly. The v2 queue methods are +# additive; legacy queue/* calls remain available during daemon migration. +animus-execution-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.7.0-rc.14" } +animus-plugin-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.7.0-rc.14" } +animus-queue-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.7.0-rc.14" } +animus-subject-protocol = { git = "https://github.com/launchapp-dev/animus-protocol", tag = "v0.7.0-rc.14" } tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "fs", "io-util", "io-std"] } serde = { version = "1", features = ["derive"] } diff --git a/README.md b/README.md index 53b48ed..299a4d3 100644 --- a/README.md +++ b/README.md @@ -1,3 +1,37 @@ # animus-queue-postgres -Standalone Animus plugin (Postgres-backed). +Durable Postgres queue plugin for Animus. + +Version 0.2 adds the generation-fenced `queue/v2/*` contract used by the +five-slot coding fleet: + +- idempotent enqueue with a monotonic generation per qualified subject; +- exact repository/base/head-ref reservations; +- Pending-only fresh leasing (expired assignments are never silently reused); +- compare-and-swap lease renewal, expired-lease recovery, completion, and + return-to-pending; +- stable workflow and subject identity across daemon restart/recovery; +- stale-owner fencing after every lease transfer. + +Legacy `queue/*` methods remain available for older daemons. They cannot lease +or mutate v2 entries, so mixed-version deployment fails closed instead of +creating a second workflow or node. + +## Configuration + +- `DATABASE_URL` or `ANIMUS_POSTGRES_URL`: Postgres connection URL. +- `ANIMUS_QUEUE_LEASE_TTL_SECS`: maximum lease TTL (default 1800 seconds). +- `ANIMUS_QUEUE_TABLE`: safe table identifier (default `queue_item`). + +The plugin applies additive, idempotent schema changes. Its health check is not +ready until both Postgres and the complete v2 schema are available. + +## Verification + +```sh +cargo test --all-targets +cargo clippy --all-targets --all-features -- -D warnings +``` + +Postgres integration tests use `TEST_DATABASE_URL`, falling back to +`postgres://postgres:postgres@localhost:55432/animus_test`. diff --git a/src/config.rs b/src/config.rs index ea9f0a5..085d1d0 100644 --- a/src/config.rs +++ b/src/config.rs @@ -11,7 +11,8 @@ pub const ENV_POSTGRES_URL: &str = "ANIMUS_POSTGRES_URL"; /// Override the lease time-to-live, in seconds. A leased entry whose /// `lease_expires_at` (= lease time + this TTL) has passed is reclaimable by a -/// later `queue/lease`. Default [`DEFAULT_LEASE_TTL_SECS`]. +/// later legacy `queue/lease`. V2 assignments instead require explicit fenced +/// recovery. Default [`DEFAULT_LEASE_TTL_SECS`]. pub const ENV_LEASE_TTL_SECS: &str = "ANIMUS_QUEUE_LEASE_TTL_SECS"; /// Override the table name (default `queue_item`). Useful when more than one @@ -20,7 +21,7 @@ pub const ENV_TABLE: &str = "ANIMUS_QUEUE_TABLE"; /// Default lease TTL: 30 minutes. Long enough that a healthy in-flight /// workflow's lease never expires under it, short enough that a crashed -/// daemon's work is reclaimed promptly on the next lease. +/// daemon's work becomes eligible for explicit recovery promptly. pub const DEFAULT_LEASE_TTL_SECS: i64 = 1800; /// Default durable queue table name. @@ -44,16 +45,8 @@ impl QueueConfig { let database_url = std::env::var(ENV_DATABASE_URL) .ok() .filter(|s| !s.is_empty()) - .or_else(|| { - std::env::var(ENV_POSTGRES_URL) - .ok() - .filter(|s| !s.is_empty()) - }) - .ok_or_else(|| { - anyhow!( - "no Postgres URL configured: set {ENV_DATABASE_URL} (or {ENV_POSTGRES_URL})" - ) - })?; + .or_else(|| std::env::var(ENV_POSTGRES_URL).ok().filter(|s| !s.is_empty())) + .ok_or_else(|| anyhow!("no Postgres URL configured: set {ENV_DATABASE_URL} (or {ENV_POSTGRES_URL})"))?; let lease_ttl_secs = std::env::var(ENV_LEASE_TTL_SECS) .ok() @@ -67,11 +60,7 @@ impl QueueConfig { .filter(|s| is_safe_ident(s)) .unwrap_or_else(|| DEFAULT_TABLE.to_string()); - Ok(Self { - database_url, - lease_ttl_secs, - table, - }) + Ok(Self { database_url, lease_ttl_secs, table }) } /// In-process builder for tests / embedders. diff --git a/src/lib.rs b/src/lib.rs index 429063d..ef72c0d 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,24 +1,22 @@ //! `animus-queue-postgres`: a durable, Postgres-backed `queue` plugin for //! Animus. //! -//! This is a drop-in replacement for `launchapp-dev/animus-queue-default`. It -//! speaks the EXACT same `queue/*` JSON-RPC contract (from -//! `animus-queue-protocol` v0.5.10, the tag queue-default v0.3.3 pins) so the -//! daemon needs no change to use it. The only difference is the storage layer: -//! the reference plugin keeps queue state in a file-locked JSON blob under the -//! project root (ephemeral on a container redeploy), whereas this plugin keeps -//! it in a shared Postgres table (`queue_item`) so the dispatch queue SURVIVES -//! daemon restarts and container redeploys. -//! -//! ## Durability + lease-expiry reclaim -//! -//! The headline feature is **lease-expiry reclaim**. Each `queue/lease` -//! stamps `lease_owner` + `lease_expires_at` on the claimed rows. A `leased` -//! row whose `lease_expires_at` has passed is treated as re-leasable: the next -//! `queue/lease` reclaims it via an atomic `SELECT ... FOR UPDATE SKIP LOCKED`. -//! So if the daemon crashes or the container redeploys mid-workflow, the -//! unfinished work is re-dispatched instead of being lost (the file backend -//! would have lost the whole queue). +//! It retains the legacy `queue/*` surface while adding the generation-fenced +//! `queue/v2/*` scheduler contract. Queue state and execution ownership live in +//! shared Postgres tables so daemon restarts do not erase or duplicate work. +//! +//! ## Generation-fenced recovery +//! +//! V2 enqueue allocates a monotonic generation for a qualified subject and can +//! reserve one exact repository/head ref. Fresh leasing only considers Pending +//! rows. When an Assigned lease expires, the daemon must reconcile the existing +//! workflow and explicitly recover its lease with a compare-and-swap fence. +//! Recovery preserves workflow/subject generations and increments only lease +//! generation, so the previous daemon can no longer mutate the entry. +//! +//! Legacy leasing keeps its historical expiry-reclaim behavior, but legacy +//! methods cannot lease or mutate v2 rows. Mixed-version rollouts therefore +//! fail closed. //! //! ## State model //! @@ -51,3 +49,6 @@ pub mod store; pub use config::QueueConfig; pub use store::Store; + +/// Maximum generation-fenced leases granted in one scheduling call. +pub const MAX_GENERATION_FENCED_LEASE_BATCH: usize = 5; diff --git a/src/plugin.rs b/src/plugin.rs index c8ec7ec..caa3a78 100644 --- a/src/plugin.rs +++ b/src/plugin.rs @@ -1,9 +1,8 @@ //! Stdio JSON-RPC loop for the `animus-queue-postgres` plugin. //! //! Handles `initialize`, `$/ping`, `health/check`, `shutdown`, `exit`, -//! `--manifest` / `--help` CLI shortcuts, and the 12 `queue/*` methods — the -//! EXACT method set + request/response shapes of `animus-queue-default` -//! (`animus-queue-protocol` v0.5.10). +//! `--manifest` / `--help` CLI shortcuts, the legacy `queue/*` surface, and the +//! generation-fenced `queue/v2/*` scheduler contract. //! //! The Postgres pool is opened lazily at `initialize` (not at process start) //! so the `--manifest` probe `animus install --locked` runs at build time @@ -13,17 +12,19 @@ use std::io::{self, IsTerminal, Write}; use std::sync::Arc; use animus_plugin_protocol::{ - error_codes as plugin_error_codes, EnvRequirement, HealthCheckResult, HealthStatus, - InitializeResult, KindCapability, PluginCapabilities, PluginInfo, PluginManifest, RpcError, - RpcRequest, RpcResponse, PLUGIN_KIND_QUEUE, PROTOCOL_VERSION, + error_codes as plugin_error_codes, EnvRequirement, HealthCheckResult, HealthStatus, InitializeResult, + KindCapability, PluginCapabilities, PluginInfo, PluginManifest, RpcError, RpcRequest, RpcResponse, + PLUGIN_KIND_QUEUE, PROTOCOL_VERSION, }; use animus_queue_protocol::{ - error_codes as queue_error_codes, QueueCapabilities, QueueCompletionRequest, QueueDropRequest, - QueueEnqueueRequest, QueueEnqueueResponse, QueueHoldRequest, QueueLeaseRequest, - QueueListRequest, QueueMarkAssignedRequest, QueueReleasePendingParams, QueueReleaseRequest, - QueueReorderRequest, KIND, METHOD_QUEUE_COMPLETION, METHOD_QUEUE_DROP, METHOD_QUEUE_ENQUEUE, - METHOD_QUEUE_HOLD, METHOD_QUEUE_LEASE, METHOD_QUEUE_LIST, METHOD_QUEUE_MARK_ASSIGNED, - METHOD_QUEUE_NEXT_DEADLINE, METHOD_QUEUE_RELEASE, METHOD_QUEUE_RELEASE_PENDING, + error_codes as queue_error_codes, QueueCapabilities, QueueCompletionRequest, QueueCompletionV2Request, + QueueDropRequest, QueueEnqueueRequest, QueueEnqueueResponse, QueueEnqueueV2Request, QueueHoldRequest, + QueueLeaseRecoverRequest, QueueLeaseRenewRequest, QueueLeaseRequest, QueueLeaseV2Request, QueueListRequest, + QueueMarkAssignedRequest, QueueReleasePendingParams, QueueReleasePendingV2Request, QueueReleaseRequest, + QueueReorderRequest, KIND, METHOD_QUEUE_COMPLETION, METHOD_QUEUE_COMPLETION_V2, METHOD_QUEUE_DROP, + METHOD_QUEUE_ENQUEUE, METHOD_QUEUE_ENQUEUE_V2, METHOD_QUEUE_HOLD, METHOD_QUEUE_LEASE, METHOD_QUEUE_LEASE_RECOVER, + METHOD_QUEUE_LEASE_RENEW, METHOD_QUEUE_LEASE_V2, METHOD_QUEUE_LIST, METHOD_QUEUE_MARK_ASSIGNED, + METHOD_QUEUE_NEXT_DEADLINE, METHOD_QUEUE_RELEASE, METHOD_QUEUE_RELEASE_PENDING, METHOD_QUEUE_RELEASE_PENDING_V2, METHOD_QUEUE_REORDER, METHOD_QUEUE_STATS, PROTOCOL_VERSION as QUEUE_PROTOCOL_VERSION, }; use anyhow::Result; @@ -33,11 +34,11 @@ use tokio::sync::{Mutex, RwLock}; use crate::config::QueueConfig; use crate::store::{QueueLeaseError, QueueReleasePendingError, Store}; +use crate::MAX_GENERATION_FENCED_LEASE_BATCH; const PLUGIN_NAME: &str = "animus-queue-postgres"; const PLUGIN_VERSION: &str = env!("CARGO_PKG_VERSION"); -const PLUGIN_DESCRIPTION: &str = - "Durable Postgres-backed queue plugin for Animus (survives restarts; reclaims crashed leases)."; +const PLUGIN_DESCRIPTION: &str = "Durable generation-fenced Postgres queue plugin for Animus."; /// Stable entrypoint. Call from `#[tokio::main]` in `main.rs`. pub async fn run() -> Result<()> { @@ -69,18 +70,14 @@ pub async fn run() -> Result<()> { buffer.extend_from_slice(&chunk[..n]); loop { - let leading_ws = buffer - .iter() - .take_while(|b| b.is_ascii_whitespace()) - .count(); + let leading_ws = buffer.iter().take_while(|b| b.is_ascii_whitespace()).count(); if leading_ws > 0 { buffer.drain(..leading_ws); } if buffer.is_empty() { break; } - let mut stream = - serde_json::Deserializer::from_slice(&buffer).into_iter::(); + let mut stream = serde_json::Deserializer::from_slice(&buffer).into_iter::(); match stream.next() { Some(Ok(request)) => { let consumed = stream.byte_offset(); @@ -120,7 +117,9 @@ fn handle_cli_args() -> bool { eprintln!("Usage:"); eprintln!(" {PLUGIN_NAME} --manifest Print plugin manifest as JSON and exit"); eprintln!(" {PLUGIN_NAME} Run JSON-RPC loop on stdin/stdout"); - eprintln!("Env: DATABASE_URL (or ANIMUS_POSTGRES_URL), ANIMUS_QUEUE_LEASE_TTL_SECS, ANIMUS_QUEUE_TABLE"); + eprintln!( + "Env: DATABASE_URL (or ANIMUS_POSTGRES_URL), ANIMUS_QUEUE_LEASE_TTL_SECS, ANIMUS_QUEUE_TABLE" + ); return true; } _ => {} @@ -134,6 +133,7 @@ fn print_manifest() { name: PLUGIN_NAME.to_string(), version: PLUGIN_VERSION.to_string(), plugin_kind: PLUGIN_KIND_QUEUE.to_string(), + plugin_kinds: Vec::new(), description: PLUGIN_DESCRIPTION.to_string(), protocol_version: PROTOCOL_VERSION.to_string(), capabilities: queue_methods().into_iter().map(|m| m.to_string()).collect(), @@ -144,13 +144,10 @@ fn print_manifest() { // Postgres plugins. env_required: env_requirements(), notification_buffer_size: None, + supports_mcp: None, }; let mut stdout = io::stdout().lock(); - let _ = writeln!( - stdout, - "{}", - serde_json::to_string(&manifest).expect("serialize manifest") - ); + let _ = writeln!(stdout, "{}", serde_json::to_string(&manifest).expect("serialize manifest")); let _ = stdout.flush(); } @@ -160,9 +157,7 @@ fn env_requirements() -> Vec { vec![ EnvRequirement { name: "DATABASE_URL".to_string(), - description: Some( - "Postgres connection URL (e.g. postgres://user:pass@host:5432/dbname).".to_string(), - ), + description: Some("Postgres connection URL (e.g. postgres://user:pass@host:5432/dbname).".to_string()), sensitive: true, required: false, }, @@ -175,7 +170,7 @@ fn env_requirements() -> Vec { EnvRequirement { name: "ANIMUS_QUEUE_LEASE_TTL_SECS".to_string(), description: Some( - "Lease TTL in seconds; an expired lease is reclaimable (default 1800).".to_string(), + "Maximum lease TTL in seconds; v2 expiry requires explicit fenced recovery (default 1800).".to_string(), ), sensitive: false, required: false, @@ -189,8 +184,7 @@ fn env_requirements() -> Vec { EnvRequirement { name: "TOKIO_WORKER_THREADS".to_string(), description: Some( - "Caps the plugin's tokio worker threads (set by the daemon to avoid PID exhaustion)." - .to_string(), + "Caps the plugin's tokio worker threads (set by the daemon to avoid PID exhaustion).".to_string(), ), sensitive: false, required: false, @@ -212,6 +206,12 @@ fn queue_methods() -> Vec<&'static str> { METHOD_QUEUE_REORDER, METHOD_QUEUE_MARK_ASSIGNED, METHOD_QUEUE_COMPLETION, + METHOD_QUEUE_ENQUEUE_V2, + METHOD_QUEUE_LEASE_V2, + METHOD_QUEUE_LEASE_RENEW, + METHOD_QUEUE_LEASE_RECOVER, + METHOD_QUEUE_COMPLETION_V2, + METHOD_QUEUE_RELEASE_PENDING_V2, "health/check", ] } @@ -237,15 +237,17 @@ async fn handle_request( METHOD_QUEUE_NEXT_DEADLINE => Some(handle_next_deadline(id, &backend).await), METHOD_QUEUE_HOLD => Some(handle_hold(id, request.params, &backend).await), METHOD_QUEUE_RELEASE => Some(handle_release(id, request.params, &backend).await), - METHOD_QUEUE_RELEASE_PENDING => { - Some(handle_release_pending(id, request.params, &backend).await) - } + METHOD_QUEUE_RELEASE_PENDING => Some(handle_release_pending(id, request.params, &backend).await), METHOD_QUEUE_DROP => Some(handle_drop(id, request.params, &backend).await), METHOD_QUEUE_REORDER => Some(handle_reorder(id, request.params, &backend).await), - METHOD_QUEUE_MARK_ASSIGNED => { - Some(handle_mark_assigned(id, request.params, &backend).await) - } + METHOD_QUEUE_MARK_ASSIGNED => Some(handle_mark_assigned(id, request.params, &backend).await), METHOD_QUEUE_COMPLETION => Some(handle_completion(id, request.params, &backend).await), + METHOD_QUEUE_ENQUEUE_V2 => Some(handle_enqueue_v2(id, request.params, &backend).await), + METHOD_QUEUE_LEASE_V2 => Some(handle_lease_v2(id, request.params, &backend).await), + METHOD_QUEUE_LEASE_RENEW => Some(handle_renew_v2(id, request.params, &backend).await), + METHOD_QUEUE_LEASE_RECOVER => Some(handle_recover_v2(id, request.params, &backend).await), + METHOD_QUEUE_COMPLETION_V2 => Some(handle_completion_v2(id, request.params, &backend).await), + METHOD_QUEUE_RELEASE_PENDING_V2 => Some(handle_release_pending_v2(id, request.params, &backend).await), other => Some(RpcResponse::err( id, RpcError { @@ -272,27 +274,17 @@ async fn write_frame(stdout: &Arc> async fn health_check(id: Option, backend: &Arc>>) -> RpcResponse { let (status, last_error) = match backend.read().await.as_ref() { - Some(store) => match store.ping().await { + Some(store) => match store.ready().await { Ok(()) => (HealthStatus::Healthy, None), - Err(error) => ( - HealthStatus::Unhealthy, - Some(format!("Postgres unreachable: {error}")), - ), + Err(error) => (HealthStatus::Unhealthy, Some(format!("Postgres unreachable: {error}"))), }, // Not yet initialized: the process is up but has no pool to probe. None => (HealthStatus::Healthy, None), }; - let result = HealthCheckResult { - status, - uptime_ms: None, - memory_usage_bytes: None, - last_error, - }; + let result = HealthCheckResult { status, uptime_ms: None, memory_usage_bytes: None, last_error }; match serde_json::to_value(result) { Ok(value) => RpcResponse::ok(id, value), - Err(error) => { - internal_error_response(id, format!("failed to encode health result: {error}")) - } + Err(error) => internal_error_response(id, format!("failed to encode health result: {error}")), } } @@ -337,18 +329,15 @@ async fn handle_initialize(id: Option, backend: &Arc // `priority` column is stored for observability but does not weight // dispatch order. priority_weighted: false, - // No backend-side cap on lease batch size. - max_lease_batch: u32::MAX, + // One daemon tick can fill the entire five-slot coding fleet without + // locking an unbounded number of queue rows. + max_lease_batch: MAX_GENERATION_FENCED_LEASE_BATCH as u32, + generation_fenced_leases_v1: true, }; let extra = serde_json::to_value(capabilities).unwrap_or(Value::Null); let mut kind_capabilities = std::collections::HashMap::new(); - kind_capabilities.insert( - KIND.to_string(), - KindCapability { - crate_version: QUEUE_PROTOCOL_VERSION.to_string(), - extra, - }, - ); + kind_capabilities + .insert(KIND.to_string(), KindCapability { crate_version: QUEUE_PROTOCOL_VERSION.to_string(), extra }); let result = InitializeResult { protocol_version: PROTOCOL_VERSION.to_string(), @@ -356,13 +345,11 @@ async fn handle_initialize(id: Option, backend: &Arc name: PLUGIN_NAME.to_string(), version: PLUGIN_VERSION.to_string(), plugin_kind: PLUGIN_KIND_QUEUE.to_string(), + plugin_kinds: Vec::new(), description: Some(PLUGIN_DESCRIPTION.to_string()), }, capabilities: PluginCapabilities { - methods: queue_methods() - .into_iter() - .map(ToString::to_string) - .collect(), + methods: queue_methods().into_iter().map(ToString::to_string).collect(), streaming: false, progress: false, cancellation: false, @@ -375,9 +362,7 @@ async fn handle_initialize(id: Option, backend: &Arc match serde_json::to_value(result) { Ok(value) => RpcResponse::ok(id, value), - Err(error) => { - internal_error_response(id, format!("failed to encode initialize result: {error}")) - } + Err(error) => internal_error_response(id, format!("failed to encode initialize result: {error}")), } } @@ -385,11 +370,7 @@ async fn handle_initialize(id: Option, backend: &Arc // queue/* handlers // ============================================================ -async fn handle_enqueue( - id: Option, - params: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_enqueue(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -398,14 +379,7 @@ async fn handle_enqueue( Ok(req) => req, Err(response) => return response, }; - match store - .enqueue( - request.subject_dispatch, - request.run_at, - request.expire_after_secs, - ) - .await - { + match store.enqueue(request.subject_dispatch, request.run_at, request.expire_after_secs).await { Ok(outcome) => to_value_response( id, &QueueEnqueueResponse { @@ -419,11 +393,29 @@ async fn handle_enqueue( } } -async fn handle_list( +async fn handle_enqueue_v2( id: Option, params: Option, backend: &Arc>>, ) -> RpcResponse { + let store = match require_backend(id.clone(), backend).await { + Ok(store) => store, + Err(response) => return response, + }; + let request: QueueEnqueueV2Request = match parse_params(id.clone(), params, METHOD_QUEUE_ENQUEUE_V2) { + Ok(request) => request, + Err(response) => return response, + }; + if let Err(error) = request.validate() { + return RpcResponse::err(id, invalid_params(error)); + } + match store.enqueue_v2(request).await { + Ok(response) => to_value_response(id, &response), + Err(error) => internal_error_response(id, format!("{METHOD_QUEUE_ENQUEUE_V2} failed: {error:#}")), + } +} + +async fn handle_list(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -432,28 +424,18 @@ async fn handle_list( Some(value) => match serde_json::from_value(value) { Ok(req) => req, Err(error) => { - return RpcResponse::err( - id, - invalid_params(format!("invalid queue/list params: {error}")), - ); + return RpcResponse::err(id, invalid_params(format!("invalid queue/list params: {error}"))); } }, None => QueueListRequest::default(), }; - match store - .list(&request.status, request.limit, request.offset) - .await - { + match store.list(&request.status, request.limit, request.offset).await { Ok(response) => to_value_response(id, &response), Err(error) => internal_error_response(id, format!("queue/list failed: {error:#}")), } } -async fn handle_lease( - id: Option, - params: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_lease(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -462,13 +444,8 @@ async fn handle_lease( Ok(req) => req, Err(response) => return response, }; - let exclude_subjects = request - .exclude_subjects - .map(|ids| ids.into_iter().map(|id| id.0).collect::>()); - match store - .lease(request.max, request.workflow_ids, exclude_subjects) - .await - { + let exclude_subjects = request.exclude_subjects.map(|ids| ids.into_iter().map(|id| id.0).collect::>()); + match store.lease(request.max, request.workflow_ids, exclude_subjects).await { Ok(response) => to_value_response(id, &response), Err(QueueLeaseError::WorkflowIdCountMismatch { expected, actual }) => RpcResponse::err( id, @@ -478,9 +455,73 @@ async fn handle_lease( data: Some(json!({ "expected": expected, "actual": actual })), }, ), - Err(QueueLeaseError::Backend(error)) => { - internal_error_response(id, format!("queue/lease failed: {error:#}")) - } + Err(QueueLeaseError::Backend(error)) => internal_error_response(id, format!("queue/lease failed: {error:#}")), + } +} + +async fn handle_lease_v2( + id: Option, + params: Option, + backend: &Arc>>, +) -> RpcResponse { + let store = match require_backend(id.clone(), backend).await { + Ok(store) => store, + Err(response) => return response, + }; + let request: QueueLeaseV2Request = match parse_params(id.clone(), params, METHOD_QUEUE_LEASE_V2) { + Ok(request) => request, + Err(response) => return response, + }; + if let Err(error) = request.validate() { + return RpcResponse::err(id, invalid_params(error)); + } + match store.lease_v2(request).await { + Ok(response) => to_value_response(id, &response), + Err(error) => internal_error_response(id, format!("{METHOD_QUEUE_LEASE_V2} failed: {error:#}")), + } +} + +async fn handle_renew_v2( + id: Option, + params: Option, + backend: &Arc>>, +) -> RpcResponse { + let store = match require_backend(id.clone(), backend).await { + Ok(store) => store, + Err(response) => return response, + }; + let request: QueueLeaseRenewRequest = match parse_params(id.clone(), params, METHOD_QUEUE_LEASE_RENEW) { + Ok(request) => request, + Err(response) => return response, + }; + if let Err(error) = request.validate() { + return RpcResponse::err(id, invalid_params(error)); + } + match store.renew_v2(request).await { + Ok(response) => to_value_response(id, &response), + Err(error) => internal_error_response(id, format!("{METHOD_QUEUE_LEASE_RENEW} failed: {error:#}")), + } +} + +async fn handle_recover_v2( + id: Option, + params: Option, + backend: &Arc>>, +) -> RpcResponse { + let store = match require_backend(id.clone(), backend).await { + Ok(store) => store, + Err(response) => return response, + }; + let request: QueueLeaseRecoverRequest = match parse_params(id.clone(), params, METHOD_QUEUE_LEASE_RECOVER) { + Ok(request) => request, + Err(response) => return response, + }; + if let Err(error) = request.validate() { + return RpcResponse::err(id, invalid_params(error)); + } + match store.recover_v2(request).await { + Ok(response) => to_value_response(id, &response), + Err(error) => internal_error_response(id, format!("{METHOD_QUEUE_LEASE_RECOVER} failed: {error:#}")), } } @@ -495,10 +536,7 @@ async fn handle_stats(id: Option, backend: &Arc>>) - } } -async fn handle_next_deadline( - id: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_next_deadline(id: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -509,11 +547,7 @@ async fn handle_next_deadline( } } -async fn handle_hold( - id: Option, - params: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_hold(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -528,11 +562,7 @@ async fn handle_hold( } } -async fn handle_release( - id: Option, - params: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_release(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -556,15 +586,11 @@ async fn handle_release_pending( Ok(s) => s, Err(response) => return response, }; - let request: QueueReleasePendingParams = - match parse_params(id.clone(), params, "queue/release_pending") { - Ok(req) => req, - Err(response) => return response, - }; - match store - .release_pending(&request.entry_id, &request.reason) - .await - { + let request: QueueReleasePendingParams = match parse_params(id.clone(), params, "queue/release_pending") { + Ok(req) => req, + Err(response) => return response, + }; + match store.release_pending(&request.entry_id, &request.reason).await { Ok(response) => to_value_response(id, &response), // Mirror queue-default: missing entry → -32602 invalid_params. Err(QueueReleasePendingError::NotFound { entry_id }) => RpcResponse::err( @@ -575,16 +601,11 @@ async fn handle_release_pending( data: None, }, ), - Err(QueueReleasePendingError::NotAssigned { - entry_id, - actual_state, - }) => RpcResponse::err( + Err(QueueReleasePendingError::NotAssigned { entry_id, actual_state }) => RpcResponse::err( id, RpcError { code: queue_error_codes::QUEUE_ENTRY_NOT_ASSIGNED, - message: format!( - "entry {entry_id} is in state '{actual_state}', expected 'assigned'" - ), + message: format!("entry {entry_id} is in state '{actual_state}', expected 'assigned'"), data: Some(json!({ "actual_state": actual_state })), }, ), @@ -594,11 +615,7 @@ async fn handle_release_pending( } } -async fn handle_drop( - id: Option, - params: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_drop(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -613,11 +630,7 @@ async fn handle_drop( } } -async fn handle_reorder( - id: Option, - params: Option, - backend: &Arc>>, -) -> RpcResponse { +async fn handle_reorder(id: Option, params: Option, backend: &Arc>>) -> RpcResponse { let store = match require_backend(id.clone(), backend).await { Ok(s) => s, Err(response) => return response, @@ -648,15 +661,11 @@ async fn handle_mark_assigned( Ok(s) => s, Err(response) => return response, }; - let request: QueueMarkAssignedRequest = - match parse_params(id.clone(), params, "queue/mark_assigned") { - Ok(req) => req, - Err(response) => return response, - }; - match store - .mark_assigned(&request.entry_id, request.workflow_id) - .await - { + let request: QueueMarkAssignedRequest = match parse_params(id.clone(), params, "queue/mark_assigned") { + Ok(req) => req, + Err(response) => return response, + }; + match store.mark_assigned(&request.entry_id, request.workflow_id).await { Ok(response) => to_value_response(id, &response), Err(error) => not_pending_or_internal(id, &error, "queue/mark_assigned"), } @@ -671,18 +680,12 @@ async fn handle_completion( Ok(s) => s, Err(response) => return response, }; - let request: QueueCompletionRequest = match parse_params(id.clone(), params, "queue/completion") - { + let request: QueueCompletionRequest = match parse_params(id.clone(), params, "queue/completion") { Ok(req) => req, Err(response) => return response, }; match store - .completion( - &request.entry_id, - &request.status, - request.workflow_ref.as_deref(), - request.workflow_id.as_deref(), - ) + .completion(&request.entry_id, &request.status, request.workflow_ref.as_deref(), request.workflow_id.as_deref()) .await { Ok(response) => to_value_response(id, &response), @@ -697,6 +700,51 @@ async fn handle_completion( } } +async fn handle_completion_v2( + id: Option, + params: Option, + backend: &Arc>>, +) -> RpcResponse { + let store = match require_backend(id.clone(), backend).await { + Ok(store) => store, + Err(response) => return response, + }; + let request: QueueCompletionV2Request = match parse_params(id.clone(), params, METHOD_QUEUE_COMPLETION_V2) { + Ok(request) => request, + Err(response) => return response, + }; + if let Err(error) = request.validate() { + return RpcResponse::err(id, invalid_params(error)); + } + match store.completion_v2(request).await { + Ok(response) => to_value_response(id, &response), + Err(error) => internal_error_response(id, format!("{METHOD_QUEUE_COMPLETION_V2} failed: {error:#}")), + } +} + +async fn handle_release_pending_v2( + id: Option, + params: Option, + backend: &Arc>>, +) -> RpcResponse { + let store = match require_backend(id.clone(), backend).await { + Ok(store) => store, + Err(response) => return response, + }; + let request: QueueReleasePendingV2Request = match parse_params(id.clone(), params, METHOD_QUEUE_RELEASE_PENDING_V2) + { + Ok(request) => request, + Err(response) => return response, + }; + if let Err(error) = request.validate() { + return RpcResponse::err(id, invalid_params(error)); + } + match store.release_pending_v2(request).await { + Ok(response) => to_value_response(id, &response), + Err(error) => internal_error_response(id, format!("{METHOD_QUEUE_RELEASE_PENDING_V2} failed: {error:#}")), + } +} + // ============================================================ // helpers // ============================================================ @@ -723,18 +771,10 @@ fn parse_params( params: Option, method: &str, ) -> std::result::Result { - let value = params.ok_or_else(|| { - RpcResponse::err( - id.clone(), - invalid_params(format!("missing params for {method}")), - ) - })?; - serde_json::from_value::(value).map_err(|error| { - RpcResponse::err( - id, - invalid_params(format!("invalid {method} params: {error}")), - ) - }) + let value = + params.ok_or_else(|| RpcResponse::err(id.clone(), invalid_params(format!("missing params for {method}"))))?; + serde_json::from_value::(value) + .map_err(|error| RpcResponse::err(id, invalid_params(format!("invalid {method} params: {error}")))) } fn to_value_response(id: Option, value: &T) -> RpcResponse { @@ -749,31 +789,16 @@ fn not_pending_or_internal(id: Option, error: &anyhow::Error, method: &st if msg.contains("not in the expected pre-mutation status") { return RpcResponse::err( id, - RpcError { - code: queue_error_codes::QUEUE_ENTRY_NOT_PENDING, - message: msg, - data: None, - }, + RpcError { code: queue_error_codes::QUEUE_ENTRY_NOT_PENDING, message: msg, data: None }, ); } internal_error_response(id, format!("{method} failed: {error:#}")) } fn invalid_params(message: impl Into) -> RpcError { - RpcError { - code: plugin_error_codes::INVALID_PARAMS, - message: message.into(), - data: None, - } + RpcError { code: plugin_error_codes::INVALID_PARAMS, message: message.into(), data: None } } fn internal_error_response(id: Option, message: impl Into) -> RpcResponse { - RpcResponse::err( - id, - RpcError { - code: plugin_error_codes::INTERNAL_ERROR, - message: message.into(), - data: None, - }, - ) + RpcResponse::err(id, RpcError { code: plugin_error_codes::INTERNAL_ERROR, message: message.into(), data: None }) } diff --git a/src/store/mod.rs b/src/store/mod.rs index 396a871..3269a80 100644 --- a/src/store/mod.rs +++ b/src/store/mod.rs @@ -9,7 +9,11 @@ //! - `queue/lease` runs inside a single transaction and claims rows with //! `SELECT ... FOR UPDATE SKIP LOCKED`, so two daemons (or a daemon and a //! manual CLI lease) never hand out the same entry twice. -//! - **Lease-expiry reclaim**: the lease-eligible set is `state = 'pending'` +//! - Legacy `queue/lease` retains its historical lease-expiry reclaim behavior. +//! - Generation-fenced `queue/v2/lease` only considers Pending rows. Expired +//! assignments require an explicit compare-and-swap recovery, preserving the +//! workflow and subject generation rather than starting duplicate work. +//! - **Legacy lease-expiry reclaim**: the lease-eligible set is `state = 'pending'` //! (and due) OR `state = 'leased' AND lease_expires_at < now()`. A crashed or //! redeployed daemon's stale lease is therefore re-dispatched on the next //! `queue/lease` instead of being lost. @@ -17,10 +21,16 @@ //! `completion` / `release_pending`) lock the target row `FOR UPDATE` inside //! a transaction for read-modify-write safety. +use animus_execution_protocol::{ + ExecutionFence, QueueLeaseFence, RepositoryReservation, SubjectGeneration, EXECUTION_FENCE_SCHEMA_ID, + EXECUTION_FENCE_VERSION, +}; use animus_queue_protocol::{ - completion_status, status, QueueEntry, QueueLeaseResponse, QueueListResponse, - QueueMutationResponse, QueueNextDeadlineResponse, QueueReleasePendingResponse, - QueueReorderResponse, QueueStats, + completion_status, status, FencedQueueEntry, QueueCompletionV2Request, QueueEnqueueV2Request, + QueueEnqueueV2Response, QueueEntry, QueueLeaseBlock, QueueLeaseBlockReason, QueueLeaseMutationOutcome, + QueueLeaseMutationResponse, QueueLeaseRecoverRequest, QueueLeaseRenewRequest, QueueLeaseResponse, + QueueLeaseV2Request, QueueLeaseV2Response, QueueListResponse, QueueMutationResponse, QueueNextDeadlineResponse, + QueueReleasePendingResponse, QueueReleasePendingV2Request, QueueReorderResponse, QueueStats, }; use animus_subject_protocol::SubjectDispatch; use anyhow::{Context, Result}; @@ -28,8 +38,10 @@ use chrono::{DateTime, Utc}; use serde_json::{json, Value}; use sqlx::postgres::{PgPool, PgPoolOptions}; use sqlx::Row; +use std::sync::Arc; +use tokio::sync::OnceCell; -use crate::config::QueueConfig; +use crate::{config::QueueConfig, MAX_GENERATION_FENCED_LEASE_BATCH}; /// Durable `state` column values. The first three map 1:1 onto the wire /// `status` vocabulary; `done` / `dropped` are terminal soft-delete states @@ -87,6 +99,7 @@ pub struct Store { pool: PgPool, table: String, lease_ttl_secs: i64, + schema_ready: Arc>, } impl Store { @@ -106,14 +119,11 @@ impl Store { pool, table: config.table.clone(), lease_ttl_secs: config.lease_ttl_secs, + schema_ready: Arc::new(OnceCell::new()), }; - let migrate_store = Self { - pool: store.pool.clone(), - table: store.table.clone(), - lease_ttl_secs: store.lease_ttl_secs, - }; + let migrate_store = store.clone(); tokio::spawn(async move { - if let Err(error) = migrate_store.migrate().await { + if let Err(error) = migrate_store.ensure_migrated().await { eprintln!("[animus-queue-postgres] background migrate failed (schema likely present; retried next spawn): {error:#}"); } }); @@ -126,6 +136,7 @@ impl Store { pool, table: config.table.clone(), lease_ttl_secs: config.lease_ttl_secs, + schema_ready: Arc::new(OnceCell::new()), } } @@ -133,6 +144,14 @@ impl Store { format!("{}_ordinal_seq", self.table) } + fn generation_table(&self) -> String { + format!("{}_subject_generation", self.table) + } + + async fn ensure_migrated(&self) -> Result<()> { + self.schema_ready.get_or_try_init(|| async { self.migrate().await }).await.map(|_| ()) + } + /// Idempotent schema migration. Safe to run on every boot. pub async fn migrate(&self) -> Result<()> { let t = &self.table; @@ -169,33 +188,95 @@ impl Store { .await .context("failed to create queue_item table")?; - // Dispatch hot path: lease/list scan live rows in ordinal order. + for column in [ + "fence_version int", + "idempotency_key text", + "subject_qualified_id text", + "subject_generation bigint", + "workflow_generation bigint", + "lease_generation bigint", + "repository text", + "base_ref text", + "head_ref text", + "terminal_status text", + ] { + sqlx::query(&format!("ALTER TABLE {t} ADD COLUMN IF NOT EXISTS {column}")) + .execute(&self.pool) + .await + .with_context(|| format!("failed to add generation-fence column {column}"))?; + } + + let generations = self.generation_table(); sqlx::query(&format!( - "CREATE INDEX IF NOT EXISTS {t}_state_ordinal_idx ON {t}(state, ordinal)" + "CREATE TABLE IF NOT EXISTS {generations} ( \ + subject_qualified_id text PRIMARY KEY, \ + last_generation bigint NOT NULL CHECK (last_generation > 0) \ + )" )) .execute(&self.pool) .await - .context("failed to create state/ordinal index")?; + .context("failed to create subject-generation allocator")?; + + // Dispatch hot path: lease/list scan live rows in ordinal order. + sqlx::query(&format!("CREATE INDEX IF NOT EXISTS {t}_state_ordinal_idx ON {t}(state, ordinal)")) + .execute(&self.pool) + .await + .context("failed to create state/ordinal index")?; // Subject collision + exclude_subjects lookups. + sqlx::query(&format!("CREATE INDEX IF NOT EXISTS {t}_subject_idx ON {t}(subject_id)")) + .execute(&self.pool) + .await + .context("failed to create subject index")?; + sqlx::query(&format!( + "CREATE UNIQUE INDEX IF NOT EXISTS {t}_idempotency_v2_uidx \ + ON {t}(idempotency_key) WHERE idempotency_key IS NOT NULL" + )) + .execute(&self.pool) + .await + .context("failed to create v2 idempotency index")?; + sqlx::query(&format!( + "CREATE UNIQUE INDEX IF NOT EXISTS {t}_subject_generation_v2_uidx \ + ON {t}(subject_qualified_id, subject_generation) \ + WHERE fence_version = 1" + )) + .execute(&self.pool) + .await + .context("failed to create v2 subject-generation index")?; sqlx::query(&format!( - "CREATE INDEX IF NOT EXISTS {t}_subject_idx ON {t}(subject_id)" + "CREATE UNIQUE INDEX IF NOT EXISTS {t}_live_subject_v2_uidx \ + ON {t}(subject_qualified_id) \ + WHERE fence_version = 1 AND state IN ('pending', 'leased', 'held')" )) .execute(&self.pool) .await - .context("failed to create subject index")?; + .context("failed to create v2 live-subject ownership index")?; + sqlx::query(&format!( + "CREATE UNIQUE INDEX IF NOT EXISTS {t}_repository_head_v2_uidx \ + ON {t}(lower(repository), head_ref) \ + WHERE fence_version = 1 AND repository IS NOT NULL \ + AND state IN ('pending', 'leased', 'held')" + )) + .execute(&self.pool) + .await + .context("failed to create v2 repository reservation index")?; Ok(()) } /// Health probe: a trivial `SELECT 1`. pub async fn ping(&self) -> Result<()> { - sqlx::query("SELECT 1") - .execute(&self.pool) - .await - .context("Postgres ping failed")?; + sqlx::query("SELECT 1").execute(&self.pool).await.context("Postgres ping failed")?; Ok(()) } + /// Readiness probe: database connectivity plus successful application of + /// the generation-fence schema. A plugin must not advertise healthy while + /// its v2 capability cannot actually be served. + pub async fn ready(&self) -> Result<()> { + self.ensure_migrated().await?; + self.ping().await + } + // ============================================================ // queue/enqueue // ============================================================ @@ -212,8 +293,8 @@ impl Store { // advisory reflects the live queue. self.sweep_expired().await?; - let subject_key = dispatch.subject_key(); - let subject_kind = dispatch.subject_kind().to_string(); + let subject_key = dispatch.subject_key().context("queue entries require a subject identity")?; + let subject_kind = dispatch.subject_kind().context("queue entries require a subject kind")?.to_string(); let workflow_ref = dispatch.workflow_ref.clone(); let priority = priority_to_int(dispatch.priority.as_deref()); @@ -263,11 +344,175 @@ impl Store { .await .context("failed to insert queue entry")?; - Ok(EnqueueOutcome { + Ok(EnqueueOutcome { enqueued: true, entry_id, subject_id: subject_key, warning }) + } + + /// Idempotently enqueue a subject and allocate its next immutable + /// generation. Repository/head reservations are write-once across all live + /// v2 entries, preventing two coding slots from owning one branch. + pub async fn enqueue_v2(&self, request: QueueEnqueueV2Request) -> Result { + request.validate().map_err(anyhow::Error::msg)?; + let run_at_ts = match request.run_at.as_deref() { + Some(raw) => Some(parse_rfc3339(raw).context("v2 run_at must be RFC 3339")?), + None => None, + }; + self.ensure_migrated().await?; + self.sweep_expired().await?; + + let subject_ref = request.subject_dispatch.subject().context("generation-fenced enqueue requires a subject")?; + let qualified_id = format!("{}:{}", subject_ref.kind(), subject_ref.id()); + let subject_key = + request.subject_dispatch.subject_key().context("generation-fenced enqueue requires a subject key")?; + let subject_kind = subject_ref.kind().to_string(); + let workflow_ref = request.subject_dispatch.workflow_ref.clone(); + let priority = priority_to_int(request.subject_dispatch.priority.as_deref()); + let payload = + serde_json::to_value(&request.subject_dispatch).context("failed to encode generation-fenced dispatch")?; + let expire_after = request.expire_after_secs.map(|seconds| seconds as i64); + let mut tx = self.pool.begin().await.context("failed to begin v2 enqueue tx")?; + + if let Some(key) = request.idempotency_key.as_deref() { + sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended($1, 1175))") + .bind(key) + .execute(&mut *tx) + .await + .context("failed to lock v2 idempotency key")?; + if let Some(row) = sqlx::query(&format!( + "SELECT id, subject_qualified_id, subject_generation, workflow_ref, \ + repository, base_ref, head_ref FROM {} \ + WHERE idempotency_key = $1", + self.table + )) + .bind(key) + .fetch_optional(&mut *tx) + .await + .context("failed to load idempotent queue entry")? + { + let existing_qualified: Option = row.get("subject_qualified_id"); + let existing_generation: Option = row.get("subject_generation"); + if existing_qualified.as_deref() != Some(qualified_id.as_str()) { + anyhow::bail!("idempotency key already belongs to a different subject"); + } + let existing_workflow_ref: String = row.get("workflow_ref"); + if existing_workflow_ref != workflow_ref { + anyhow::bail!("idempotency key already belongs to a different workflow_ref"); + } + let existing_repository = row_to_repository(&row)?; + let requested_repository = request.repository.as_ref().map(|reservation| RepositoryReservation { + repository: reservation.repository.trim().to_ascii_lowercase(), + base_ref: reservation.base_ref.clone(), + head_ref: reservation.head_ref.clone(), + }); + if existing_repository != requested_repository { + anyhow::bail!("idempotency key already belongs to a different repository reservation"); + } + tx.commit().await.context("failed to commit v2 enqueue replay")?; + return Ok(QueueEnqueueV2Response { + enqueued: false, + entry_id: row.get("id"), + subject: SubjectGeneration { + qualified_id, + generation: positive_u64(existing_generation, "subject_generation")?, + }, + warning: None, + }); + } + } + + sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended($1, 1176))") + .bind(&qualified_id) + .execute(&mut *tx) + .await + .context("failed to lock qualified subject for v2 enqueue")?; + if let Some(row) = sqlx::query(&format!( + "SELECT id, subject_generation FROM {} \ + WHERE subject_qualified_id = $1 AND state = ANY($2) \ + ORDER BY subject_generation DESC LIMIT 1", + self.table + )) + .bind(&qualified_id) + .bind(&db_state::LIVE[..]) + .fetch_optional(&mut *tx) + .await + .context("failed to check live subject ownership")? + { + let generation = positive_u64(row.get("subject_generation"), "subject_generation")?; + tx.commit().await.context("failed to commit v2 subject collision")?; + return Ok(QueueEnqueueV2Response { + enqueued: false, + entry_id: row.get("id"), + subject: SubjectGeneration { qualified_id: qualified_id.clone(), generation }, + warning: Some(format!("subject {qualified_id} already has an active generation; enqueue rejected")), + }); + } + + let generation: i64 = sqlx::query_scalar(&format!( + "INSERT INTO {generations} (subject_qualified_id, last_generation) VALUES ($1, 1) \ + ON CONFLICT (subject_qualified_id) DO UPDATE \ + SET last_generation = {generations}.last_generation + 1 \ + RETURNING last_generation", + generations = self.generation_table(), + )) + .bind(&qualified_id) + .fetch_one(&mut *tx) + .await + .context("failed to allocate subject generation")?; + + let entry_id = uuid::Uuid::new_v4().to_string(); + let (repository, base_ref, head_ref) = request + .repository + .as_ref() + .map(|reservation| { + ( + Some(reservation.repository.trim().to_ascii_lowercase()), + Some(reservation.base_ref.clone()), + Some(reservation.head_ref.clone()), + ) + }) + .unwrap_or((None, None, None)); + let insert = sqlx::query(&format!( + "INSERT INTO {t} ( \ + id, subject_kind, subject_id, workflow_ref, state, priority, ordinal, \ + enqueued_at, run_at, expire_after_secs, updated_at, payload, fence_version, \ + idempotency_key, subject_qualified_id, subject_generation, \ + repository, base_ref, head_ref \ + ) VALUES ( \ + $1, $2, $3, $4, '{pending}', $5, nextval('{seq}'), \ + now(), $6, $7, now(), $8, 1, $9, $10, $11, $12, $13, $14 \ + )", + t = self.table, + pending = db_state::PENDING, + seq = self.seq(), + )) + .bind(&entry_id) + .bind(&subject_kind) + .bind(&subject_key) + .bind(&workflow_ref) + .bind(priority) + .bind(run_at_ts) + .bind(expire_after) + .bind(&payload) + .bind(request.idempotency_key.as_deref()) + .bind(&qualified_id) + .bind(generation) + .bind(repository) + .bind(base_ref) + .bind(head_ref) + .execute(&mut *tx) + .await; + if let Err(error) = insert { + if error.as_database_error().is_some_and(|db| db.is_unique_violation()) { + anyhow::bail!("generation-fenced queue ownership collision: {error}"); + } + return Err(anyhow::Error::from(error).context("failed to insert v2 queue entry")); + } + tx.commit().await.context("failed to commit v2 enqueue")?; + + Ok(QueueEnqueueV2Response { enqueued: true, entry_id, - subject_id: subject_key, - warning, + subject: SubjectGeneration { qualified_id, generation: generation as u64 }, + warning: None, }) } @@ -301,11 +546,7 @@ impl Store { filtered.truncate(limit); } - Ok(QueueListResponse { - entries: filtered, - total, - stats, - }) + Ok(QueueListResponse { entries: filtered, total, stats }) } /// Aggregate counts over live rows. @@ -318,6 +559,7 @@ impl Store { /// (ordinal ASC), mapped to wire [`QueueEntry`]s. Rows whose payload can no /// longer decode into a `SubjectDispatch` are logged and skipped. async fn fetch_live_entries(&self) -> Result> { + self.ensure_migrated().await?; let rows = sqlx::query(&format!( "SELECT id, subject_id, state, workflow_id, enqueued_at, assigned_at, held_at, \ run_at, expire_after_secs, payload \ @@ -347,23 +589,20 @@ impl Store { ) -> std::result::Result { if let Some(ids) = workflow_ids.as_ref() { if ids.len() != max { - return Err(QueueLeaseError::WorkflowIdCountMismatch { - expected: max, - actual: ids.len(), - }); + return Err(QueueLeaseError::WorkflowIdCountMismatch { expected: max, actual: ids.len() }); } } if max == 0 { return Ok(QueueLeaseResponse { leased: Vec::new() }); } - self.sweep_expired() - .await - .map_err(QueueLeaseError::Backend)?; + self.sweep_expired().await.map_err(QueueLeaseError::Backend)?; - let mut tx = self.pool.begin().await.map_err(|e| { - QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to begin lease tx")) - })?; + let mut tx = self + .pool + .begin() + .await + .map_err(|e| QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to begin lease tx")))?; // Lock every lease-eligible row in dispatch order. SKIP LOCKED lets a // concurrent lease proceed on rows we don't hold. The reclaim arm @@ -374,7 +613,7 @@ impl Store { WHERE ( \ (state = '{pending}' AND (run_at IS NULL OR run_at <= now())) \ OR (state = '{leased}' AND lease_expires_at IS NOT NULL AND lease_expires_at < now()) \ - ) \ + ) AND COALESCE(fence_version, 0) <> 1 \ ORDER BY ordinal ASC \ FOR UPDATE SKIP LOCKED", t = self.table, @@ -383,12 +622,9 @@ impl Store { )) .fetch_all(&mut *tx) .await - .map_err(|e| { - QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to select lease candidates")) - })?; + .map_err(|e| QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to select lease candidates")))?; - let mut exclude_set: std::collections::HashSet = - exclude_subjects.into_iter().flatten().collect(); + let mut exclude_set: std::collections::HashSet = exclude_subjects.into_iter().flatten().collect(); let mut chosen: Vec<(String, String)> = Vec::new(); // (entry_id, subject_key) for row in &candidate_rows { @@ -434,23 +670,416 @@ impl Store { .bind(entry_id) .fetch_one(&mut *tx) .await - .map_err(|e| { - QueueLeaseError::Backend( - anyhow::Error::from(e).context("failed to mark entry leased"), - ) - })?; + .map_err(|e| QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to mark entry leased")))?; if let Some(entry) = row_to_entry(&row) { leased.push(entry); } } - tx.commit().await.map_err(|e| { - QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to commit lease tx")) - })?; + tx.commit() + .await + .map_err(|e| QueueLeaseError::Backend(anyhow::Error::from(e).context("failed to commit lease tx")))?; Ok(QueueLeaseResponse { leased }) } + /// Lease only fresh Pending v2 entries. Expired Assigned rows are surfaced + /// as recovery-required and remain untouched until an exact CAS recovery. + pub async fn lease_v2(&self, request: QueueLeaseV2Request) -> Result { + request.validate().map_err(anyhow::Error::msg)?; + if request.max > MAX_GENERATION_FENCED_LEASE_BATCH { + anyhow::bail!( + "queue/v2/lease max {} exceeds fleet limit {}", + request.max, + MAX_GENERATION_FENCED_LEASE_BATCH + ); + } + if request.exclude.len() > MAX_GENERATION_FENCED_LEASE_BATCH { + anyhow::bail!( + "queue/v2/lease exclude count {} exceeds fleet limit {}", + request.exclude.len(), + MAX_GENERATION_FENCED_LEASE_BATCH + ); + } + self.ensure_migrated().await?; + self.sweep_expired().await?; + + let mut blocked = self.v2_recovery_blocks().await?; + let mut tx = self.pool.begin().await.context("failed to begin v2 lease tx")?; + let rows = sqlx::query(&format!( + "SELECT id, subject_id, state, workflow_id, enqueued_at, assigned_at, held_at, \ + run_at, expire_after_secs, payload, fence_version, subject_qualified_id, \ + subject_generation, workflow_generation, lease_owner, lease_generation, \ + lease_expires_at, repository, base_ref, head_ref \ + FROM {t} \ + WHERE state = '{pending}' AND fence_version = 1 \ + AND (run_at IS NULL OR run_at <= now()) \ + ORDER BY ordinal ASC LIMIT $1 FOR UPDATE SKIP LOCKED", + t = self.table, + pending = db_state::PENDING, + )) + .bind(request.max.saturating_add(request.exclude.len()) as i64) + .fetch_all(&mut *tx) + .await + .context("failed to select v2 lease candidates")?; + + let mut leased = Vec::new(); + let mut workflow_id_index = 0usize; + for row in rows { + let entry_id: String = row.get("id"); + let state: String = row.get("state"); + let fence_version: Option = row.get("fence_version"); + let stored = row_to_execution(&row); + if fence_version != Some(1) { + blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::MissingExecutionIdentity, + conflicts_with: None, + }); + continue; + } + if state == db_state::LEASED { + match stored { + Ok(execution) => blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::ExpiredLeaseRecoveryRequired, + conflicts_with: Some(execution), + }), + Err(_) => blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::MissingExecutionIdentity, + conflicts_with: None, + }), + } + continue; + } + let subject_generation: Option = row.get("subject_generation"); + let qualified_id: Option = row.get("subject_qualified_id"); + if subject_generation.is_none() || qualified_id.as_deref().is_none_or(str::is_empty) { + blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::MissingExecutionIdentity, + conflicts_with: None, + }); + continue; + } + + let subject = SubjectGeneration { + qualified_id: qualified_id.expect("checked"), + generation: positive_u64(subject_generation, "subject_generation")?, + }; + let repository = row_to_repository(&row)?; + let subject_conflict = + request.exclude.iter().find(|execution| execution.subject.as_ref() == Some(&subject)); + if let Some(conflict) = subject_conflict { + blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::SubjectGenerationActive, + conflicts_with: Some(conflict.clone()), + }); + continue; + } + let repository_conflict = repository.as_ref().and_then(|candidate| { + request.exclude.iter().find(|execution| { + execution + .repository + .as_ref() + .is_some_and(|active| active.collision_key() == candidate.collision_key()) + }) + }); + if let Some(conflict) = repository_conflict { + blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::RepositoryRefCollision, + conflicts_with: Some(conflict.clone()), + }); + continue; + } + if leased.len() == request.max { + break; + } + + let existing_workflow_id: Option = row.get("workflow_id"); + let workflow_id = match existing_workflow_id { + Some(workflow_id) => workflow_id, + None => { + let workflow_id = request.workflow_ids[workflow_id_index].clone(); + workflow_id_index += 1; + workflow_id + } + }; + let updated = sqlx::query(&format!( + "UPDATE {t} SET state = '{leased}', workflow_id = $1, \ + workflow_generation = COALESCE(workflow_generation, 1), \ + lease_owner = $2, lease_generation = COALESCE(lease_generation, 0) + 1, \ + lease_expires_at = now() + make_interval(secs => $3::int), \ + assigned_at = now(), updated_at = now() \ + WHERE id = $4 AND state = '{pending}' AND fence_version = 1 \ + RETURNING id, subject_id, state, workflow_id, enqueued_at, assigned_at, held_at, \ + run_at, expire_after_secs, payload, fence_version, subject_qualified_id, \ + subject_generation, workflow_generation, lease_owner, lease_generation, \ + lease_expires_at, repository, base_ref, head_ref", + t = self.table, + leased = db_state::LEASED, + pending = db_state::PENDING, + )) + .bind(&workflow_id) + .bind(&request.owner_id) + .bind(self.lease_ttl_secs) + .bind(&entry_id) + .fetch_one(&mut *tx) + .await + .context("failed to assign generation-fenced lease")?; + leased.push(row_to_fenced_entry(&updated)?); + } + tx.commit().await.context("failed to commit v2 lease")?; + Ok(QueueLeaseV2Response { leased, blocked }) + } + + /// Renew a still-live lease only when every ownership field matches. + pub async fn renew_v2(&self, request: QueueLeaseRenewRequest) -> Result { + request.validate().map_err(anyhow::Error::msg)?; + self.ensure_migrated().await?; + let expected = request.execution; + let entry_id = expected.queue_lease.as_ref().expect("validated queue lease").entry_id.clone(); + let ttl = clamp_ttl(request.ttl_secs, self.lease_ttl_secs); + let mut tx = self.pool.begin().await.context("failed to begin v2 renew tx")?; + let Some(row) = self.fenced_row_for_update(&mut tx, &entry_id).await? else { + return Ok(mutation_response(QueueLeaseMutationOutcome::NotFound, None, "entry not found")); + }; + let state: String = row.get("state"); + if state != db_state::LEASED { + return Ok(mutation_response(QueueLeaseMutationOutcome::NotAssigned, None, "entry is not assigned")); + } + let current = match row_to_execution(&row) { + Ok(current) => current, + Err(error) => return Ok(mutation_response(QueueLeaseMutationOutcome::StaleFence, None, error.to_string())), + }; + if current != expected { + return Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + "execution fence does not match", + )); + } + let backend_now: DateTime = row.get("backend_now"); + if current.queue_lease.as_ref().expect("queue lease").expires_at <= backend_now { + return Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + "lease expired; explicit recovery required", + )); + } + let updated = sqlx::query(&format!( + "UPDATE {t} SET lease_expires_at = now() + make_interval(secs => $1::int), updated_at = now() \ + WHERE id = $2 RETURNING id, subject_id, state, workflow_id, enqueued_at, assigned_at, \ + held_at, run_at, expire_after_secs, payload, fence_version, subject_qualified_id, \ + subject_generation, workflow_generation, lease_owner, lease_generation, \ + lease_expires_at, repository, base_ref, head_ref", + t = self.table, + )) + .bind(ttl) + .bind(&entry_id) + .fetch_one(&mut *tx) + .await + .context("failed to renew generation-fenced lease")?; + let execution = row_to_execution(&updated)?; + tx.commit().await.context("failed to commit v2 renew")?; + Ok(mutation_response(QueueLeaseMutationOutcome::Applied, Some(execution), "lease renewed")) + } + + /// Transfer one expired lease to a new owner while preserving workflow, + /// subject, and repository identity and incrementing only lease generation. + pub async fn recover_v2(&self, request: QueueLeaseRecoverRequest) -> Result { + request.validate().map_err(anyhow::Error::msg)?; + self.ensure_migrated().await?; + let expected = request.execution; + let entry_id = expected.queue_lease.as_ref().expect("validated queue lease").entry_id.clone(); + let ttl = clamp_ttl(request.ttl_secs, self.lease_ttl_secs); + let mut tx = self.pool.begin().await.context("failed to begin v2 recovery tx")?; + let Some(row) = self.fenced_row_for_update(&mut tx, &entry_id).await? else { + return Ok(mutation_response(QueueLeaseMutationOutcome::NotFound, None, "entry not found")); + }; + let state: String = row.get("state"); + if state != db_state::LEASED { + return Ok(mutation_response(QueueLeaseMutationOutcome::NotAssigned, None, "entry is not assigned")); + } + let current = match row_to_execution(&row) { + Ok(current) => current, + Err(error) => return Ok(mutation_response(QueueLeaseMutationOutcome::StaleFence, None, error.to_string())), + }; + if current != expected { + return Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + "execution fence does not match", + )); + } + let backend_now: DateTime = row.get("backend_now"); + if current.queue_lease.as_ref().expect("queue lease").expires_at > backend_now { + return Ok(mutation_response( + QueueLeaseMutationOutcome::LeaseStillLive, + Some(current), + "lease has not expired", + )); + } + let updated = sqlx::query(&format!( + "UPDATE {t} SET lease_owner = $1, lease_generation = lease_generation + 1, \ + lease_expires_at = now() + make_interval(secs => $2::int), updated_at = now() \ + WHERE id = $3 RETURNING id, subject_id, state, workflow_id, enqueued_at, assigned_at, \ + held_at, run_at, expire_after_secs, payload, fence_version, subject_qualified_id, \ + subject_generation, workflow_generation, lease_owner, lease_generation, \ + lease_expires_at, repository, base_ref, head_ref", + t = self.table, + )) + .bind(&request.new_owner_id) + .bind(ttl) + .bind(&entry_id) + .fetch_one(&mut *tx) + .await + .context("failed to recover generation-fenced lease")?; + let execution = row_to_execution(&updated)?; + tx.commit().await.context("failed to commit v2 recovery")?; + Ok(mutation_response(QueueLeaseMutationOutcome::Applied, Some(execution), "lease ownership transferred")) + } + + /// Complete only the exact execution generation. Replays from the same + /// fence are idempotent; stale owners cannot terminalize newer work. + pub async fn completion_v2(&self, request: QueueCompletionV2Request) -> Result { + request.validate().map_err(anyhow::Error::msg)?; + self.ensure_migrated().await?; + let expected = request.execution; + let entry_id = expected.queue_lease.as_ref().expect("validated queue lease").entry_id.clone(); + let mut tx = self.pool.begin().await.context("failed to begin v2 completion tx")?; + let Some(row) = self.fenced_row_for_update(&mut tx, &entry_id).await? else { + return Ok(mutation_response(QueueLeaseMutationOutcome::NotFound, None, "entry not found")); + }; + let current = match row_to_execution(&row) { + Ok(current) => current, + Err(error) => return Ok(mutation_response(QueueLeaseMutationOutcome::StaleFence, None, error.to_string())), + }; + if current != expected { + return Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + "execution fence does not match", + )); + } + if let Some(want_ref) = request.workflow_ref.as_deref() { + let payload: Value = row.get("payload"); + if payload.get("workflow_ref").and_then(Value::as_str) != Some(want_ref) { + return Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + "workflow_ref does not match", + )); + } + } + let state: String = row.get("state"); + if state == db_state::DONE { + let terminal_status: Option = row.get("terminal_status"); + return if terminal_status.as_deref() == Some(request.status.as_str()) { + Ok(mutation_response( + QueueLeaseMutationOutcome::AlreadyApplied, + Some(current), + "completion already recorded", + )) + } else { + Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + format!( + "completion outcome conflicts with stored terminal status '{}'", + terminal_status.as_deref().unwrap_or("unknown") + ), + )) + }; + } + if state != db_state::LEASED { + return Ok(mutation_response( + QueueLeaseMutationOutcome::NotAssigned, + Some(current), + "entry is not assigned", + )); + } + sqlx::query(&format!( + "UPDATE {} SET state = '{}', terminal_status = $1, updated_at = now() WHERE id = $2", + self.table, + db_state::DONE + )) + .bind(&request.status) + .bind(&entry_id) + .execute(&mut *tx) + .await + .context("failed to complete generation-fenced entry")?; + tx.commit().await.context("failed to commit v2 completion")?; + Ok(mutation_response(QueueLeaseMutationOutcome::Applied, Some(current), "completion recorded")) + } + + /// Return the exact execution to Pending without erasing its canonical + /// workflow id/generation. The next lease increments the lease generation. + pub async fn release_pending_v2( + &self, + request: QueueReleasePendingV2Request, + ) -> Result { + request.validate().map_err(anyhow::Error::msg)?; + self.ensure_migrated().await?; + let expected = request.execution; + let entry_id = expected.queue_lease.as_ref().expect("validated queue lease").entry_id.clone(); + let mut tx = self.pool.begin().await.context("failed to begin v2 release-pending tx")?; + let Some(row) = self.fenced_row_for_update(&mut tx, &entry_id).await? else { + return Ok(mutation_response(QueueLeaseMutationOutcome::NotFound, None, "entry not found")); + }; + let current = match row_to_execution(&row) { + Ok(current) => current, + Err(error) => return Ok(mutation_response(QueueLeaseMutationOutcome::StaleFence, None, error.to_string())), + }; + if current != expected { + return Ok(mutation_response( + QueueLeaseMutationOutcome::StaleFence, + Some(current), + "execution fence does not match", + )); + } + let state: String = row.get("state"); + if state == db_state::PENDING { + return Ok(mutation_response( + QueueLeaseMutationOutcome::AlreadyApplied, + Some(current), + "entry already pending", + )); + } + if state != db_state::LEASED { + return Ok(mutation_response( + QueueLeaseMutationOutcome::NotAssigned, + Some(current), + "entry is not assigned", + )); + } + let audit = json!({ + "at": Utc::now().to_rfc3339(), + "method": "queue/v2/release_pending", + "from_status": status::ASSIGNED, + "to_status": status::PENDING, + "reason": request.reason, + "workflow_id": current.workflow_id, + "workflow_generation": current.workflow_generation, + }); + sqlx::query(&format!( + "UPDATE {t} SET state = '{pending}', assigned_at = NULL, \ + audit_log = audit_log || $1::jsonb, updated_at = now() WHERE id = $2", + t = self.table, + pending = db_state::PENDING, + )) + .bind(&audit) + .bind(&entry_id) + .execute(&mut *tx) + .await + .context("failed to return generation-fenced entry to pending")?; + tx.commit().await.context("failed to commit v2 release-pending")?; + Ok(mutation_response(QueueLeaseMutationOutcome::Applied, Some(current), "entry returned to pending")) + } + // ============================================================ // queue/next_deadline // ============================================================ @@ -466,9 +1095,7 @@ impl Store { .fetch_one(&self.pool) .await .context("failed to compute next deadline")?; - Ok(QueueNextDeadlineResponse { - next_run_at: next.map(|ts| ts.to_rfc3339()), - }) + Ok(QueueNextDeadlineResponse { next_run_at: next.map(|ts| ts.to_rfc3339()) }) } // ============================================================ @@ -516,20 +1143,13 @@ impl Store { } /// Transition a single Pending entry to Assigned (leased). - pub async fn mark_assigned( - &self, - entry_id: &str, - workflow_id: Option, - ) -> Result { + pub async fn mark_assigned(&self, entry_id: &str, workflow_id: Option) -> Result { + self.ensure_migrated().await?; let wid = workflow_id.unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); - let mut tx = self - .pool - .begin() - .await - .context("failed to begin mark_assigned tx")?; + let mut tx = self.pool.begin().await.context("failed to begin mark_assigned tx")?; let row = sqlx::query(&format!( - "SELECT state FROM {} WHERE id = $1 AND state = ANY($2) FOR UPDATE", + "SELECT state FROM {} WHERE id = $1 AND state = ANY($2) AND COALESCE(fence_version, 0) <> 1 FOR UPDATE", self.table )) .bind(entry_id) @@ -540,20 +1160,14 @@ impl Store { let Some(row) = row else { tx.rollback().await.ok(); - return Ok(QueueMutationResponse { - changed: false, - not_found: true, - }); + return Ok(QueueMutationResponse { changed: false, not_found: true }); }; let state: String = row.get("state"); match state.as_str() { // Already leased: idempotent no-op. db_state::LEASED => { tx.rollback().await.ok(); - Ok(QueueMutationResponse { - changed: false, - not_found: false, - }) + Ok(QueueMutationResponse { changed: false, not_found: false }) } db_state::PENDING => { sqlx::query(&format!( @@ -569,13 +1183,8 @@ impl Store { .execute(&mut *tx) .await .context("failed to mark entry assigned")?; - tx.commit() - .await - .context("failed to commit mark_assigned tx")?; - Ok(QueueMutationResponse { - changed: true, - not_found: false, - }) + tx.commit().await.context("failed to commit mark_assigned tx")?; + Ok(QueueMutationResponse { changed: true, not_found: false }) } // held (or anything else live) cannot jump straight to assigned. _ => { @@ -594,22 +1203,14 @@ impl Store { workflow_ref: Option<&str>, workflow_id: Option<&str>, ) -> Result { - if !matches!( - status, - completion_status::COMPLETED | completion_status::FAILED | completion_status::CANCELLED - ) { - anyhow::bail!( - "invalid completion status: '{status}' (expected one of: completed, failed, cancelled)" - ); + self.ensure_migrated().await?; + if !matches!(status, completion_status::COMPLETED | completion_status::FAILED | completion_status::CANCELLED) { + anyhow::bail!("invalid completion status: '{status}' (expected one of: completed, failed, cancelled)"); } - let mut tx = self - .pool - .begin() - .await - .context("failed to begin completion tx")?; + let mut tx = self.pool.begin().await.context("failed to begin completion tx")?; let row = sqlx::query(&format!( - "SELECT state, workflow_id, payload FROM {} WHERE id = $1 FOR UPDATE", + "SELECT state, workflow_id, payload, fence_version FROM {} WHERE id = $1 FOR UPDATE", self.table )) .bind(entry_id) @@ -619,21 +1220,19 @@ impl Store { let Some(row) = row else { tx.rollback().await.ok(); - return Ok(QueueMutationResponse { - changed: false, - not_found: true, - }); + return Ok(QueueMutationResponse { changed: false, not_found: true }); }; let state: String = row.get("state"); + if row.get::, _>("fence_version") == Some(1) { + tx.rollback().await.ok(); + return Ok(QueueMutationResponse { changed: false, not_found: true }); + } // Only leased (assigned) entries are prunable. A completion frame for a // pending/held/terminal entry must NOT delete un-leased work. if state != db_state::LEASED { tx.rollback().await.ok(); - return Ok(QueueMutationResponse { - changed: false, - not_found: true, - }); + return Ok(QueueMutationResponse { changed: false, not_found: true }); } // Optional workflow_ref / workflow_id match guards. if let Some(want_ref) = workflow_ref { @@ -641,20 +1240,14 @@ impl Store { let dispatch_ref = payload.get("workflow_ref").and_then(Value::as_str); if dispatch_ref.is_some_and(|r| r != want_ref) { tx.rollback().await.ok(); - return Ok(QueueMutationResponse { - changed: false, - not_found: true, - }); + return Ok(QueueMutationResponse { changed: false, not_found: true }); } } if let Some(want_wid) = workflow_id { let existing: Option = row.get("workflow_id"); if existing.as_deref().is_some_and(|w| w != want_wid) { tx.rollback().await.ok(); - return Ok(QueueMutationResponse { - changed: false, - not_found: true, - }); + return Ok(QueueMutationResponse { changed: false, not_found: true }); } } @@ -667,14 +1260,9 @@ impl Store { .execute(&mut *tx) .await .context("failed to mark entry done")?; - tx.commit() - .await - .context("failed to commit completion tx")?; + tx.commit().await.context("failed to commit completion tx")?; - Ok(QueueMutationResponse { - changed: true, - not_found: false, - }) + Ok(QueueMutationResponse { changed: true, not_found: false }) } // ============================================================ @@ -685,11 +1273,8 @@ impl Store { /// order; the rest keep their relative order. Returns the count of entries /// whose absolute position changed. pub async fn reorder(&self, entry_ids: &[String]) -> Result { - let mut tx = self - .pool - .begin() - .await - .context("failed to begin reorder tx")?; + self.ensure_migrated().await?; + let mut tx = self.pool.begin().await.context("failed to begin reorder tx")?; // Current live order. Lock the rows so a concurrent enqueue/lease can't // interleave new ordinals while we rewrite them. @@ -720,11 +1305,7 @@ impl Store { } } - let reordered_count = new_order - .iter() - .zip(original.iter()) - .filter(|(after, before)| after != before) - .count(); + let reordered_count = new_order.iter().zip(original.iter()).filter(|(after, before)| after != before).count(); if reordered_count == 0 { tx.rollback().await.ok(); @@ -760,32 +1341,25 @@ impl Store { entry_id: &str, reason: &str, ) -> std::result::Result { - let mut tx = self - .pool - .begin() + self.ensure_migrated().await.map_err(QueueReleasePendingError::Backend)?; + let mut tx = self.pool.begin().await.map_err(|e| QueueReleasePendingError::Backend(anyhow::Error::from(e)))?; + + let row = sqlx::query(&format!("SELECT state, fence_version FROM {} WHERE id = $1 FOR UPDATE", self.table)) + .bind(entry_id) + .fetch_optional(&mut *tx) .await .map_err(|e| QueueReleasePendingError::Backend(anyhow::Error::from(e)))?; - let row = sqlx::query(&format!( - "SELECT state FROM {} WHERE id = $1 FOR UPDATE", - self.table - )) - .bind(entry_id) - .fetch_optional(&mut *tx) - .await - .map_err(|e| QueueReleasePendingError::Backend(anyhow::Error::from(e)))?; - let Some(row) = row else { - return Err(QueueReleasePendingError::NotFound { - entry_id: entry_id.to_string(), - }); + return Err(QueueReleasePendingError::NotFound { entry_id: entry_id.to_string() }); }; let state: String = row.get("state"); + if row.get::, _>("fence_version") == Some(1) { + return Err(QueueReleasePendingError::NotFound { entry_id: entry_id.to_string() }); + } // Terminal rows are "gone" for the purposes of this contract. if state == db_state::DONE || state == db_state::DROPPED { - return Err(QueueReleasePendingError::NotFound { - entry_id: entry_id.to_string(), - }); + return Err(QueueReleasePendingError::NotFound { entry_id: entry_id.to_string() }); } if state != db_state::LEASED { return Err(QueueReleasePendingError::NotAssigned { @@ -816,23 +1390,88 @@ impl Store { .await .map_err(|e| QueueReleasePendingError::Backend(anyhow::Error::from(e)))?; - tx.commit() - .await - .map_err(|e| QueueReleasePendingError::Backend(anyhow::Error::from(e)))?; + tx.commit().await.map_err(|e| QueueReleasePendingError::Backend(anyhow::Error::from(e)))?; - Ok(QueueReleasePendingResponse { - entry_id: entry_id.to_string(), - status: status::PENDING.to_string(), - }) + Ok(QueueReleasePendingResponse { entry_id: entry_id.to_string(), status: status::PENDING.to_string() }) } // ============================================================ // internal helpers // ============================================================ + async fn v2_recovery_blocks(&self) -> Result> { + let mut blocked = Vec::new(); + let expired = sqlx::query(&format!( + "SELECT id, subject_id, state, workflow_id, enqueued_at, assigned_at, held_at, \ + run_at, expire_after_secs, payload, fence_version, subject_qualified_id, \ + subject_generation, workflow_generation, lease_owner, lease_generation, \ + lease_expires_at, repository, base_ref, head_ref \ + FROM {t} WHERE state = '{leased}' AND lease_expires_at IS NOT NULL \ + AND lease_expires_at < now() ORDER BY ordinal ASC LIMIT 32", + t = self.table, + leased = db_state::LEASED, + )) + .fetch_all(&self.pool) + .await + .context("failed to inspect expired v2 assignments")?; + for row in expired { + let entry_id: String = row.get("id"); + match row_to_execution(&row) { + Ok(execution) => blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::ExpiredLeaseRecoveryRequired, + conflicts_with: Some(execution), + }), + Err(_) => blocked.push(QueueLeaseBlock { + entry_id, + reason: QueueLeaseBlockReason::MissingExecutionIdentity, + conflicts_with: None, + }), + } + } + let missing = sqlx::query(&format!( + "SELECT id FROM {t} WHERE state = '{pending}' \ + AND COALESCE(fence_version, 0) <> 1 \ + AND (run_at IS NULL OR run_at <= now()) \ + ORDER BY ordinal ASC LIMIT 32", + t = self.table, + pending = db_state::PENDING, + )) + .fetch_all(&self.pool) + .await + .context("failed to inspect unfenced pending entries")?; + blocked.extend(missing.into_iter().map(|row| QueueLeaseBlock { + entry_id: row.get("id"), + reason: QueueLeaseBlockReason::MissingExecutionIdentity, + conflicts_with: None, + })); + Ok(blocked) + } + + async fn fenced_row_for_update( + &self, + tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, + entry_id: &str, + ) -> Result> { + sqlx::query(&format!( + "SELECT id, subject_id, state, workflow_id, enqueued_at, assigned_at, held_at, \ + run_at, expire_after_secs, payload, fence_version, subject_qualified_id, \ + subject_generation, workflow_generation, lease_owner, lease_generation, \ + lease_expires_at, repository, base_ref, head_ref, terminal_status, \ + now() AS backend_now \ + FROM {} WHERE id = $1 FOR UPDATE", + self.table + )) + .bind(entry_id) + .fetch_optional(&mut **tx) + .await + .context("failed to load generation-fenced queue entry") + } + /// Soft-delete pending deferred entries past their `run_at + /// expire_after_secs` window (→ `dropped`) instead of dispatching late. async fn sweep_expired(&self) -> Result<()> { + self.ensure_migrated().await?; sqlx::query(&format!( "UPDATE {t} SET state = '{dropped}', updated_at = now() \ WHERE state = '{pending}' AND run_at IS NOT NULL AND expire_after_secs IS NOT NULL \ @@ -853,13 +1492,10 @@ impl Store { where F: FnOnce(&str) -> std::result::Result, { - let mut tx = self - .pool - .begin() - .await - .context("failed to begin mutate tx")?; + self.ensure_migrated().await?; + let mut tx = self.pool.begin().await.context("failed to begin mutate tx")?; let row = sqlx::query(&format!( - "SELECT state FROM {} WHERE id = $1 AND state = ANY($2) FOR UPDATE", + "SELECT state FROM {} WHERE id = $1 AND state = ANY($2) AND COALESCE(fence_version, 0) <> 1 FOR UPDATE", self.table )) .bind(entry_id) @@ -870,32 +1506,19 @@ impl Store { let Some(row) = row else { tx.rollback().await.ok(); - return Ok(QueueMutationResponse { - changed: false, - not_found: true, - }); + return Ok(QueueMutationResponse { changed: false, not_found: true }); }; let state: String = row.get("state"); match plan(&state) { Ok(MutationPlan::NoChange) => { tx.rollback().await.ok(); - Ok(QueueMutationResponse { - changed: false, - not_found: false, - }) + Ok(QueueMutationResponse { changed: false, not_found: false }) } Ok(MutationPlan::Sql(sql)) => { - sqlx::query(&sql) - .bind(entry_id) - .execute(&mut *tx) - .await - .context("failed to apply entry mutation")?; + sqlx::query(&sql).bind(entry_id).execute(&mut *tx).await.context("failed to apply entry mutation")?; tx.commit().await.context("failed to commit mutate tx")?; - Ok(QueueMutationResponse { - changed: true, - not_found: false, - }) + Ok(QueueMutationResponse { changed: true, not_found: false }) } Err(MutationError::NotPending) => { tx.rollback().await.ok(); @@ -907,11 +1530,7 @@ impl Store { /// Test helper: total row count (any state). #[doc(hidden)] pub async fn total_rows(&self) -> Result { - Ok( - sqlx::query_scalar(&format!("SELECT count(*) FROM {}", self.table)) - .fetch_one(&self.pool) - .await?, - ) + Ok(sqlx::query_scalar(&format!("SELECT count(*) FROM {}", self.table)).fetch_one(&self.pool).await?) } } @@ -970,9 +1589,7 @@ fn priority_to_int(priority: Option<&str>) -> i32 { /// Parse an RFC 3339 string into a UTC instant; `None` on malformed input /// (treated as immediate, matching the reference plugin's lenient handling). fn parse_rfc3339(raw: &str) -> Option> { - DateTime::parse_from_rfc3339(raw) - .ok() - .map(|dt| dt.with_timezone(&Utc)) + DateTime::parse_from_rfc3339(raw).ok().map(|dt| dt.with_timezone(&Utc)) } /// Build a wire [`QueueEntry`] from a queue row. Returns `None` (logging) when @@ -1011,16 +1628,81 @@ fn row_to_entry(row: &sqlx::postgres::PgRow) -> Option { }) } +fn row_to_repository(row: &sqlx::postgres::PgRow) -> Result> { + let repository: Option = row.get("repository"); + let base_ref: Option = row.get("base_ref"); + let head_ref: Option = row.get("head_ref"); + match (repository, base_ref, head_ref) { + (None, None, None) => Ok(None), + (Some(repository), Some(base_ref), Some(head_ref)) => { + let reservation = RepositoryReservation { repository, base_ref, head_ref }; + reservation.validate().map_err(anyhow::Error::msg)?; + Ok(Some(reservation)) + } + _ => anyhow::bail!("stored repository reservation is incomplete"), + } +} + +fn row_to_execution(row: &sqlx::postgres::PgRow) -> Result { + if row.get::, _>("fence_version") != Some(1) { + anyhow::bail!("entry lacks generation-fenced execution identity"); + } + let workflow_id: Option = row.get("workflow_id"); + let qualified_id: Option = row.get("subject_qualified_id"); + let subject_generation: Option = row.get("subject_generation"); + let lease_owner: Option = row.get("lease_owner"); + let lease_generation: Option = row.get("lease_generation"); + let lease_expires_at: Option> = row.get("lease_expires_at"); + let execution = ExecutionFence { + schema: EXECUTION_FENCE_SCHEMA_ID.to_string(), + version: EXECUTION_FENCE_VERSION, + workflow_id: workflow_id.context("entry lacks workflow_id")?, + workflow_generation: positive_u64(row.get("workflow_generation"), "workflow_generation")?, + subject: Some(SubjectGeneration { + qualified_id: qualified_id.context("entry lacks subject_qualified_id")?, + generation: positive_u64(subject_generation, "subject_generation")?, + }), + queue_lease: Some(QueueLeaseFence { + entry_id: row.get("id"), + owner_id: lease_owner.context("entry lacks lease_owner")?, + generation: positive_u64(lease_generation, "lease_generation")?, + expires_at: lease_expires_at.context("entry lacks lease_expires_at")?, + }), + repository: row_to_repository(row)?, + }; + execution.validate_queue_backed().map_err(anyhow::Error::msg)?; + Ok(execution) +} + +fn row_to_fenced_entry(row: &sqlx::postgres::PgRow) -> Result { + let entry = row_to_entry(row).context("stored queue dispatch is undecodable")?; + let fenced = FencedQueueEntry { entry, execution: row_to_execution(row)? }; + fenced.validate().map_err(anyhow::Error::msg)?; + Ok(fenced) +} + +fn positive_u64(value: Option, field: &str) -> Result { + let value = value.with_context(|| format!("entry lacks {field}"))?; + u64::try_from(value).ok().filter(|value| *value > 0).with_context(|| format!("stored {field} must be positive")) +} + +fn clamp_ttl(requested: Option, configured: i64) -> i64 { + let configured = configured.max(1) as u64; + requested.unwrap_or(configured).clamp(1, configured) as i64 +} + +fn mutation_response( + outcome: QueueLeaseMutationOutcome, + execution: Option, + reason: impl Into, +) -> QueueLeaseMutationResponse { + QueueLeaseMutationResponse { outcome, execution, reason: Some(reason.into()) } +} + /// Compute aggregate stats from already-mapped live entries. fn stats_from_entries(entries: &[QueueEntry]) -> QueueStats { let now = Utc::now(); - let mut stats = QueueStats { - total: entries.len(), - pending: 0, - assigned: 0, - held: 0, - deferred: 0, - }; + let mut stats = QueueStats { total: entries.len(), pending: 0, assigned: 0, held: 0, deferred: 0 }; for entry in entries { match entry.status.as_str() { s if s == status::PENDING => { diff --git a/tests/postgres.rs b/tests/postgres.rs index d44b72c..f3961a4 100644 --- a/tests/postgres.rs +++ b/tests/postgres.rs @@ -9,8 +9,13 @@ //! Each test uses a unique table name so they can run in parallel and never //! clobber a shared queue. +use animus_execution_protocol::RepositoryReservation; use animus_queue_postgres::config::QueueConfig; use animus_queue_postgres::store::Store; +use animus_queue_protocol::{ + QueueCompletionV2Request, QueueEnqueueV2Request, QueueLeaseMutationOutcome, QueueLeaseRecoverRequest, + QueueLeaseRenewRequest, QueueLeaseV2Request, QueueReleasePendingV2Request, +}; use animus_subject_protocol::{SubjectDispatch, SubjectRef}; use chrono::Utc; @@ -22,11 +27,7 @@ fn test_database_url() -> String { } fn config_with_table(table: &str, lease_ttl_secs: i64) -> QueueConfig { - QueueConfig { - database_url: test_database_url(), - lease_ttl_secs, - table: table.to_string(), - } + QueueConfig { database_url: test_database_url(), lease_ttl_secs, table: table.to_string() } } fn unique_table(prefix: &str) -> String { @@ -37,7 +38,13 @@ fn unique_table(prefix: &str) -> String { /// Open a store, or `None` (with a skip note) when no Postgres is reachable. async fn open_or_skip(config: &QueueConfig) -> Option { match Store::open(config).await { - Ok(store) => Some(store), + Ok(store) => match store.ping().await { + Ok(()) => Some(store), + Err(error) => { + eprintln!("SKIP: no reachable Postgres for integration test: {error}"); + None + } + }, Err(error) => { eprintln!("SKIP: no reachable Postgres for integration test: {error}"); None @@ -46,27 +53,42 @@ async fn open_or_skip(config: &QueueConfig) -> Option { } async fn cleanup(table: &str) { - if let Ok(pool) = sqlx::postgres::PgPoolOptions::new() - .max_connections(1) - .connect(&test_database_url()) - .await - { - let _ = sqlx::query(&format!("DROP TABLE IF EXISTS {table}")) - .execute(&pool) - .await; - let _ = sqlx::query(&format!("DROP SEQUENCE IF EXISTS {table}_ordinal_seq")) - .execute(&pool) - .await; + if let Ok(pool) = sqlx::postgres::PgPoolOptions::new().max_connections(1).connect(&test_database_url()).await { + let _ = sqlx::query(&format!("DROP TABLE IF EXISTS {table}")).execute(&pool).await; + let _ = sqlx::query(&format!("DROP TABLE IF EXISTS {table}_subject_generation")).execute(&pool).await; + let _ = sqlx::query(&format!("DROP SEQUENCE IF EXISTS {table}_ordinal_seq")).execute(&pool).await; } } fn task(task_id: &str, workflow_ref: &str) -> SubjectDispatch { - SubjectDispatch::for_subject_with_metadata( - SubjectRef::task(task_id), - workflow_ref, - "integration-test", - Utc::now(), - ) + SubjectDispatch::for_subject_with_metadata(SubjectRef::task(task_id), workflow_ref, "integration-test", Utc::now()) +} + +fn reservation(task_id: &str) -> RepositoryReservation { + RepositoryReservation { + repository: "https://github.com/launchapp-dev/animus-cli.git".to_string(), + base_ref: "refs/heads/main".to_string(), + head_ref: format!("refs/heads/animus/{task_id}"), + } +} + +fn enqueue_v2_request(task_id: &str) -> QueueEnqueueV2Request { + QueueEnqueueV2Request { + subject_dispatch: task(task_id, "coding"), + idempotency_key: Some(format!("delivery-{task_id}")), + repository: Some(reservation(task_id)), + run_at: None, + expire_after_secs: None, + } +} + +fn lease_one(owner_id: &str, workflow_id: &str) -> QueueLeaseV2Request { + QueueLeaseV2Request { + max: 1, + owner_id: owner_id.to_string(), + workflow_ids: vec![workflow_id.to_string()], + exclude: Vec::new(), + } } #[tokio::test] @@ -77,14 +99,8 @@ async fn enqueue_lease_and_list_round_trip() { return; }; - let a = store - .enqueue(task("TASK-1", "standard"), None, None) - .await - .unwrap(); - let b = store - .enqueue(task("TASK-2", "standard"), None, None) - .await - .unwrap(); + let a = store.enqueue(task("TASK-1", "standard"), None, None).await.unwrap(); + let b = store.enqueue(task("TASK-2", "standard"), None, None).await.unwrap(); assert!(a.enqueued && b.enqueued); assert_ne!(a.entry_id, b.entry_id); @@ -114,39 +130,24 @@ async fn lease_expiry_reclaims_crashed_daemon_work() { return; }; - let enq = store - .enqueue(task("TASK-1", "standard"), None, None) - .await - .unwrap(); + let enq = store.enqueue(task("TASK-1", "standard"), None, None).await.unwrap(); // First lease claims it. - let first = store - .lease(1, Some(vec!["wf-1".into()]), None) - .await - .unwrap(); + let first = store.lease(1, Some(vec!["wf-1".into()]), None).await.unwrap(); assert_eq!(first.leased.len(), 1); assert_eq!(first.leased[0].entry_id, enq.entry_id); // Immediately re-leasing finds nothing — the lease is still live. let still_held = store.lease(1, None, None).await.unwrap(); - assert!( - still_held.leased.is_empty(), - "a live lease must not be re-handed-out" - ); + assert!(still_held.leased.is_empty(), "a live lease must not be re-handed-out"); // Let the 1s lease TTL expire (simulating the holder dying). tokio::time::sleep(std::time::Duration::from_millis(1300)).await; // The expired lease is reclaimed and re-dispatched. - let reclaimed = store - .lease(1, Some(vec!["wf-2".into()]), None) - .await - .unwrap(); + let reclaimed = store.lease(1, Some(vec!["wf-2".into()]), None).await.unwrap(); assert_eq!(reclaimed.leased.len(), 1, "expired lease must be reclaimed"); - assert_eq!( - reclaimed.leased[0].entry_id, enq.entry_id, - "same entry re-dispatched" - ); + assert_eq!(reclaimed.leased[0].entry_id, enq.entry_id, "same entry re-dispatched"); assert_eq!(reclaimed.leased[0].workflow_id.as_deref(), Some("wf-2")); cleanup(&table).await; @@ -160,10 +161,7 @@ async fn hold_release_drop_and_release_pending() { return; }; - let e = store - .enqueue(task("TASK-1", "standard"), None, None) - .await - .unwrap(); + let e = store.enqueue(task("TASK-1", "standard"), None, None).await.unwrap(); // hold → held, idempotent assert!(store.hold(&e.entry_id).await.unwrap().changed); @@ -177,15 +175,9 @@ async fn hold_release_drop_and_release_pending() { assert_eq!(store.stats().await.unwrap().pending, 1); // lease then release_pending → back to pending - let leased = store - .lease(1, Some(vec!["wf-x".into()]), None) - .await - .unwrap(); + let leased = store.lease(1, Some(vec!["wf-x".into()]), None).await.unwrap(); assert_eq!(leased.leased.len(), 1); - let rp = store - .release_pending(&e.entry_id, "duplicate-in-flight") - .await - .unwrap(); + let rp = store.release_pending(&e.entry_id, "duplicate-in-flight").await.unwrap(); assert_eq!(rp.status, "pending"); assert_eq!(store.stats().await.unwrap().pending, 1); @@ -207,28 +199,16 @@ async fn completion_prunes_only_leased_entries() { return; }; - let e = store - .enqueue(task("TASK-1", "standard"), None, None) - .await - .unwrap(); + let e = store.enqueue(task("TASK-1", "standard"), None, None).await.unwrap(); // completion on a pending (not-leased) entry is a no-op. - let noop = store - .completion(&e.entry_id, "completed", None, None) - .await - .unwrap(); + let noop = store.completion(&e.entry_id, "completed", None, None).await.unwrap(); assert!(!noop.changed && noop.not_found); assert_eq!(store.stats().await.unwrap().pending, 1); // lease then complete → pruned. - store - .lease(1, Some(vec!["wf-1".into()]), None) - .await - .unwrap(); - let done = store - .completion(&e.entry_id, "completed", Some("standard"), Some("wf-1")) - .await - .unwrap(); + store.lease(1, Some(vec!["wf-1".into()]), None).await.unwrap(); + let done = store.completion(&e.entry_id, "completed", Some("standard"), Some("wf-1")).await.unwrap(); assert!(done.changed); assert_eq!(store.stats().await.unwrap().total, 0); @@ -245,10 +225,7 @@ async fn deferred_entry_not_leased_until_due_and_swept_when_expired() { // Future run_at → pending but not leasable; contributes to next_deadline. let future = (Utc::now() + chrono::Duration::hours(1)).to_rfc3339(); - store - .enqueue(task("FUTURE", "standard"), Some(future.clone()), None) - .await - .unwrap(); + store.enqueue(task("FUTURE", "standard"), Some(future.clone()), None).await.unwrap(); assert!(store.lease(5, None, None).await.unwrap().leased.is_empty()); let stats = store.stats().await.unwrap(); assert_eq!(stats.pending, 1); @@ -257,10 +234,7 @@ async fn deferred_entry_not_leased_until_due_and_swept_when_expired() { // Past run_at + tiny grace → swept (dropped), never dispatched late. let past = (Utc::now() - chrono::Duration::hours(2)).to_rfc3339(); - store - .enqueue(task("EXPIRED", "standard"), Some(past), Some(60)) - .await - .unwrap(); + store.enqueue(task("EXPIRED", "standard"), Some(past), Some(60)).await.unwrap(); let leased = store.lease(5, None, None).await.unwrap(); assert!(leased.leased.is_empty()); // FUTURE remains pending; EXPIRED swept. @@ -277,23 +251,11 @@ async fn reorder_moves_named_entries_to_front() { return; }; - let a = store - .enqueue(task("TASK-1", "standard"), None, None) - .await - .unwrap(); - let _b = store - .enqueue(task("TASK-2", "standard"), None, None) - .await - .unwrap(); - let c = store - .enqueue(task("TASK-3", "standard"), None, None) - .await - .unwrap(); + let a = store.enqueue(task("TASK-1", "standard"), None, None).await.unwrap(); + let _b = store.enqueue(task("TASK-2", "standard"), None, None).await.unwrap(); + let c = store.enqueue(task("TASK-3", "standard"), None, None).await.unwrap(); - let res = store - .reorder(&[c.entry_id.clone(), a.entry_id.clone()]) - .await - .unwrap(); + let res = store.reorder(&[c.entry_id.clone(), a.entry_id.clone()]).await.unwrap(); assert!(res.reordered_count > 0); let listed = store.list(&[], None, None).await.unwrap(); @@ -312,22 +274,298 @@ async fn exclude_subjects_skips_in_flight_subjects() { return; }; - store - .enqueue(task("TASK-1", "standard"), None, None) + store.enqueue(task("TASK-1", "standard"), None, None).await.unwrap(); + store.enqueue(task("TASK-2", "standard"), None, None).await.unwrap(); + + // Exclude TASK-1: only TASK-2 should be leased. + let leased = store.lease(5, None, Some(vec!["TASK-1".to_string()])).await.unwrap(); + assert_eq!(leased.leased.len(), 1); + assert_eq!(leased.leased[0].subject_id, "TASK-2"); + + cleanup(&table).await; +} + +#[tokio::test] +async fn v2_leases_five_independent_slots_and_replays_enqueue_idempotently() { + let table = unique_table("v2five"); + let config = config_with_table(&table, 1800); + let Some(store) = open_or_skip(&config).await else { + return; + }; + + let mut original_entry = String::new(); + for index in 1..=5 { + let response = store.enqueue_v2(enqueue_v2_request(&format!("TASK-{index}"))).await.unwrap(); + assert!(response.enqueued); + assert_eq!(response.subject.generation, 1); + if index == 1 { + original_entry = response.entry_id; + } + } + let replay = store.enqueue_v2(enqueue_v2_request("TASK-1")).await.unwrap(); + assert!(!replay.enqueued); + assert_eq!(replay.entry_id, original_entry); + assert_eq!(replay.subject.generation, 1); + let mut mismatched_replay = enqueue_v2_request("TASK-1"); + mismatched_replay.repository = Some(reservation("DIFFERENT-BRANCH")); + assert!(store + .enqueue_v2(mismatched_replay) .await - .unwrap(); - store - .enqueue(task("TASK-2", "standard"), None, None) + .unwrap_err() + .to_string() + .contains("different repository reservation")); + let mut malformed_schedule = enqueue_v2_request("TASK-BAD-TIME"); + malformed_schedule.run_at = Some("tomorrow morning".to_string()); + assert!(store.enqueue_v2(malformed_schedule).await.unwrap_err().to_string().contains("RFC 3339")); + + assert!( + store.lease(5, None, None).await.unwrap().leased.is_empty(), + "legacy leasing must fail closed for v2 entries" + ); + + let store_one = store.clone(); + let store_two = store.clone(); + let store_three = store.clone(); + let store_four = store.clone(); + let store_five = store.clone(); + let (one, two, three, four, five) = tokio::join!( + store_one.lease_v2(lease_one("daemon-1", "workflow-1")), + store_two.lease_v2(lease_one("daemon-2", "workflow-2")), + store_three.lease_v2(lease_one("daemon-3", "workflow-3")), + store_four.lease_v2(lease_one("daemon-4", "workflow-4")), + store_five.lease_v2(lease_one("daemon-5", "workflow-5")), + ); + let leased: Vec<_> = + [one, two, three, four, five].into_iter().flat_map(|response| response.unwrap().leased).collect(); + assert_eq!(leased.len(), 5); + assert!(leased.iter().all(|entry| { + entry.execution.workflow_generation == 1 && entry.execution.queue_lease.as_ref().unwrap().generation == 1 + })); + assert!(store + .lease_v2(QueueLeaseV2Request { + max: 6, + owner_id: "daemon-a".to_string(), + workflow_ids: (6..=11).map(|index| format!("workflow-{index}")).collect(), + exclude: Vec::new(), + }) .await - .unwrap(); + .unwrap_err() + .to_string() + .contains("exceeds fleet limit")); + + cleanup(&table).await; +} + +#[tokio::test] +async fn v2_concurrent_duplicate_subject_enqueue_creates_one_live_generation() { + let table = unique_table("v2subject"); + let config = config_with_table(&table, 1800); + let Some(store) = open_or_skip(&config).await else { + return; + }; + + let mut first_request = enqueue_v2_request("TASK-1"); + first_request.idempotency_key = Some("delivery-a".to_string()); + let mut second_request = enqueue_v2_request("TASK-1"); + second_request.idempotency_key = Some("delivery-b".to_string()); + let first_store = store.clone(); + let second_store = store.clone(); + let (first, second) = tokio::join!(first_store.enqueue_v2(first_request), second_store.enqueue_v2(second_request),); + let first = first.unwrap(); + let second = second.unwrap(); + assert_eq!(usize::from(first.enqueued) + usize::from(second.enqueued), 1); + assert_eq!(first.entry_id, second.entry_id); + assert_eq!(first.subject, second.subject); + assert!(first.warning.is_some() || second.warning.is_some()); - // Exclude TASK-1: only TASK-2 should be leased. let leased = store - .lease(5, None, Some(vec!["TASK-1".to_string()])) + .lease_v2(QueueLeaseV2Request { + max: 5, + owner_id: "daemon-a".to_string(), + workflow_ids: (1..=5).map(|index| format!("workflow-{index}")).collect(), + exclude: Vec::new(), + }) .await .unwrap(); assert_eq!(leased.leased.len(), 1); - assert_eq!(leased.leased[0].subject_id, "TASK-2"); + assert_eq!(leased.leased[0].entry.subject_id, "TASK-1"); + + cleanup(&table).await; +} + +#[tokio::test] +async fn v2_repository_head_collision_fails_closed() { + let table = unique_table("v2ref"); + let config = config_with_table(&table, 1800); + let Some(store) = open_or_skip(&config).await else { + return; + }; + + store.enqueue_v2(enqueue_v2_request("TASK-1")).await.unwrap(); + let mut collision = enqueue_v2_request("TASK-2"); + collision.repository = Some(reservation("TASK-1")); + let error = store.enqueue_v2(collision).await.unwrap_err(); + assert!(error.to_string().contains("ownership collision")); + assert_eq!(store.stats().await.unwrap().pending, 1); + + cleanup(&table).await; +} + +#[tokio::test] +async fn v2_expiry_requires_recovery_and_fences_the_previous_owner() { + let table = unique_table("v2recover"); + let config = config_with_table(&table, 1); + let Some(store) = open_or_skip(&config).await else { + return; + }; + + store.enqueue_v2(enqueue_v2_request("TASK-1")).await.unwrap(); + let first = store + .lease_v2(QueueLeaseV2Request { + max: 1, + owner_id: "daemon-a".to_string(), + workflow_ids: vec!["workflow-stable".to_string()], + exclude: Vec::new(), + }) + .await + .unwrap() + .leased + .remove(0) + .execution; + tokio::time::sleep(std::time::Duration::from_millis(1300)).await; + + let ordinary_lease = store + .lease_v2(QueueLeaseV2Request { + max: 1, + owner_id: "daemon-b".to_string(), + workflow_ids: vec!["workflow-would-duplicate".to_string()], + exclude: Vec::new(), + }) + .await + .unwrap(); + assert!(ordinary_lease.leased.is_empty()); + assert_eq!( + ordinary_lease.blocked[0].reason, + animus_queue_protocol::QueueLeaseBlockReason::ExpiredLeaseRecoveryRequired + ); + + let recovered = store + .recover_v2(QueueLeaseRecoverRequest { + execution: first.clone(), + new_owner_id: "daemon-b".to_string(), + ttl_secs: None, + }) + .await + .unwrap(); + assert_eq!(recovered.outcome, QueueLeaseMutationOutcome::Applied); + let recovered_fence = recovered.execution.unwrap(); + assert!(first.same_execution_generation(&recovered_fence)); + assert_eq!(recovered_fence.workflow_id, "workflow-stable"); + assert_eq!(recovered_fence.queue_lease.as_ref().unwrap().generation, 2); + assert_eq!(recovered_fence.queue_lease.as_ref().unwrap().owner_id, "daemon-b"); + + let stale = store.renew_v2(QueueLeaseRenewRequest { execution: first.clone(), ttl_secs: None }).await.unwrap(); + assert_eq!(stale.outcome, QueueLeaseMutationOutcome::StaleFence); + + let stale_completion = store + .completion_v2(QueueCompletionV2Request { + execution: first, + status: "completed".to_string(), + workflow_ref: Some("coding".to_string()), + }) + .await + .unwrap(); + assert_eq!(stale_completion.outcome, QueueLeaseMutationOutcome::StaleFence); + let completion = store + .completion_v2(QueueCompletionV2Request { + execution: recovered_fence.clone(), + status: "completed".to_string(), + workflow_ref: Some("coding".to_string()), + }) + .await + .unwrap(); + assert_eq!(completion.outcome, QueueLeaseMutationOutcome::Applied); + let replayed_completion = store + .completion_v2(QueueCompletionV2Request { + execution: recovered_fence, + status: "completed".to_string(), + workflow_ref: Some("coding".to_string()), + }) + .await + .unwrap(); + assert_eq!(replayed_completion.outcome, QueueLeaseMutationOutcome::AlreadyApplied); + let conflicting_completion = store + .completion_v2(QueueCompletionV2Request { + execution: replayed_completion.execution.unwrap(), + status: "failed".to_string(), + workflow_ref: Some("coding".to_string()), + }) + .await + .unwrap(); + assert_eq!(conflicting_completion.outcome, QueueLeaseMutationOutcome::StaleFence); + + cleanup(&table).await; +} + +#[tokio::test] +async fn v2_release_pending_preserves_canonical_workflow_identity() { + let table = unique_table("v2release"); + let config = config_with_table(&table, 1800); + let Some(store) = open_or_skip(&config).await else { + return; + }; + + store.enqueue_v2(enqueue_v2_request("TASK-1")).await.unwrap(); + let first = store + .lease_v2(QueueLeaseV2Request { + max: 1, + owner_id: "daemon-a".to_string(), + workflow_ids: vec!["workflow-stable".to_string()], + exclude: Vec::new(), + }) + .await + .unwrap() + .leased + .remove(0) + .execution; + let released = store + .release_pending_v2(QueueReleasePendingV2Request { + execution: first.clone(), + reason: "slot handoff".to_string(), + }) + .await + .unwrap(); + assert_eq!(released.outcome, QueueLeaseMutationOutcome::Applied); + + store.enqueue_v2(enqueue_v2_request("TASK-2")).await.unwrap(); + let advanced = store + .lease_v2(QueueLeaseV2Request { + max: 1, + owner_id: "daemon-b".to_string(), + workflow_ids: vec!["workflow-task-2".to_string()], + exclude: vec![first.clone()], + }) + .await + .unwrap(); + assert_eq!(advanced.leased[0].entry.subject_id, "TASK-2"); + assert!(advanced.blocked.iter().any(|block| block.entry_id == first.queue_lease.as_ref().unwrap().entry_id)); + + let second = store + .lease_v2(QueueLeaseV2Request { + max: 1, + owner_id: "daemon-b".to_string(), + workflow_ids: vec!["must-not-replace".to_string()], + exclude: Vec::new(), + }) + .await + .unwrap() + .leased + .remove(0) + .execution; + assert_eq!(second.workflow_id, first.workflow_id); + assert_eq!(second.workflow_generation, first.workflow_generation); + assert_eq!(second.subject, first.subject); + assert_eq!(second.queue_lease.as_ref().unwrap().generation, 2); cleanup(&table).await; }