use std::{ fs, io::Cursor, path::{Path, PathBuf}, process::Command, sync::{Arc, OnceLock}, }; use axum::{ Router, body::{Body, to_bytes}, http::{Request, StatusCode, header}, }; use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64}; use serde_json::{Value, json}; use tempfile::TempDir; use tower::ServiceExt; use wasmeld_console::{Console, ConsoleConfig, app, gateway_app}; use wasmeld_package::{ module::{ModuleLock, sync_dependencies}, wit_package::build_wit_package, write_package, }; static COMPONENTS_BUILT: OnceLock<()> = OnceLock::new(); #[tokio::test] async fn manages_a_resident_component_over_http() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let application = test_app(artifact_dir.path()).await; let component = fs::read(echo_component()).expect("echo component should be readable"); let response = application .clone() .oneshot(package_request("echo", "0.1.0", &component)) .await .expect("register request should complete"); let status = response.status(); let registered = response_json(response).await; assert_eq!(status, StatusCode::CREATED, "{registered}"); assert_eq!(registered["id"], "echo"); assert_eq!(registered["status"], "stopped"); let response = application .clone() .oneshot(empty_post("/api/v1/services/echo/0.1.0/start")) .await .expect("start request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = application .clone() .oneshot(json_request( "/api/v1/services/echo/0.1.0/invoke", json!({ "input_base64": BASE64.encode(b"resident") }), )) .await .expect("invoke request should complete"); assert_eq!(response.status(), StatusCode::OK); let invocation = response_json(response).await; assert_eq!( BASE64 .decode(invocation["output_base64"].as_str().unwrap()) .unwrap(), b"resident" ); let response = application .clone() .oneshot(empty_post("/api/v1/services/echo/0.1.0/stop")) .await .expect("stop request should complete"); assert_eq!(response.status(), StatusCode::OK); assert_eq!(response_json(response).await["status"], "stopped"); let response = application .clone() .oneshot(empty_post("/api/v1/services/echo/0.1.0/restart")) .await .expect("restart request should complete"); assert_eq!(response.status(), StatusCode::OK); assert_eq!(response_json(response).await["status"], "running"); let response = application .oneshot(get_request("/api/v1/events")) .await .expect("events request should complete"); let events = response_json(response).await; assert!(events["events"].as_array().unwrap().len() >= 5); } #[tokio::test] async fn routes_raw_gateway_calls_through_the_active_deployment() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let (management, gateway) = test_apps(artifact_dir.path()).await; let component = fs::read(echo_component()).expect("echo component should be readable"); let response = gateway .clone() .oneshot(raw_request( "/v1/services/echo/invoke", "application/octet-stream", b"not deployed", )) .await .expect("gateway request should complete"); assert_eq!(response.status(), StatusCode::NOT_FOUND); assert_eq!( response_json(response).await["error"]["code"], "service_not_found" ); let response = gateway .clone() .oneshot(get_request("/api/v1/services")) .await .expect("gateway route isolation request should complete"); assert_eq!(response.status(), StatusCode::NOT_FOUND); for revision in ["0.1.0", "0.2.0"] { let response = management .clone() .oneshot(package_request("echo", revision, &component)) .await .expect("register request should complete"); assert_eq!(response.status(), StatusCode::CREATED); } let response = management .clone() .oneshot(json_request( "/api/v1/deployments/echo/activate", json!({ "revision": "0.1.0" }), )) .await .expect("activation request should complete"); assert_eq!(response.status(), StatusCode::OK); let deployment = response_json(response).await; assert_eq!(deployment["active_revision"], "0.1.0"); assert_eq!(deployment["status"], "running"); let response = gateway .clone() .oneshot(raw_request( "/v1/services/echo/invoke", "application/octet-stream", b"public bytes", )) .await .expect("gateway invocation should complete"); assert_eq!(response.status(), StatusCode::OK); assert_eq!( response.headers()["content-type"], "application/octet-stream" ); assert_eq!(response.headers()["x-wasmeld-revision"], "0.1.0"); assert_eq!( to_bytes(response.into_body(), usize::MAX).await.unwrap(), b"public bytes".as_slice() ); let response = management .clone() .oneshot(json_request( "/api/v1/deployments/echo/activate", json!({ "revision": "0.2.0" }), )) .await .expect("switch request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = gateway .clone() .oneshot(raw_request( "/v1/services/echo/invoke", "application/octet-stream", b"switched", )) .await .expect("gateway invocation should complete"); assert_eq!(response.status(), StatusCode::OK); assert_eq!(response.headers()["x-wasmeld-revision"], "0.2.0"); let response = gateway .clone() .oneshot(raw_request( "/v1/services/echo/invoke", "application/json", b"{}", )) .await .expect("unsupported media request should complete"); assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE); assert_eq!( response_json(response).await["error"]["code"], "unsupported_media_type" ); let response = management .clone() .oneshot(empty_post("/api/v1/services/echo/0.2.0/stop")) .await .expect("active actor stop request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = gateway .oneshot(raw_request( "/v1/services/echo/invoke", "application/octet-stream", b"stopped", )) .await .expect("stopped deployment request should complete"); assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE); assert_eq!( response_json(response).await["error"]["code"], "service_unavailable" ); let response = management .oneshot(get_request("/api/v1/deployments")) .await .expect("deployment list request should complete"); let deployments = response_json(response).await; assert_eq!(deployments["deployments"].as_array().unwrap().len(), 1); assert_eq!(deployments["deployments"][0]["active_revision"], "0.2.0"); } #[tokio::test] async fn restores_the_active_deployment_with_a_fresh_actor() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let component = fs::read(echo_component()).expect("echo component should be readable"); { let (management, gateway) = test_apps(artifact_dir.path()).await; let response = management .clone() .oneshot(package_request("restored", "1.0.0", &component)) .await .expect("register request should complete"); assert_eq!(response.status(), StatusCode::CREATED); let response = management .oneshot(json_request( "/api/v1/deployments/restored/activate", json!({ "revision": "1.0.0" }), )) .await .expect("activation request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = gateway .oneshot(raw_request( "/v1/services/restored/invoke", "application/octet-stream", b"before restart", )) .await .expect("gateway invocation should complete"); assert_eq!(response.status(), StatusCode::OK); } let (management, gateway) = test_apps(artifact_dir.path()).await; let response = management .oneshot(get_request("/api/v1/deployments/restored")) .await .expect("restored deployment request should complete"); let deployment = response_json(response).await; assert_eq!(deployment["active_revision"], "1.0.0"); assert_eq!(deployment["status"], "running"); let response = gateway .oneshot(raw_request( "/v1/services/restored/invoke", "application/octet-stream", b"after restart", )) .await .expect("restored gateway invocation should complete"); assert_eq!(response.status(), StatusCode::OK); assert_eq!( to_bytes(response.into_body(), usize::MAX).await.unwrap(), b"after restart".as_slice() ); } #[tokio::test] async fn reloads_manifests_without_restoring_actor_memory() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let component = fs::read(echo_component()).expect("echo component should be readable"); { let application = test_app(artifact_dir.path()).await; let response = application .clone() .oneshot(package_request("persisted", "0.1.0", &component)) .await .expect("register request should complete"); assert_eq!(response.status(), StatusCode::CREATED); let response = application .clone() .oneshot(empty_post("/api/v1/services/persisted/0.1.0/start")) .await .expect("start request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = application .oneshot(json_request( "/api/v1/services/persisted/0.1.0/invoke", json!({ "input_base64": BASE64.encode(b"persist metrics") }), )) .await .expect("invoke request should complete"); assert_eq!(response.status(), StatusCode::OK); } let application = test_app(artifact_dir.path()).await; let response = application .clone() .oneshot(get_request("/api/v1/services")) .await .expect("list request should complete"); assert_eq!(response.status(), StatusCode::OK); let body = response_json(response).await; assert_eq!(body["services"][0]["id"], "persisted"); assert_eq!(body["services"][0]["status"], "stopped"); assert_eq!(body["services"][0]["calls"], 1); let response = application .oneshot(get_request("/api/v1/events")) .await .expect("events request should complete"); let body = response_json(response).await; assert!(body["events"].as_array().unwrap().len() >= 3); } #[tokio::test] async fn manages_the_embedded_runtime_lifecycle() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let application = test_app(artifact_dir.path()).await; let component = fs::read(echo_component()).expect("echo component should be readable"); let response = application .clone() .oneshot(get_request("/api/v1/runtime")) .await .expect("runtime status request should complete"); let body = response_json(response).await; assert_eq!(body["status"], "running"); assert_eq!(body["managed_services"], 0); let response = application .clone() .oneshot(package_request("managed-runtime", "0.1.0", &component)) .await .expect("register request should complete"); assert_eq!(response.status(), StatusCode::CREATED); let response = application .clone() .oneshot(empty_post("/api/v1/runtime/stop")) .await .expect("runtime stop request should complete"); let body = response_json(response).await; assert_eq!(body["status"], "stopped"); assert_eq!(body["managed_services"], 1); assert_eq!(body["registered_services"], 0); assert_eq!(body["running_services"], 0); let response = application .clone() .oneshot(json_request( "/api/v1/services/managed-runtime/0.1.0/invoke", json!({ "input_base64": BASE64.encode(b"stopped") }), )) .await .expect("invoke request should complete"); assert_eq!(response.status(), StatusCode::CONFLICT); assert_eq!( response_json(response).await["error"]["code"], "runtime_not_running" ); let response = application .clone() .oneshot(empty_post("/api/v1/runtime/start")) .await .expect("runtime start request should complete"); let body = response_json(response).await; assert_eq!(body["status"], "running"); assert_eq!(body["registered_services"], 1); assert_eq!(body["running_services"], 0); let response = application .clone() .oneshot(empty_post("/api/v1/services/managed-runtime/0.1.0/start")) .await .expect("actor start request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = application .clone() .oneshot(json_request( "/api/v1/services/managed-runtime/0.1.0/invoke", json!({ "input_base64": BASE64.encode(b"running again") }), )) .await .expect("invoke request should complete"); assert_eq!(response.status(), StatusCode::OK); let response = application .clone() .oneshot(empty_post("/api/v1/runtime/restart")) .await .expect("runtime restart request should complete"); let body = response_json(response).await; assert_eq!(body["status"], "running"); assert_eq!(body["registered_services"], 1); assert_eq!(body["running_services"], 0); let response = application .oneshot(get_request("/api/v1/services")) .await .expect("services request should complete"); let body = response_json(response).await; assert_eq!(body["services"][0]["status"], "stopped"); } #[tokio::test] async fn preserves_counter_memory_across_http_calls() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let application = test_app(artifact_dir.path()).await; let component = fs::read(component_artifact("counter_component.wasm")) .expect("counter component should be readable"); let response = application .clone() .oneshot(package_request("counter", "0.1.0", &component)) .await .expect("register request should complete"); let status = response.status(); let body = response_json(response).await; assert_eq!(status, StatusCode::CREATED, "{body}"); let response = application .clone() .oneshot(empty_post("/api/v1/services/counter/0.1.0/start")) .await .expect("counter start should complete"); assert_eq!(response.status(), StatusCode::OK); for expected in [1_u64, 2] { let response = application .clone() .oneshot(json_request( "/api/v1/services/counter/0.1.0/invoke", json!({ "input_base64": "" }), )) .await .expect("counter invocation should complete"); assert_eq!(response.status(), StatusCode::OK); let body = response_json(response).await; let bytes = BASE64 .decode(body["output_base64"].as_str().unwrap()) .unwrap(); assert_eq!( u64::from_le_bytes(bytes.try_into().expect("counter returns eight bytes")), expected ); } } #[tokio::test] async fn rejects_invalid_component_packages() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let application = test_app(artifact_dir.path()).await; let response = application .oneshot(package_bytes_request(b"not-a-package")) .await .expect("register request should complete"); assert_eq!(response.status(), StatusCode::UNPROCESSABLE_ENTITY); assert_eq!( response_json(response).await["error"]["code"], "invalid_package" ); } #[tokio::test] async fn publishes_and_downloads_immutable_wit_packages() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let package_source = TempDir::new().expect("temporary WIT source"); fs::write( package_source.path().join("package.wit"), "package wasmeld:test-clock@1.2.3;\ninterface clock { now: func() -> u64; }\n", ) .unwrap(); let package = build_wit_package(package_source.path()).unwrap(); let application = test_app(artifact_dir.path()).await; let response = application .clone() .oneshot(wit_package_request(&package.bytes)) .await .expect("publish request should complete"); let status = response.status(); let body = response_json(response).await; assert_eq!(status, StatusCode::CREATED, "{body}"); assert_eq!(body["name"], "wasmeld:test-clock"); assert_eq!(body["version"], "1.2.3"); assert_eq!(body["sha256"], package.metadata.sha256); let response = application .clone() .oneshot(wit_package_request(&package.bytes)) .await .expect("duplicate publish request should complete"); assert_eq!(response.status(), StatusCode::CONFLICT); let response = application .clone() .oneshot(get_request("/api/v1/wit/packages")) .await .unwrap(); let list = response_json(response).await; assert_eq!(list["packages"].as_array().unwrap().len(), 1); let response = application .clone() .oneshot(get_request("/api/v1/wit/packages/wasmeld/test-clock/1.2.3")) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); assert_eq!(response_json(response).await["name"], "wasmeld:test-clock"); let response = application .oneshot(get_request( "/api/v1/wit/packages/wasmeld/test-clock/1.2.3/content", )) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); assert_eq!( to_bytes(response.into_body(), usize::MAX).await.unwrap(), package.bytes ); let restarted = test_app(artifact_dir.path()).await; let response = restarted .oneshot(get_request("/api/v1/wit/packages")) .await .unwrap(); assert_eq!( response_json(response).await["packages"] .as_array() .unwrap() .len(), 1 ); } #[tokio::test] async fn rejects_invalid_binary_wit_packages() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let application = test_app(artifact_dir.path()).await; let response = application .oneshot(wit_package_request(b"not-a-wit-package")) .await .expect("publish request should complete"); assert_eq!(response.status(), StatusCode::UNPROCESSABLE_ENTITY); assert_eq!( response_json(response).await["error"]["code"], "invalid_wit_package" ); } #[tokio::test] async fn client_fetches_transitive_wit_dependencies_from_the_registry() { let artifact_dir = TempDir::new().expect("temporary artifact directory"); let sources = TempDir::new().expect("temporary WIT sources"); let clock = sources.path().join("clock"); let timer = sources.path().join("timer"); fs::create_dir(&clock).unwrap(); fs::create_dir_all(timer.join("deps/clock")).unwrap(); fs::write( clock.join("package.wit"), "package wasmeld:clock@1.0.0;\ninterface clock { now: func() -> u64; }\n", ) .unwrap(); fs::copy( clock.join("package.wit"), timer.join("deps/clock/package.wit"), ) .unwrap(); fs::write( timer.join("package.wit"), "package wasmeld:timer@1.0.0;\nworld timer-host {\n import wasmeld:clock/clock@1.0.0;\n}\n", ) .unwrap(); let clock_package = build_wit_package(&clock).unwrap(); let timer_package = build_wit_package(&timer).unwrap(); assert_eq!(timer_package.metadata.dependencies.len(), 1); let application = test_app(artifact_dir.path()).await; for package in [&clock_package.bytes, &timer_package.bytes] { let response = application .clone() .oneshot(wit_package_request(package)) .await .unwrap(); assert_eq!(response.status(), StatusCode::CREATED); } let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server = tokio::spawn(async move { axum::serve(listener, application).await.unwrap(); }); let component = sources.path().join("component"); fs::create_dir_all(component.join("wit")).unwrap(); fs::write( component.join("wit/world.wit"), "package example:registry-client@0.1.0;\nworld client {\n include wasmeld:timer/timer-host@1.0.0;\n}\n", ) .unwrap(); fs::write( component.join("wasmeld.toml"), format!( "schema_version = 1\n\n[registry]\nurl = \"http://{address}\"\n\n[dependencies]\n\"wasmeld:timer\" = \"1.0.0\"\n" ), ) .unwrap(); let manifest = component.join("wasmeld.toml"); let sync_manifest = manifest.clone(); let report = tokio::task::spawn_blocking(move || sync_dependencies(sync_manifest, false)) .await .unwrap() .unwrap(); assert_eq!(report.packages.len(), 2); assert!( component .join("wit/deps/wasmeld-clock-1.0.0.wasm") .is_file() ); assert!( component .join("wit/deps/wasmeld-timer-1.0.0.wasm") .is_file() ); let lock = ModuleLock::read(component.join("wit.lock")).unwrap(); assert_eq!(lock.packages.len(), 2); assert!( lock.packages .iter() .all(|package| package.source == format!("registry+http://{address}")) ); let locked_manifest = manifest.clone(); tokio::task::spawn_blocking(move || sync_dependencies(locked_manifest, true)) .await .unwrap() .unwrap(); server.abort(); } async fn test_app(artifact_dir: &Path) -> Router { test_apps(artifact_dir).await.0 } async fn test_apps(artifact_dir: &Path) -> (Router, Router) { let console = Console::new(ConsoleConfig { artifact_dir: artifact_dir.to_path_buf(), wit_registry_dir: artifact_dir.join("wit-packages"), database_path: artifact_dir.join("console.db"), ..ConsoleConfig::default() }) .await .expect("console should start"); let console = Arc::new(console); (app(Arc::clone(&console), Vec::new()), gateway_app(console)) } fn package_request(id: &str, revision: &str, component: &[u8]) -> Request
{ let mut package = Cursor::new(Vec::new()); let world = if id == "counter" { "component:counter/counter-component@0.1.0" } else { "component:echo/echo-component@0.1.0" }; write_package(&mut package, id, revision, world, component) .expect("test package should be valid"); package_bytes_request(package.get_ref()) } fn package_bytes_request(package: &[u8]) -> Request { multipart_package_request( "/api/v1/services", "component.wasmpkg", "application/vnd.wasm.component-package", package, ) } fn wit_package_request(package: &[u8]) -> Request { multipart_package_request( "/api/v1/wit/packages", "package.wasm", "application/wasm", package, ) } fn multipart_package_request( uri: &str, filename: &str, content_type: &str, package: &[u8], ) -> Request { const BOUNDARY: &str = "wasmeld-console-test-boundary"; let mut body = Vec::new(); body.extend_from_slice(format!("--{BOUNDARY}\r\n").as_bytes()); body.extend_from_slice( format!("Content-Disposition: form-data; name=\"package\"; filename=\"{filename}\"\r\n") .as_bytes(), ); body.extend_from_slice(format!("Content-Type: {content_type}\r\n\r\n").as_bytes()); body.extend_from_slice(package); body.extend_from_slice(b"\r\n"); body.extend_from_slice(format!("--{BOUNDARY}--\r\n").as_bytes()); Request::builder() .method("POST") .uri(uri) .header( header::CONTENT_TYPE, format!("multipart/form-data; boundary={BOUNDARY}"), ) .body(Body::from(body)) .unwrap() } fn json_request(uri: &str, value: Value) -> Request { Request::builder() .method("POST") .uri(uri) .header(header::CONTENT_TYPE, "application/json") .body(Body::from(value.to_string())) .unwrap() } fn raw_request(uri: &str, content_type: &str, body: &[u8]) -> Request { Request::builder() .method("POST") .uri(uri) .header(header::CONTENT_TYPE, content_type) .body(Body::from(body.to_vec())) .unwrap() } fn empty_post(uri: &str) -> Request { Request::builder() .method("POST") .uri(uri) .body(Body::empty()) .unwrap() } fn get_request(uri: &str) -> Request { Request::builder().uri(uri).body(Body::empty()).unwrap() } async fn response_json(response: axum::response::Response) -> Value { let bytes = to_bytes(response.into_body(), usize::MAX) .await .expect("response body should be readable"); serde_json::from_slice(&bytes).expect("response should be JSON") } fn echo_component() -> PathBuf { component_artifact("echo_component.wasm") } fn component_artifact(name: &str) -> PathBuf { COMPONENTS_BUILT.get_or_init(|| { for manifest in [ "components/echo/Cargo.toml", "components/counter/Cargo.toml", ] { let status = Command::new("rustup") .args([ "run", "1.90.0", "cargo", "build", "--manifest-path", manifest, "--target", "wasm32-wasip2", "--release", ]) .current_dir(workspace_root()) .status() .expect("component build command should start"); assert!(status.success(), "component build should succeed"); } }); workspace_root() .join("target/wasm32-wasip2/release") .join(name) } fn workspace_root() -> &'static Path { static ROOT: OnceLock