//! 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"); }