d0d2e06f68
Expose a management DELETE endpoint that refuses to remove an active Deployment, delegates Actor and compiled Component cleanup to Runtime, and removes stored artifacts. Reconcile Toasty snapshots by deleting absent service and Deployment rows so removed revisions do not reappear after restart. Persist unregister lifecycle events and verify active-revision protection, artifact deletion, Runtime cleanup, and restart behavior.
1042 lines
34 KiB
Rust
1042 lines
34 KiB
Rust
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);
|
|
}
|
|
|
|
#[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"
|
|
};
|
|
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")
|
|
}
|
|
|
|
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",
|
|
"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()
|
|
}
|