`dispatcher_e2e` failed about one run in three, on a different test each time, and never under `--test-threads=1`. Root cause, not timing: Each test calls `build_app`, which spawns a REAL dispatcher, and they all shared one database — so they shared one `outbox`. `OutboxRepo::claim_due` is deliberately not app-scoped (in production one dispatcher serves the whole instance, so claiming any due row is correct). Test A's dispatcher would therefore claim test B's row, test A would finish, its server would be dropped mid-dispatch, and the claim was stranded. Test B polled for its handler's side effect until it timed out. The harness was modelling something production never does: N independent instances sharing one database. So the fix belongs in the harness. Weakening the dispatcher (app-scoping `claim_due`) would be wrong, and shortening `PICLOUD_OUTBOX_CLAIM_TIMEOUT_SEC` to fit a 10s test window would make production abandon the live claims of scripts that may legitimately run 300s. `tests/common` now hands each test its own database. Replaying 73 migrations per test is too slow, so the first caller builds a migrated TEMPLATE database once and every test clones it (`CREATE DATABASE ... TEMPLATE` is a file copy), guarded by a cluster-wide advisory lock because Postgres will not clone a template that has an open connection. Names are derived from (suite, test), so a test reclaims its own database on entry: a crashed run leaves nothing to clean up. Only the five local `pool_or_skip` wrappers change — all 33 call sites are untouched, and the skip-when-`DATABASE_URL`-is-unset contract is preserved. (`#[sqlx::test]` already does this and is what `api.rs`/`authz.rs` use, but it requires `DATABASE_URL` and would force `#[ignore]`, silently dropping these tests from the default `cargo test --workspace` gate.) This also stops the e2e tests leaving stranded claims behind in the shared dev database, which is the likely source of the occasional lone journey failure. Before: dispatcher_e2e red ~1 run in 3. After: workspace green twice over, journeys 157/157. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
159 lines
5.1 KiB
Rust
159 lines
5.1 KiB
Rust
//! 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 the admin bypass POST /api/v1/execute/{id}
|
|
//! (sidesteps the per-app domain matcher).
|
|
//!
|
|
//! Skips when DATABASE_URL is unset.
|
|
|
|
#![allow(clippy::needless_pass_by_value)]
|
|
|
|
use axum_test::TestServer;
|
|
use serde_json::{json, Value};
|
|
use sqlx::PgPool;
|
|
use uuid::Uuid;
|
|
|
|
mod common;
|
|
|
|
async fn pool_or_skip() -> Option<PgPool> {
|
|
common::test_pool("retry_e2e").await
|
|
}
|
|
|
|
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::<Value>()["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_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()
|
|
}
|
|
|
|
/// Execute via the admin bypass (same pattern as dispatcher_e2e).
|
|
async fn execute(server: &TestServer, script_id: &str) -> Value {
|
|
server
|
|
.post(&format!("/api/v1/execute/{script_id}"))
|
|
.json(&json!({}))
|
|
.await
|
|
.json()
|
|
}
|
|
|
|
#[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 }
|
|
"#;
|
|
let script_id = create_script(&server, &app_id, "succ", src).await;
|
|
let body = execute(&server, &script_id).await;
|
|
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 });
|
|
let out = "did not throw";
|
|
try {
|
|
retry::run(p, || { throw "boom" });
|
|
} catch(e) {
|
|
out = e;
|
|
}
|
|
#{ statusCode: 200, body: out }
|
|
"#;
|
|
let script_id = create_script(&server, &app_id, "surf", src).await;
|
|
let body = execute(&server, &script_id).await;
|
|
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 }
|
|
"#;
|
|
let script_id = create_script(&server, &app_id, "codes", src).await;
|
|
let body = execute(&server, &script_id).await;
|
|
assert_eq!(
|
|
body,
|
|
json!(1),
|
|
"non-matching error should surface immediately"
|
|
);
|
|
}
|