Files
wasmeld/crates/wasmeld-console/tests/api.rs
T

1042 lines
34 KiB
Rust
Raw Normal View History

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");
assert_eq!(registered["capabilities"], json!([]));
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 reports_exact_component_host_capabilities() {
let artifact_dir = TempDir::new().expect("temporary artifact directory");
let application = test_app(artifact_dir.path()).await;
let component = fs::read(component_artifact("clock_probe_component.wasm"))
.expect("clock probe component should be readable");
let response = application
.clone()
.oneshot(package_request("clock-probe", "0.1.0", &component))
.await
.expect("register request should complete");
assert_eq!(response.status(), StatusCode::CREATED);
let registered = response_json(response).await;
assert_eq!(
registered["capabilities"],
json!([{
"interface": "wasmeld:clock/monotonic-clock@0.1.0",
"package": "wasmeld:clock",
"name": "monotonic-clock",
"version": "0.1.0"
}])
);
let response = application
.oneshot(get_request("/api/v1/services"))
.await
.expect("service list request should complete");
let services = response_json(response).await;
assert_eq!(
services["services"][0]["capabilities"],
registered["capabilities"]
);
}
#[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 unregisters_only_inactive_service_revisions_and_prunes_persistence() {
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;
for revision in ["0.1.0", "0.2.0"] {
let response = application
.clone()
.oneshot(package_request("dev-echo", revision, &component))
.await
.expect("register request should complete");
assert_eq!(response.status(), StatusCode::CREATED);
}
let response = application
.clone()
.oneshot(json_request(
"/api/v1/deployments/dev-echo/activate",
json!({ "revision": "0.1.0" }),
))
.await
.expect("activation request should complete");
assert_eq!(response.status(), StatusCode::OK);
let response = application
.clone()
.oneshot(delete_request("/api/v1/services/dev-echo/0.1.0"))
.await
.expect("active unregister request should complete");
assert_eq!(response.status(), StatusCode::CONFLICT);
assert_eq!(
response_json(response).await["error"]["code"],
"active_deployment"
);
let response = application
.clone()
.oneshot(json_request(
"/api/v1/deployments/dev-echo/activate",
json!({ "revision": "0.2.0" }),
))
.await
.expect("replacement activation should complete");
assert_eq!(response.status(), StatusCode::OK);
let response = application
.clone()
.oneshot(delete_request("/api/v1/services/dev-echo/0.1.0"))
.await
.expect("inactive unregister request should complete");
assert_eq!(response.status(), StatusCode::OK);
assert!(!artifact_dir.path().join("dev-echo@0.1.0.wasm").exists());
assert!(!artifact_dir.path().join("dev-echo@0.1.0.toml").exists());
}
let application = test_app(artifact_dir.path()).await;
let response = application
.oneshot(get_request("/api/v1/services"))
.await
.expect("restored service list request should complete");
let body = response_json(response).await;
let services = body["services"].as_array().unwrap();
assert_eq!(services.len(), 1);
assert_eq!(services[0]["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);
}
2026-07-29 19:37:17 +08:00
#[tokio::test]
async fn persists_service_scoped_kv_across_revisions_and_restarts() {
const KV_WORLD: &str = "component:kv-probe/kv-probe-component@0.1.0";
let artifact_dir = TempDir::new().expect("temporary artifact directory");
let component = fs::read(component_artifact("kv_probe_component.wasm"))
.expect("KV probe component should be readable");
{
let application = test_app(artifact_dir.path()).await;
for (id, revision) in [
("shared-kv", "0.1.0"),
("shared-kv", "0.2.0"),
("isolated-kv", "0.1.0"),
] {
let response = application
.clone()
.oneshot(package_request_with_world(
id, revision, KV_WORLD, &component,
))
.await
.expect("KV component registration should complete");
assert_eq!(response.status(), StatusCode::CREATED);
let registered = response_json(response).await;
assert_eq!(
registered["capabilities"][0]["interface"],
"wasmeld:kv/store@0.1.0"
);
let response = application
.clone()
.oneshot(empty_post(&format!(
"/api/v1/services/{id}/{revision}/start"
)))
.await
.expect("KV actor start should complete");
assert_eq!(response.status(), StatusCode::OK);
}
assert_eq!(
invoke_service(&application, "shared-kv", "0.1.0", b"set:answer:forty-two").await,
b""
);
assert_eq!(
invoke_service(&application, "shared-kv", "0.2.0", b"get:answer").await,
b"forty-two"
);
assert_eq!(
invoke_service(&application, "isolated-kv", "0.1.0", b"get:answer").await,
b""
);
}
let application = test_app(artifact_dir.path()).await;
let response = application
.clone()
.oneshot(empty_post("/api/v1/services/shared-kv/0.2.0/start"))
.await
.expect("restored KV actor start should complete");
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
invoke_service(&application, "shared-kv", "0.2.0", b"get:answer").await,
b"forty-two"
);
assert_eq!(
invoke_service(&application, "shared-kv", "0.2.0", b"delete:answer").await,
b""
);
assert_eq!(
invoke_service(&application, "shared-kv", "0.2.0", b"get:answer").await,
b""
);
}
#[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<Body> {
let world = if id == "counter" {
"component:counter/counter-component@0.1.0"
} else if id == "clock-probe" {
"component:clock-probe/clock-probe-component@0.1.0"
} else {
"component:echo/echo-component@0.1.0"
};
2026-07-29 19:37:17 +08:00
package_request_with_world(id, revision, world, component)
}
fn package_request_with_world(
id: &str,
revision: &str,
world: &str,
component: &[u8],
) -> Request<Body> {
let mut package = Cursor::new(Vec::new());
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<Body> {
multipart_package_request(
"/api/v1/services",
"component.wasmpkg",
"application/vnd.wasm.component-package",
package,
)
}
fn wit_package_request(package: &[u8]) -> Request<Body> {
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<Body> {
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<Body> {
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<Body> {
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<Body> {
Request::builder()
.method("POST")
.uri(uri)
.body(Body::empty())
.unwrap()
}
fn delete_request(uri: &str) -> Request<Body> {
Request::builder()
.method("DELETE")
.uri(uri)
.body(Body::empty())
.unwrap()
}
fn get_request(uri: &str) -> Request<Body> {
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")
}
2026-07-29 19:37:17 +08:00
async fn invoke_service(application: &Router, id: &str, revision: &str, input: &[u8]) -> Vec<u8> {
let response = application
.clone()
.oneshot(json_request(
&format!("/api/v1/services/{id}/{revision}/invoke"),
json!({ "input_base64": BASE64.encode(input) }),
))
.await
.expect("component invocation should complete");
let status = response.status();
let body = response_json(response).await;
assert_eq!(status, StatusCode::OK, "{body}");
BASE64
.decode(body["output_base64"].as_str().unwrap())
.unwrap()
}
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",
"components/clock-probe/Cargo.toml",
2026-07-29 19:37:17 +08:00
"components/kv-probe/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<PathBuf> = OnceLock::new();
ROOT.get_or_init(|| {
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.parent()
.and_then(Path::parent)
.expect("console crate must live in the workspace")
.to_path_buf()
})
.as_path()
}