Files
wasmeld/crates/wasmeld-console/src/main.rs
T
Maofeng 68218b2eae perf(console): decouple invocation telemetry writes
Move Axum management and gateway adapters into a dedicated module while preserving the existing public routers and response contracts.

Replace per-invocation full snapshots with a bounded asynchronous writer that batches metric deltas and events. Serialize it with control snapshots, flush before lifecycle persistence and graceful shutdown, and keep invocation results independent from telemetry failures.

Enable the Runtime compilation cache beside the Console artifact directory and verify metric recovery across a lifecycle flush.
2026-07-30 07:18:25 +08:00

91 lines
3.3 KiB
Rust

//! Standalone Wasmeld control-plane and data-plane 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, gateway_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 management_address = env::var("WASMELD_ADDR")
.unwrap_or_else(|_| "127.0.0.1:8080".to_owned())
.parse::<SocketAddr>()?;
let gateway_address = env::var("WASMELD_GATEWAY_ADDR")
.unwrap_or_else(|_| "0.0.0.0:8081".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 management = app(Arc::clone(&console), allowed_origins);
let gateway = gateway_app(Arc::clone(&console));
let management_listener = tokio::net::TcpListener::bind(management_address).await?;
let gateway_listener = tokio::net::TcpListener::bind(gateway_address).await?;
let (shutdown_sender, shutdown_receiver) = tokio::sync::watch::channel(false);
let signal_task = tokio::spawn(async move {
shutdown_signal().await;
let _ = shutdown_sender.send(true);
});
info!(%management_address, "Wasmeld management API listening");
info!(%gateway_address, "Wasmeld Gateway listening");
let result = tokio::try_join!(
axum::serve(management_listener, management)
.with_graceful_shutdown(wait_for_shutdown(shutdown_receiver.clone())),
axum::serve(gateway_listener, gateway)
.with_graceful_shutdown(wait_for_shutdown(shutdown_receiver)),
);
signal_task.abort();
result?;
console.flush_invocation_telemetry().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");
}
}
async fn wait_for_shutdown(mut receiver: tokio::sync::watch::Receiver<bool>) {
while !*receiver.borrow() {
if receiver.changed().await.is_err() {
break;
}
}
}