feat(console): add runtime management and WIT Registry
- embed wasmeld-runtime behind an Axum management API\n- persist service metadata, metrics, and events with Toasty and libSQL\n- register immutable component artifacts with platform-owned limits\n- publish, list, inspect, and download binary WIT package versions\n- verify Runtime lifecycle, persistence, Registry immutability, and transitive fetch
This commit is contained in:
Generated
+2369
-22
File diff suppressed because it is too large
Load Diff
@@ -2,6 +2,7 @@
|
||||
members = [
|
||||
"crates/wasmeld-package",
|
||||
"crates/wasmeld-runtime",
|
||||
"crates/wasmeld-console",
|
||||
"components/echo",
|
||||
"components/counter",
|
||||
"components/fault",
|
||||
@@ -18,13 +19,21 @@ rust-version = "1.90"
|
||||
license = "Apache-2.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
axum = { version = "0.8.9", features = ["multipart"] }
|
||||
base64 = "0.23.0"
|
||||
reqwest = { version = "0.13.4", default-features = false, features = ["blocking", "json", "multipart", "rustls"] }
|
||||
serde = { version = "1.0.228", features = ["derive"] }
|
||||
serde_json = "1.0.149"
|
||||
semver = "1.0.28"
|
||||
sha2 = "0.10.9"
|
||||
thiserror = "2.0.17"
|
||||
toasty = "=0.9.0"
|
||||
toasty-driver-turso = "=0.9.0"
|
||||
tokio = { version = "1.53.1", features = ["macros", "net", "rt-multi-thread", "signal", "sync"] }
|
||||
toml = "0.9.8"
|
||||
tower-http = { version = "0.7.0", features = ["cors", "trace"] }
|
||||
tracing = "0.1.44"
|
||||
tracing-subscriber = { version = "0.3.23", features = ["env-filter", "fmt"] }
|
||||
wasmtime = { version = "=41.0.0", default-features = false, features = ["component-model", "cranelift", "runtime", "std"] }
|
||||
wasmtime-wasi = { version = "=41.0.0", default-features = false, features = ["p2"] }
|
||||
wit-bindgen = "=0.41.0"
|
||||
|
||||
@@ -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"] }
|
||||
@@ -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
@@ -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");
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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}")
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
Reference in New Issue
Block a user