feat(console): add runtime management and WIT Registry

- embed wasmeld-runtime behind an Axum management API
- persist service metadata, metrics, and events with Toasty and libSQL
- register immutable component artifacts with platform-owned limits
- publish, list, inspect, and download binary WIT package versions
- verify Runtime lifecycle, persistence, Registry immutability, and transitive fetch
This commit is contained in:
Maofeng
2026-07-27 05:01:26 +08:00
parent cfac7d85f0
commit 12992186a9
9 changed files with 4976 additions and 22 deletions
+27
View File
@@ -0,0 +1,27 @@
[package]
name = "wasmeld-console"
version = "0.1.0"
edition.workspace = true
rust-version = "1.95"
license.workspace = true
[dependencies]
axum.workspace = true
base64.workspace = true
serde.workspace = true
serde_json.workspace = true
semver.workspace = true
thiserror.workspace = true
toasty.workspace = true
toasty-driver-turso.workspace = true
tokio.workspace = true
toml.workspace = true
tower-http.workspace = true
tracing.workspace = true
tracing-subscriber.workspace = true
wasmeld-package = { path = "../wasmeld-package" }
wasmeld-runtime = { path = "../wasmeld-runtime" }
[dev-dependencies]
tempfile = "3.23.0"
tower = { version = "0.5.2", features = ["util"] }
+65
View File
@@ -0,0 +1,65 @@
# Wasmeld Console
`wasmeld-console``wasmeld-runtime` 的本地 HTTP 管理适配器。它负责:
- 接收 `.wasmpkg`,校验并保存 `component.wasm` 与运行时 manifest
- 发布、查询和下载不可变的二进制 WIT Package
- 启动、停止和重启内嵌的 `wasmeld-runtime`
- 启动、停止和重启常驻 Actor
- 转发二进制调用
- 暴露调用计数、状态和最近事件
Wasm 制品保存在本地文件系统。服务 manifest、调用计数和最近 256 条事件通过
[Toasty](https://github.com/tokio-rs/toasty) 写入本地 libSQL 数据库,默认路径是
`var/wasmeld/console.db`
`wasmeld-console` 启动时会在同一进程内创建 `wasmeld-runtime`、加载 Wasmtime Engine
并重新注册已保存的 Component。Runtime 停止后 Console HTTP 和数据库仍然在线;
再次启动 Runtime 会重新注册 Component,但 Actor 内存不做快照,服务运行状态统一
恢复为 `stopped`
## 启动
`wasmeld-console` 需要 Rust 1.95 或更高版本;Wasm 组件继续使用项目约定的 Rust 1.90
工具链构建。
```bash
cargo +stable run -p wasmeld-console
```
默认监听 `127.0.0.1:8080`
可用环境变量:
- `WASMELD_ADDR`
- `WASMELD_ARTIFACT_DIR`
- `WASMELD_DATABASE_PATH`
- `WASMELD_WIT_REGISTRY_DIR`
- `WASMELD_ALLOWED_ORIGINS`,多个 Origin 使用逗号分隔
- `RUST_LOG`
## API
| 方法 | 路径 | 用途 |
| --- | --- | --- |
| `GET` | `/healthz` | 健康检查 |
| `GET` | `/api/v1/runtime` | Runtime 状态 |
| `POST` | `/api/v1/runtime/start` | 启动 Runtime 并重新注册受管 Component |
| `POST` | `/api/v1/runtime/stop` | 停止 Runtime 并释放所有 Actor |
| `POST` | `/api/v1/runtime/restart` | 重建 Runtime,服务恢复为停止状态 |
| `GET` | `/api/v1/services` | 服务版本列表 |
| `POST` | `/api/v1/services` | multipart 上传单个 `package` 字段(`.wasmpkg` |
| `POST` | `/api/v1/services/{id}/{revision}/start` | 启动 Actor |
| `POST` | `/api/v1/services/{id}/{revision}/stop` | 停止 Actor |
| `POST` | `/api/v1/services/{id}/{revision}/restart` | 重启 Actor |
| `POST` | `/api/v1/services/{id}/{revision}/invoke` | Base64 二进制调用 |
| `GET` | `/api/v1/events` | 最近 256 条持久化事件 |
| `GET` | `/api/v1/wit/packages` | WIT 包版本列表 |
| `POST` | `/api/v1/wit/packages` | multipart 发布单个二进制 WIT `package` 字段 |
| `GET` | `/api/v1/wit/packages/{namespace}/{name}/{version}` | WIT 包元数据 |
| `GET` | `/api/v1/wit/packages/{namespace}/{name}/{version}/content` | 下载 WIT 包 |
组件包的构建与格式说明见
[`docs/design/wasmeld-component-package.md`](../../docs/design/wasmeld-component-package.md)。
WIT 依赖、Registry 和本地 replace 说明见
[`docs/design/wit-package-management.md`](../../docs/design/wit-package-management.md)。
File diff suppressed because it is too large Load Diff
+65
View File
@@ -0,0 +1,65 @@
//! Standalone Wasmeld management server entry point.
use std::{env, net::SocketAddr, path::PathBuf, sync::Arc};
use axum::http::HeaderValue;
use tracing::info;
use tracing_subscriber::EnvFilter;
use wasmeld_console::{Console, ConsoleConfig, app};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt()
.with_env_filter(
EnvFilter::try_from_default_env()
.unwrap_or_else(|_| EnvFilter::new("wasmeld_console=info")),
)
.init();
let address = env::var("WASMELD_ADDR")
.unwrap_or_else(|_| "127.0.0.1:8080".to_owned())
.parse::<SocketAddr>()?;
let artifact_dir = env::var_os("WASMELD_ARTIFACT_DIR")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("var/wasmeld/components"));
let database_path = env::var_os("WASMELD_DATABASE_PATH")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("var/wasmeld/console.db"));
let wit_registry_dir = env::var_os("WASMELD_WIT_REGISTRY_DIR")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("var/wasmeld/wit-packages"));
let allowed_origins = allowed_origins()?;
let console = Arc::new(
Console::new(ConsoleConfig {
artifact_dir,
database_path,
wit_registry_dir,
..ConsoleConfig::default()
})
.await?,
);
let application = app(console, allowed_origins);
let listener = tokio::net::TcpListener::bind(address).await?;
info!(%address, "Wasmeld listening");
axum::serve(listener, application)
.with_graceful_shutdown(shutdown_signal())
.await?;
Ok(())
}
fn allowed_origins() -> Result<Vec<HeaderValue>, Box<dyn std::error::Error>> {
let value = env::var("WASMELD_ALLOWED_ORIGINS")
.unwrap_or_else(|_| "http://localhost:3000,http://127.0.0.1:3000".to_owned());
value
.split(',')
.filter(|origin| !origin.trim().is_empty())
.map(|origin| origin.trim().parse::<HeaderValue>().map_err(Into::into))
.collect()
}
async fn shutdown_signal() {
if tokio::signal::ctrl_c().await.is_err() {
tracing::warn!("failed to install Ctrl+C handler");
}
}
+101
View File
@@ -0,0 +1,101 @@
//! Toasty models and transactional libSQL snapshots for Console control data.
//!
//! Component bytes and WIT packages remain filesystem artifacts. This module
//! persists only service manifests, metrics, and the bounded event history.
use std::path::Path;
use toasty::Db;
use toasty_driver_turso::Turso;
use tokio::sync::Mutex;
#[derive(Clone, Debug, toasty::Model)]
pub(crate) struct StoredService {
#[key]
pub service_key: String,
pub manifest_toml: String,
pub updated_at_ms: u64,
pub calls: u64,
pub errors: u64,
pub last_latency_ms: Option<u64>,
}
#[derive(Clone, Debug, toasty::Model)]
pub(crate) struct StoredEvent {
#[key]
pub id: u64,
pub timestamp_ms: u64,
pub kind: String,
pub service_id: Option<String>,
pub revision: Option<String>,
pub message: String,
}
/// Serialized access to the Toasty database connection.
pub(crate) struct Persistence {
db: Mutex<Db>,
}
impl Persistence {
/// Opens the database and installs the schema for a new file.
pub async fn open(path: &Path) -> toasty::Result<Self> {
let is_new_database = !path.exists();
let db = Db::builder()
.models(toasty::models!(StoredService, StoredEvent))
.build(Turso::file(path))
.await?;
if is_new_database {
db.push_schema().await?;
}
Ok(Self { db: Mutex::new(db) })
}
/// Loads the complete service set and retained event history.
pub async fn load(&self) -> toasty::Result<(Vec<StoredService>, Vec<StoredEvent>)> {
let mut db = self.db.lock().await;
let services = StoredService::all().exec(&mut *db).await?;
let events = StoredEvent::all().exec(&mut *db).await?;
Ok((services, events))
}
/// Atomically upserts a Console snapshot and prunes expired events.
pub async fn sync(
&self,
services: Vec<StoredService>,
events: Vec<StoredEvent>,
) -> toasty::Result<()> {
let mut db = self.db.lock().await;
let mut tx = db.transaction().await?;
for service in services {
StoredService::upsert_by_service_key(service.service_key)
.manifest_toml(service.manifest_toml)
.updated_at_ms(service.updated_at_ms)
.calls(service.calls)
.errors(service.errors)
.last_latency_ms(service.last_latency_ms)
.exec(&mut tx)
.await?;
}
if let Some(oldest_event) = events.last() {
StoredEvent::filter(StoredEvent::fields().id().lt(oldest_event.id))
.delete()
.exec(&mut tx)
.await?;
}
for event in events {
StoredEvent::upsert_by_id(event.id)
.timestamp_ms(event.timestamp_ms)
.kind(event.kind)
.service_id(event.service_id)
.revision(event.revision)
.message(event.message)
.exec(&mut tx)
.await?;
}
tx.commit().await
}
}
+338
View File
@@ -0,0 +1,338 @@
//! Filesystem-backed immutable WIT package Registry.
//!
//! Artifacts are stored as
//! `<root>/<namespace>/<name>/<version>/package.wasm`. The in-memory index is
//! rebuilt from artifact contents at startup; directory names and cached
//! metadata are never treated as the source of truth.
use std::{
collections::BTreeMap,
fs::{self, OpenOptions},
io::{self, Write},
path::{Path, PathBuf},
sync::{RwLock, RwLockReadGuard, RwLockWriteGuard},
};
use thiserror::Error;
use wasmeld_package::wit_package::{WitPackageError, WitPackageMetadata, inspect_wit_package};
const ARTIFACT_FILE: &str = "package.wasm";
/// Process-local index over immutable binary WIT packages.
///
/// The write lock covers both the duplicate check and file creation so
/// concurrent publishers cannot replace the same package version.
pub(crate) struct WitRegistry {
root: PathBuf,
max_package_bytes: usize,
packages: RwLock<BTreeMap<String, StoredWitPackage>>,
}
#[derive(Clone)]
struct StoredWitPackage {
metadata: WitPackageMetadata,
artifact: PathBuf,
}
/// Errors produced by WIT Registry storage and integrity checks.
#[derive(Debug, Error)]
pub enum WitRegistryError {
#[error("WIT package exceeds the {limit}-byte limit")]
TooLarge { limit: usize },
#[error("WIT package {0} is already published and cannot be overwritten")]
AlreadyExists(String),
#[error("WIT package {0} was not found")]
NotFound(String),
#[error(transparent)]
InvalidPackage(#[from] WitPackageError),
#[error("invalid WIT registry data: {0}")]
InvalidStorage(String),
#[error("failed to access {path}: {source}")]
Storage {
path: PathBuf,
#[source]
source: io::Error,
},
#[error("WIT registry lock was poisoned")]
LockPoisoned,
}
impl WitRegistry {
/// Opens a Registry and validates every persisted package before serving it.
pub(crate) fn open(
root: impl AsRef<Path>,
max_package_bytes: usize,
) -> Result<Self, WitRegistryError> {
let root = root.as_ref();
fs::create_dir_all(root).map_err(|source| WitRegistryError::Storage {
path: root.to_path_buf(),
source,
})?;
let root = fs::canonicalize(root).map_err(|source| WitRegistryError::Storage {
path: root.to_path_buf(),
source,
})?;
let registry = Self {
root,
max_package_bytes,
packages: RwLock::new(BTreeMap::new()),
};
registry.load()?;
Ok(registry)
}
pub(crate) fn max_package_bytes(&self) -> usize {
self.max_package_bytes
}
/// Publishes one immutable binary WIT package.
///
/// Identity and dependency metadata are decoded from `bytes`; callers do
/// not provide trusted sidecar metadata.
pub(crate) fn publish(&self, bytes: &[u8]) -> Result<WitPackageMetadata, WitRegistryError> {
if bytes.len() > self.max_package_bytes {
return Err(WitRegistryError::TooLarge {
limit: self.max_package_bytes,
});
}
let metadata = inspect_wit_package(bytes)?;
let (namespace, name) = package_segments(&metadata.name)?;
let key = package_key(&metadata.name, &metadata.version);
let artifact = self
.root
.join(namespace)
.join(name)
.join(&metadata.version)
.join(ARTIFACT_FILE);
let mut packages = self.packages_write()?;
if packages.contains_key(&key) || artifact.exists() {
return Err(WitRegistryError::AlreadyExists(key));
}
let parent = artifact
.parent()
.expect("a package artifact always has a parent");
fs::create_dir_all(parent).map_err(|source| WitRegistryError::Storage {
path: parent.to_path_buf(),
source,
})?;
// create_new is the filesystem-level immutability guard. It also
// protects against an artifact created outside the current process.
let write_result = (|| {
let mut file = OpenOptions::new()
.write(true)
.create_new(true)
.open(&artifact)
.map_err(|source| WitRegistryError::Storage {
path: artifact.clone(),
source,
})?;
file.write_all(bytes)
.and_then(|()| file.sync_all())
.map_err(|source| WitRegistryError::Storage {
path: artifact.clone(),
source,
})
})();
if let Err(error) = write_result {
let _ = fs::remove_file(&artifact);
return Err(error);
}
packages.insert(
key,
StoredWitPackage {
metadata: metadata.clone(),
artifact,
},
);
Ok(metadata)
}
/// Returns all package versions in deterministic key order.
pub(crate) fn list(&self) -> Result<Vec<WitPackageMetadata>, WitRegistryError> {
Ok(self
.packages_read()?
.values()
.map(|package| package.metadata.clone())
.collect())
}
/// Looks up metadata for one exact package version.
pub(crate) fn metadata(
&self,
name: &str,
version: &str,
) -> Result<WitPackageMetadata, WitRegistryError> {
let key = checked_key(name, version)?;
self.packages_read()?
.get(&key)
.map(|package| package.metadata.clone())
.ok_or(WitRegistryError::NotFound(key))
}
/// Reads an exact package version and revalidates its content before serving it.
pub(crate) fn download(
&self,
name: &str,
version: &str,
) -> Result<(WitPackageMetadata, Vec<u8>), WitRegistryError> {
let key = checked_key(name, version)?;
let package = self
.packages_read()?
.get(&key)
.cloned()
.ok_or_else(|| WitRegistryError::NotFound(key.clone()))?;
let bytes = fs::read(&package.artifact).map_err(|source| WitRegistryError::Storage {
path: package.artifact,
source,
})?;
// Detect manual file replacement or corruption after startup instead
// of returning bytes inconsistent with the indexed digest.
let inspected = inspect_wit_package(&bytes)?;
if inspected != package.metadata {
return Err(WitRegistryError::InvalidStorage(format!(
"stored artifact metadata changed for {key}"
)));
}
Ok((package.metadata, bytes))
}
fn load(&self) -> Result<(), WitRegistryError> {
let mut loaded = BTreeMap::new();
for namespace in directories(&self.root)? {
for name in directories(&namespace)? {
for version in directories(&name)? {
let artifact = version.join(ARTIFACT_FILE);
if !artifact.is_file() {
continue;
}
let bytes =
fs::read(&artifact).map_err(|source| WitRegistryError::Storage {
path: artifact.clone(),
source,
})?;
if bytes.len() > self.max_package_bytes {
return Err(WitRegistryError::TooLarge {
limit: self.max_package_bytes,
});
}
// The artifact determines identity. Directory segments are
// checked afterward only as a storage-layout invariant.
let metadata = inspect_wit_package(&bytes)?;
let (expected_namespace, expected_name) = package_segments(&metadata.name)?;
let stored_namespace = namespace
.file_name()
.and_then(|value| value.to_str())
.unwrap_or_default();
let stored_name = name
.file_name()
.and_then(|value| value.to_str())
.unwrap_or_default();
let stored_version = version
.file_name()
.and_then(|value| value.to_str())
.unwrap_or_default();
if stored_namespace != expected_namespace
|| stored_name != expected_name
|| stored_version != metadata.version
{
return Err(WitRegistryError::InvalidStorage(format!(
"{} does not match package {}@{}",
artifact.display(),
metadata.name,
metadata.version
)));
}
let key = package_key(&metadata.name, &metadata.version);
if loaded
.insert(key.clone(), StoredWitPackage { metadata, artifact })
.is_some()
{
return Err(WitRegistryError::InvalidStorage(format!(
"duplicate package {key}"
)));
}
}
}
}
*self.packages_write()? = loaded;
Ok(())
}
fn packages_read(
&self,
) -> Result<RwLockReadGuard<'_, BTreeMap<String, StoredWitPackage>>, WitRegistryError> {
self.packages
.read()
.map_err(|_| WitRegistryError::LockPoisoned)
}
fn packages_write(
&self,
) -> Result<RwLockWriteGuard<'_, BTreeMap<String, StoredWitPackage>>, WitRegistryError> {
self.packages
.write()
.map_err(|_| WitRegistryError::LockPoisoned)
}
}
fn directories(path: &Path) -> Result<Vec<PathBuf>, WitRegistryError> {
let mut directories = fs::read_dir(path)
.map_err(|source| WitRegistryError::Storage {
path: path.to_path_buf(),
source,
})?
.filter_map(|entry| entry.ok())
.filter_map(|entry| {
entry
.file_type()
.ok()
.filter(|file_type| file_type.is_dir())
.map(|_| entry.path())
})
.collect::<Vec<_>>();
directories.sort();
Ok(directories)
}
fn checked_key(name: &str, version: &str) -> Result<String, WitRegistryError> {
package_segments(name)?;
semver::Version::parse(version)
.map_err(|error| WitRegistryError::InvalidStorage(error.to_string()))?;
Ok(package_key(name, version))
}
fn package_segments(name: &str) -> Result<(&str, &str), WitRegistryError> {
let Some((namespace, package)) = name.split_once(':') else {
return Err(WitRegistryError::InvalidStorage(format!(
"invalid WIT package name {name:?}"
)));
};
if package.contains(':') || !valid_segment(namespace) || !valid_segment(package) {
return Err(WitRegistryError::InvalidStorage(format!(
"invalid WIT package name {name:?}"
)));
}
Ok((namespace, package))
}
fn valid_segment(value: &str) -> bool {
!value.is_empty()
&& value.len() <= 128
&& !value.starts_with('-')
&& !value.ends_with('-')
&& value
.bytes()
.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
}
fn package_key(name: &str, version: &str) -> String {
format!("{name}@{version}")
}
+618
View File
@@ -0,0 +1,618 @@
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};
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 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 {
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");
app(Arc::new(console), Vec::new())
}
fn package_request(id: &str, revision: &str, component: &[u8]) -> Request<Body> {
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<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 empty_post(uri: &str) -> Request<Body> {
Request::builder()
.method("POST")
.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")
}
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<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()
}