diff --git a/crates/manager-core/tests/migration_queue_messages.rs b/crates/manager-core/tests/migration_queue_messages.rs new file mode 100644 index 0000000..74c9492 --- /dev/null +++ b/crates/manager-core/tests/migration_queue_messages.rs @@ -0,0 +1,150 @@ +//! v1.1.9 migration smoke test: applies all manager-core migrations to +//! `DATABASE_URL` and verifies the `queue_messages` + `queue_trigger_details` +//! shape (table existence, columns, partial indexes). Skips cleanly when +//! `DATABASE_URL` is unset. +//! +//! Heavier schema assertions live in the schema_snapshot golden — this +//! test is a quick gate for the migration itself + an easy place to +//! exercise the partial-index conditions. + +#![allow(clippy::needless_pass_by_value)] + +use sqlx::postgres::PgPoolOptions; +use sqlx::PgPool; + +async fn pool_or_skip() -> Option { + let Ok(url) = std::env::var("DATABASE_URL") else { + eprintln!("migration_queue_messages: DATABASE_URL unset — skipping"); + return None; + }; + let pool = PgPoolOptions::new() + .max_connections(2) + .connect(&url) + .await + .expect("connect to DATABASE_URL"); + sqlx::migrate!("./migrations") + .run(&pool) + .await + .expect("apply migrations"); + Some(pool) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn queue_messages_table_exists_with_expected_columns() { + let Some(pool) = pool_or_skip().await else { return }; + let rows: Vec<(String, String, String)> = sqlx::query_as( + "SELECT column_name, data_type, is_nullable \ + FROM information_schema.columns \ + WHERE table_name = 'queue_messages' \ + ORDER BY ordinal_position", + ) + .fetch_all(&pool) + .await + .expect("column read"); + + let names: Vec<&str> = rows.iter().map(|(n, _, _)| n.as_str()).collect(); + for required in [ + "id", + "app_id", + "queue_name", + "payload", + "enqueued_at", + "deliver_after", + "claim_token", + "claimed_at", + "attempt", + "max_attempts", + "enqueued_by_principal", + ] { + assert!( + names.contains(&required), + "queue_messages should have column {required}, columns are {names:?}" + ); + } + + // Verify nullability on the right columns. + let nullable: std::collections::HashMap = rows + .into_iter() + .map(|(n, _, nl)| (n, nl)) + .collect(); + assert_eq!(nullable["app_id"], "NO"); + assert_eq!(nullable["queue_name"], "NO"); + assert_eq!(nullable["payload"], "NO"); + assert_eq!(nullable["claim_token"], "YES"); + assert_eq!(nullable["deliver_after"], "YES"); + assert_eq!(nullable["enqueued_by_principal"], "YES"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn queue_trigger_details_table_exists() { + let Some(pool) = pool_or_skip().await else { return }; + let rows: Vec<(String,)> = sqlx::query_as( + "SELECT column_name FROM information_schema.columns \ + WHERE table_name = 'queue_trigger_details' \ + ORDER BY ordinal_position", + ) + .fetch_all(&pool) + .await + .expect("read"); + let names: Vec<&str> = rows.iter().map(|(n,)| n.as_str()).collect(); + for required in ["trigger_id", "queue_name", "visibility_timeout_secs", "last_fired_at"] { + assert!( + names.contains(&required), + "queue_trigger_details should have {required}, columns are {names:?}" + ); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn queue_widens_trigger_kind_and_outbox_source_kind() { + let Some(pool) = pool_or_skip().await else { return }; + // The triggers.kind constraint must accept 'queue'. + sqlx::query( + "DO $$ BEGIN PERFORM 1 FROM pg_constraint WHERE conname = 'triggers_kind_check'; END $$;", + ) + .execute(&pool) + .await + .expect("constraint exists"); + + let (def,): (String,) = sqlx::query_as( + "SELECT pg_get_constraintdef(oid) FROM pg_constraint \ + WHERE conname = 'triggers_kind_check'", + ) + .fetch_one(&pool) + .await + .expect("read constraint"); + assert!( + def.contains("'queue'"), + "triggers_kind_check should admit 'queue', got: {def}" + ); + + let (outbox_def,): (String,) = sqlx::query_as( + "SELECT pg_get_constraintdef(oid) FROM pg_constraint \ + WHERE conname = 'outbox_source_kind_check'", + ) + .fetch_one(&pool) + .await + .expect("read constraint"); + assert!( + outbox_def.contains("'invoke'"), + "outbox_source_kind_check should admit 'invoke', got: {outbox_def}" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn queue_messages_dispatch_index_is_partial() { + let Some(pool) = pool_or_skip().await else { return }; + let (defn,): (String,) = sqlx::query_as( + "SELECT indexdef FROM pg_indexes \ + WHERE indexname = 'idx_queue_messages_dispatch'", + ) + .fetch_one(&pool) + .await + .expect("read index"); + // Partial WHERE (claim_token IS NULL) is the key feature; assert + // it's there so a future migration can't quietly drop it. + assert!( + defn.contains("claim_token IS NULL"), + "idx_queue_messages_dispatch should be partial on claim_token IS NULL, got: {defn}" + ); +} diff --git a/crates/picloud/tests/invoke_e2e.rs b/crates/picloud/tests/invoke_e2e.rs new file mode 100644 index 0000000..31fb0d0 --- /dev/null +++ b/crates/picloud/tests/invoke_e2e.rs @@ -0,0 +1,261 @@ +//! v1.1.9 invoke end-to-end tests. invoke()/invoke_async() are +//! exercised against the full all-in-one app via build_app, so the +//! invoke service + Engine back-reference + InvokeServiceImpl all +//! participate. +//! +//! Skips when DATABASE_URL is unset. + +#![allow(clippy::needless_pass_by_value)] + +use std::time::Duration; + +use axum_test::TestServer; +use serde_json::{json, Value}; +use sqlx::postgres::PgPoolOptions; +use sqlx::PgPool; +use uuid::Uuid; + +async fn pool_or_skip() -> Option { + let Ok(url) = std::env::var("DATABASE_URL") else { + eprintln!("invoke_e2e: DATABASE_URL unset — skipping"); + return None; + }; + let pool = PgPoolOptions::new() + .max_connections(5) + .connect(&url) + .await + .expect("connect to DATABASE_URL"); + sqlx::migrate!("../manager-core/migrations") + .run(&pool) + .await + .expect("apply migrations"); + Some(pool) +} + +async fn server_for(pool: PgPool, suffix: &str) -> (TestServer, String) { + use picloud_manager_core::auth::hash_password; + use picloud_shared::InstanceRole; + + let unique = format!("{suffix}-{}", Uuid::new_v4().simple()); + let auth = picloud::AuthDeps::from_pool(pool.clone()); + let username = format!("e2e-{unique}"); + let hash = hash_password("pw").expect("hash"); + auth.users + .create(&username, &hash, InstanceRole::Owner, None) + .await + .expect("seed admin"); + + let app = picloud::build_app( + pool, + auth, + picloud_shared::MasterKey::from_bytes([0x42u8; 32]), + ) + .await + .expect("build_app"); + let mut server = TestServer::new(app).expect("TestServer"); + let resp = server + .post("/api/v1/admin/auth/login") + .json(&json!({ "username": username, "password": "pw" })) + .await; + resp.assert_status_ok(); + let token = resp.json::()["token"] + .as_str() + .expect("login token") + .to_string(); + server.add_header("authorization", format!("Bearer {token}")); + + let slug = format!("e2e-{unique}"); + let created: Value = server + .post("/api/v1/admin/apps") + .json(&json!({ "slug": slug, "name": slug })) + .await + .json(); + let app_id = created["id"].as_str().expect("app id").to_string(); + (server, app_id) +} + +async fn create_script(server: &TestServer, app_id: &str, name: &str, source: &str) -> String { + let created: Value = server + .post("/api/v1/admin/scripts") + .json(&json!({ "app_id": app_id, "name": name, "source": source })) + .await + .json(); + created["id"].as_str().expect("script id").to_string() +} + +/// Create a route bound to a script so we can invoke it via HTTP (and +/// any script can invoke() by route path). +async fn create_route(server: &TestServer, app_id: &str, script_id: &str, path: &str) { + let resp = server + .post(&format!("/api/v1/admin/apps/{app_id}/routes")) + .json(&json!({ + "script_id": script_id, + "method": "GET", + "path": path, + "host": null, + "dispatch_mode": "sync" + })) + .await; + resp.assert_status(axum::http::StatusCode::CREATED); +} + +async fn poll_marker(pool: &PgPool, app_id: &str) -> Option { + for _ in 0..150 { + let row: Option<(Value,)> = sqlx::query_as( + "SELECT value FROM kv_entries WHERE app_id = $1 \ + AND collection = 'e2e_markers' AND key = 'marker'", + ) + .bind(Uuid::parse_str(app_id).expect("uuid")) + .fetch_optional(pool) + .await + .expect("kv read"); + if let Some((v,)) = row { + return Some(v); + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + None +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn invoke_by_name_same_app_returns_value() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "invoke-name").await; + + // Callee: returns its body's `x` + 1 as the response body. + create_script( + &server, + &app_id, + "callee", + r#"#{ statusCode: 200, body: ctx.request.body.x + 1 }"#, + ) + .await; + + // Caller: invokes by name, writes result to KV marker. + let caller = create_script( + &server, + &app_id, + "caller", + r#" + let r = invoke("callee", #{ x: 41 }); + kv::collection("e2e_markers").set("marker", #{ r: r }); + #{ statusCode: 200 } + "#, + ) + .await; + create_route(&server, &app_id, &caller, "/caller").await; + + let resp = server.get("/caller").await; + resp.assert_status_ok(); + + let marker = poll_marker(&pool, &app_id).await.expect("marker fired"); + assert_eq!(marker["r"], 42); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn invoke_cross_app_rejects() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_a) = server_for(pool.clone(), "invoke-A").await; + let (_, app_b) = server_for(pool.clone(), "invoke-B").await; + + // Callee lives in app_b. + let callee_b = create_script( + &server, + &app_b, + "callee_b", + r#"#{ statusCode: 200, body: 0 }"#, + ) + .await; + + // Caller in app_a tries to invoke callee_b by id; must reject. + let caller_src = format!( + r#" + let err = "no error"; + try {{ + invoke("{callee_b}", #{{}}); + }} catch(e) {{ + err = e; + }} + kv::collection("e2e_markers").set("marker", #{{ err: err }}); + #{{ statusCode: 200 }} + "# + ); + let caller_a = create_script(&server, &app_a, "caller_a", &caller_src).await; + create_route(&server, &app_a, &caller_a, "/x").await; + + server.get("/x").await.assert_status_ok(); + + let marker = poll_marker(&pool, &app_a).await.expect("marker fired"); + let err_str = marker["err"].as_str().unwrap_or_default(); + assert!( + err_str.contains("different app") + || err_str.contains("CrossApp") + || err_str.contains("not found"), + "expected cross-app or not-found rejection, got: {err_str}" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn invoke_depth_limit_exceeds_cleanly() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "invoke-depth").await; + + // Recursive script — invoke itself by name; loop bounded by + // Limits::trigger_depth_max (default 8). Records the error string. + let recurser_src = r#" + let err = "ok"; + try { + invoke("recurser", #{}); + } catch(e) { + err = e; + } + kv::collection("e2e_markers").set("marker", #{ err: err }); + #{ statusCode: 200 } + "#; + let recurser = create_script(&server, &app_id, "recurser", recurser_src).await; + create_route(&server, &app_id, &recurser, "/recurse").await; + + server.get("/recurse").await.assert_status_ok(); + + let marker = poll_marker(&pool, &app_id).await.expect("marker fired"); + let err = marker["err"].as_str().unwrap_or_default(); + assert!( + err.contains("depth"), + "expected depth-limit error, got: {err}" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn invoke_async_enqueues_outbox_row() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "invoke-async").await; + + // Callee writes a marker. + create_script( + &server, + &app_id, + "async_callee", + r#" + kv::collection("e2e_markers").set("marker", #{ from: "async" }); + #{ statusCode: 200 } + "#, + ) + .await; + + let caller = create_script( + &server, + &app_id, + "async_caller", + r#" + let id = invoke_async("async_callee", #{}); + #{ statusCode: 200, body: id } + "#, + ) + .await; + create_route(&server, &app_id, &caller, "/start").await; + server.get("/start").await.assert_status_ok(); + + // Dispatcher fires the OutboxSourceKind::Invoke row → callee runs. + let marker = poll_marker(&pool, &app_id).await.expect("marker fired"); + assert_eq!(marker["from"], "async"); +} diff --git a/crates/picloud/tests/queue_e2e.rs b/crates/picloud/tests/queue_e2e.rs new file mode 100644 index 0000000..87cb89a --- /dev/null +++ b/crates/picloud/tests/queue_e2e.rs @@ -0,0 +1,337 @@ +//! v1.1.9 queue end-to-end tests. +//! +//! Wires the full all-in-one app via `build_app` (spawns the real +//! dispatcher + queue arm + reclaim task), creates an app + a consumer +//! script + a `queue:receive` trigger, enqueues messages through the +//! admin path, and observes the handler's side effect via a KV marker. +//! +//! ## Gating +//! +//! Skips when `DATABASE_URL` is unset — mirrors `dispatcher_e2e.rs`. +//! +//! ## Handler observation +//! +//! Consumer handlers write `ctx.event` to a KV marker on a collection +//! no trigger watches. Tests poll the marker for the expected attempt +//! count / value. + +#![allow(clippy::needless_pass_by_value)] + +use std::time::Duration; + +use axum_test::TestServer; +use serde_json::{json, Value}; +use sqlx::postgres::PgPoolOptions; +use sqlx::PgPool; +use uuid::Uuid; + +async fn pool_or_skip() -> Option { + let Ok(url) = std::env::var("DATABASE_URL") else { + eprintln!("queue_e2e: DATABASE_URL unset — skipping"); + return None; + }; + let pool = PgPoolOptions::new() + .max_connections(5) + .connect(&url) + .await + .expect("connect to DATABASE_URL"); + sqlx::migrate!("../manager-core/migrations") + .run(&pool) + .await + .expect("apply migrations"); + Some(pool) +} + +async fn server_for(pool: PgPool, suffix: &str) -> (TestServer, String) { + use picloud_manager_core::auth::hash_password; + use picloud_shared::InstanceRole; + + let unique = format!("{suffix}-{}", Uuid::new_v4().simple()); + let auth = picloud::AuthDeps::from_pool(pool.clone()); + let username = format!("e2e-{unique}"); + let hash = hash_password("pw").expect("hash"); + auth.users + .create(&username, &hash, InstanceRole::Owner, None) + .await + .expect("seed admin"); + + let app = picloud::build_app( + pool, + auth, + picloud_shared::MasterKey::from_bytes([0x42u8; 32]), + ) + .await + .expect("build_app"); + let mut server = TestServer::new(app).expect("TestServer"); + let resp = server + .post("/api/v1/admin/auth/login") + .json(&json!({ "username": username, "password": "pw" })) + .await; + resp.assert_status_ok(); + let token = resp.json::()["token"] + .as_str() + .expect("login token") + .to_string(); + server.add_header("authorization", format!("Bearer {token}")); + + let slug = format!("e2e-{unique}"); + let created: Value = server + .post("/api/v1/admin/apps") + .json(&json!({ "slug": slug, "name": slug })) + .await + .json(); + let app_id = created["id"].as_str().expect("app id").to_string(); + (server, app_id) +} + +async fn create_script(server: &TestServer, app_id: &str, name: &str, source: &str) -> String { + let created: Value = server + .post("/api/v1/admin/scripts") + .json(&json!({ "app_id": app_id, "name": name, "source": source })) + .await + .json(); + created["id"].as_str().expect("script id").to_string() +} + +/// The acker handler records ctx.event into KV and returns successfully. +const ACK_HANDLER: &str = r#" + let e = ctx.event; + kv::collection("e2e_markers").set("marker", e); + #{ ok: true } +"#; + +/// The throwing handler records its event and then throws — used to +/// exercise nack → retry → dead-letter. +const THROW_HANDLER: &str = r#" + let e = ctx.event; + kv::collection("e2e_markers").set("marker", e); + throw "boom" +"#; + +async fn poll_marker(pool: &PgPool, app_id: &str) -> Option { + for _ in 0..200 { + let row: Option<(Value,)> = sqlx::query_as( + "SELECT value FROM kv_entries WHERE app_id = $1 \ + AND collection = 'e2e_markers' AND key = 'marker'", + ) + .bind(Uuid::parse_str(app_id).expect("uuid")) + .fetch_optional(pool) + .await + .expect("kv read"); + if let Some((v,)) = row { + return Some(v); + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + None +} + +async fn count_queue_messages(pool: &PgPool, app_id: &str, queue: &str) -> i64 { + let (n,): (i64,) = sqlx::query_as( + "SELECT COUNT(*) FROM queue_messages WHERE app_id = $1 AND queue_name = $2", + ) + .bind(Uuid::parse_str(app_id).expect("uuid")) + .bind(queue) + .fetch_one(pool) + .await + .expect("queue count"); + n +} + +async fn count_dead_letters_for_queue(pool: &PgPool, app_id: &str, queue: &str) -> i64 { + let (n,): (i64,) = sqlx::query_as( + "SELECT COUNT(*) FROM dead_letters \ + WHERE app_id = $1 AND source = 'queue' \ + AND payload->>'queue_name' = $2", + ) + .bind(Uuid::parse_str(app_id).expect("uuid")) + .bind(queue) + .fetch_one(pool) + .await + .expect("dl count"); + n +} + +/// Helper: directly INSERT a queue message via SQL (no producer-script +/// path). The producer script path is covered in invoke_e2e + the unit +/// tests on the SDK bridge. +async fn enqueue_directly(pool: &PgPool, app_id: &str, queue: &str, payload: Value) { + sqlx::query( + "INSERT INTO queue_messages (app_id, queue_name, payload, max_attempts) \ + VALUES ($1, $2, $3, 3)", + ) + .bind(Uuid::parse_str(app_id).expect("uuid")) + .bind(queue) + .bind(payload) + .execute(pool) + .await + .expect("enqueue"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn queue_receive_acks_on_success() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "ack").await; + + let script_id = create_script(&server, &app_id, "worker", ACK_HANDLER).await; + let resp = server + .post(&format!("/api/v1/admin/apps/{app_id}/triggers/queue")) + .json(&json!({ + "script_id": script_id, + "queue_name": "jobs", + "visibility_timeout_secs": 30 + })) + .await; + resp.assert_status(axum::http::StatusCode::CREATED); + + enqueue_directly(&pool, &app_id, "jobs", json!({ "x": 42 })).await; + + let marker = poll_marker(&pool, &app_id) + .await + .expect("handler fired"); + assert_eq!(marker["source"], "queue"); + assert_eq!(marker["queue"]["queue_name"], "jobs"); + assert_eq!(marker["queue"]["message"]["x"], 42); + + // Ack deleted the row. + assert_eq!(count_queue_messages(&pool, &app_id, "jobs").await, 0); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn queue_receive_dead_letters_after_max_attempts() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "dl").await; + + let script_id = create_script(&server, &app_id, "fail", THROW_HANDLER).await; + let resp = server + .post(&format!("/api/v1/admin/apps/{app_id}/triggers/queue")) + .json(&json!({ + "script_id": script_id, + "queue_name": "failing", + "visibility_timeout_secs": 30 + })) + .await; + resp.assert_status(axum::http::StatusCode::CREATED); + + enqueue_directly(&pool, &app_id, "failing", json!({ "x": 1 })).await; + + // Allow time for attempts to exhaust (3 retries × exponential + // backoff ≈ 7s in default config; we cap the poll loop generously). + for _ in 0..150 { + if count_dead_letters_for_queue(&pool, &app_id, "failing").await > 0 { + break; + } + tokio::time::sleep(Duration::from_millis(200)).await; + } + assert!( + count_dead_letters_for_queue(&pool, &app_id, "failing").await >= 1, + "expected a dead-letter row for the queue after max attempts" + ); + assert_eq!( + count_queue_messages(&pool, &app_id, "failing").await, + 0, + "queue row should be deleted after dead-letter" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn queue_one_consumer_per_queue_rejected() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "onecon").await; + + let s1 = create_script(&server, &app_id, "first", ACK_HANDLER).await; + let s2 = create_script(&server, &app_id, "second", ACK_HANDLER).await; + let resp = server + .post(&format!("/api/v1/admin/apps/{app_id}/triggers/queue")) + .json(&json!({ + "script_id": s1, + "queue_name": "single_owner", + "visibility_timeout_secs": 30 + })) + .await; + resp.assert_status(axum::http::StatusCode::CREATED); + + let resp2 = server + .post(&format!("/api/v1/admin/apps/{app_id}/triggers/queue")) + .json(&json!({ + "script_id": s2, + "queue_name": "single_owner", + "visibility_timeout_secs": 30 + })) + .await; + // The repo returns Invalid (422) for the duplicate consumer. + assert!( + resp2.status_code().is_client_error(), + "expected 4xx for duplicate consumer, got {}", + resp2.status_code() + ); + let body = resp2.text(); + assert!( + body.contains("already has a consumer") + || body.contains("queue 'single_owner'") + || body.contains("consumer trigger"), + "expected friendly conflict message, got: {body}" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn queue_visibility_timeout_reclaim() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool.clone(), "vt").await; + + let script_id = create_script(&server, &app_id, "slow", ACK_HANDLER).await; + // Minimum visibility timeout per the repo's validation. + let resp = server + .post(&format!("/api/v1/admin/apps/{app_id}/triggers/queue")) + .json(&json!({ + "script_id": script_id, + "queue_name": "vt_queue", + "visibility_timeout_secs": 5 + })) + .await; + resp.assert_status(axum::http::StatusCode::CREATED); + + // Insert a message and immediately fake a stale claim (claimed_at + // far enough in the past that the reclaim task picks it up). + let msg_id = Uuid::new_v4(); + sqlx::query( + "INSERT INTO queue_messages (id, app_id, queue_name, payload, max_attempts, \ + claim_token, claimed_at) \ + VALUES ($1, $2, 'vt_queue', $3, 3, $4, NOW() - INTERVAL '60 seconds')", + ) + .bind(msg_id) + .bind(Uuid::parse_str(&app_id).expect("uuid")) + .bind(json!({ "x": 99 })) + .bind(Uuid::new_v4()) + .execute(&pool) + .await + .expect("insert stale claim"); + + // The reclaim task default cadence is 30s; we wait up to 45s for it + // to fire OR observe the marker (whichever comes first — the marker + // shows up only after a successful claim by the live dispatcher). + let marker = poll_marker(&pool, &app_id).await; + if marker.is_some() { + return; // dispatcher already re-claimed and delivered + } + // Otherwise the reclaim path is the only way the message would + // have been picked up. Poll the claim state. + for _ in 0..50 { + let row: Option<(Option,)> = sqlx::query_as( + "SELECT claim_token FROM queue_messages WHERE id = $1", + ) + .bind(msg_id) + .fetch_optional(&pool) + .await + .expect("read"); + if let Some((token,)) = row { + if token.is_none() { + return; // reclaim succeeded — claim cleared + } + } else { + return; // already delivered + acked + } + tokio::time::sleep(Duration::from_secs(1)).await; + } + panic!("visibility timeout reclaim did not fire within 50s"); +} diff --git a/crates/picloud/tests/retry_e2e.rs b/crates/picloud/tests/retry_e2e.rs new file mode 100644 index 0000000..088ea21 --- /dev/null +++ b/crates/picloud/tests/retry_e2e.rs @@ -0,0 +1,175 @@ +//! v1.1.9 retry::* end-to-end test against a real engine inside the +//! all-in-one binary. Smaller surface than queue/invoke (no async +//! plumbing) — covers policy clamping, success-on-Nth-attempt, and +//! on_codes filtering through a real HTTP route. +//! +//! Skips when DATABASE_URL is unset. + +#![allow(clippy::needless_pass_by_value)] + +use axum_test::TestServer; +use serde_json::{json, Value}; +use sqlx::postgres::PgPoolOptions; +use sqlx::PgPool; +use uuid::Uuid; + +async fn pool_or_skip() -> Option { + let Ok(url) = std::env::var("DATABASE_URL") else { + eprintln!("retry_e2e: DATABASE_URL unset — skipping"); + return None; + }; + let pool = PgPoolOptions::new() + .max_connections(2) + .connect(&url) + .await + .expect("connect"); + sqlx::migrate!("../manager-core/migrations") + .run(&pool) + .await + .expect("migrate"); + Some(pool) +} + +async fn server_for(pool: PgPool, suffix: &str) -> (TestServer, String) { + use picloud_manager_core::auth::hash_password; + use picloud_shared::InstanceRole; + + let unique = format!("{suffix}-{}", Uuid::new_v4().simple()); + let auth = picloud::AuthDeps::from_pool(pool.clone()); + let username = format!("e2e-{unique}"); + let hash = hash_password("pw").expect("hash"); + auth.users + .create(&username, &hash, InstanceRole::Owner, None) + .await + .expect("seed admin"); + + let app = picloud::build_app( + pool, + auth, + picloud_shared::MasterKey::from_bytes([0x42u8; 32]), + ) + .await + .expect("build_app"); + let mut server = TestServer::new(app).expect("TestServer"); + let resp = server + .post("/api/v1/admin/auth/login") + .json(&json!({ "username": username, "password": "pw" })) + .await; + resp.assert_status_ok(); + let token = resp.json::()["token"] + .as_str() + .expect("token") + .to_string(); + server.add_header("authorization", format!("Bearer {token}")); + + let slug = format!("e2e-{unique}"); + let created: Value = server + .post("/api/v1/admin/apps") + .json(&json!({ "slug": slug, "name": slug })) + .await + .json(); + let app_id = created["id"].as_str().expect("app id").to_string(); + (server, app_id) +} + +async fn create_route_for( + server: &TestServer, + app_id: &str, + name: &str, + source: &str, + path: &str, +) { + let created: Value = server + .post("/api/v1/admin/scripts") + .json(&json!({ "app_id": app_id, "name": name, "source": source })) + .await + .json(); + let script_id = created["id"].as_str().expect("script id"); + let resp = server + .post(&format!("/api/v1/admin/apps/{app_id}/routes")) + .json(&json!({ + "script_id": script_id, + "method": "GET", + "path": path, + "host": null, + "dispatch_mode": "sync" + })) + .await; + resp.assert_status(axum::http::StatusCode::CREATED); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn retry_run_eventually_succeeds_inside_http_handler() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool, "retry-succ").await; + + // attempt counter mutates across retries; the 3rd attempt returns + // the count. + let src = r#" + let attempts = 0; + let p = retry::policy(#{ max_attempts: 5, base_ms: 1, jitter_pct: 0 }); + let v = retry::run(p, || { + attempts += 1; + if attempts < 3 { throw "transient" } else { attempts } + }); + #{ statusCode: 200, body: v } + "#; + create_route_for(&server, &app_id, "succ", src, "/r/succ").await; + + let resp = server.get("/r/succ").await; + resp.assert_status_ok(); + let body: Value = resp.json(); + assert_eq!(body, json!(3)); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn retry_run_surfaces_last_error_after_max_attempts() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool, "retry-surf").await; + + let src = r#" + let p = retry::policy(#{ max_attempts: 2, base_ms: 1, jitter_pct: 0 }); + try { + retry::run(p, || { throw "boom" }); + #{ statusCode: 200, body: "did not throw" } + } catch(e) { + #{ statusCode: 200, body: e } + } + "#; + create_route_for(&server, &app_id, "surf", src, "/r/surf").await; + + let resp = server.get("/r/surf").await; + resp.assert_status_ok(); + let body: Value = resp.json(); + let s = body.as_str().unwrap_or_default(); + assert!(s.contains("boom"), "expected 'boom' in error, got: {s}"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn retry_on_codes_filters_unmatched_errors() { + let Some(pool) = pool_or_skip().await else { return }; + let (server, app_id) = server_for(pool, "retry-codes").await; + + // attempts is mutated on every entry; we throw a non-matching code + // → must surface immediately (attempts == 1). + let src = r#" + let attempts = 0; + let p = retry::policy(#{ max_attempts: 5, base_ms: 1, jitter_pct: 0 }); + let p2 = retry::on_codes(p, ["http: 503"]); + try { + retry::run(p2, || { + attempts += 1; + throw "other failure" + }); + } catch(_e) { + // swallow — assert via attempts count. + } + #{ statusCode: 200, body: attempts } + "#; + create_route_for(&server, &app_id, "codes", src, "/r/codes").await; + + let resp = server.get("/r/codes").await; + resp.assert_status_ok(); + let body: Value = resp.json(); + assert_eq!(body, json!(1), "non-matching error should surface immediately"); +}