From d0d2e06f689699a9fe1aae2f1a9f4a71f7559ccf Mon Sep 17 00:00:00 2001 From: Maofeng Date: Wed, 29 Jul 2026 20:26:06 +0800 Subject: [PATCH] feat(console): unregister inactive revisions 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. --- crates/wasmeld-console/src/lib.rs | 70 ++++++++++++++++++++- crates/wasmeld-console/src/persistence.rs | 19 ++++++ crates/wasmeld-console/tests/api.rs | 76 +++++++++++++++++++++++ 3 files changed, 163 insertions(+), 2 deletions(-) diff --git a/crates/wasmeld-console/src/lib.rs b/crates/wasmeld-console/src/lib.rs index 229233b..df08951 100644 --- a/crates/wasmeld-console/src/lib.rs +++ b/crates/wasmeld-console/src/lib.rs @@ -24,7 +24,7 @@ use axum::{ extract::{DefaultBodyLimit, Multipart, Path as AxumPath, State, rejection::BytesRejection}, http::{HeaderMap, HeaderValue, Method, StatusCode, header}, response::{IntoResponse, Response}, - routing::{get, post}, + routing::{delete, get, post}, }; use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64}; use serde::{Deserialize, Serialize}; @@ -218,6 +218,7 @@ pub struct EventView { #[serde(rename_all = "snake_case")] pub enum EventKind { Registered, + Unregistered, Deployed, Started, Stopped, @@ -229,6 +230,7 @@ impl EventKind { fn as_database_value(self) -> &'static str { match self { Self::Registered => "registered", + Self::Unregistered => "unregistered", Self::Deployed => "deployed", Self::Started => "started", Self::Stopped => "stopped", @@ -240,6 +242,7 @@ impl EventKind { fn from_database_value(value: &str) -> Result { match value { "registered" => Ok(Self::Registered), + "unregistered" => Ok(Self::Unregistered), "deployed" => Ok(Self::Deployed), "started" => Ok(Self::Started), "stopped" => Ok(Self::Stopped), @@ -351,6 +354,9 @@ pub enum ConsoleError { #[error("service {0} has no active deployment")] DeploymentNotFound(String), + #[error("service revision {0} is the active deployment")] + ActiveDeployment(ServiceKey), + #[error("Wasmeld Runtime is not running")] RuntimeNotRunning, @@ -790,6 +796,50 @@ impl Console { Ok(view) } + fn unregister(&self, key: &ServiceKey) -> Result { + let _lifecycle = self + .lifecycle + .lock() + .map_err(|_| ConsoleError::LockPoisoned)?; + let view = { + let state = self.state()?; + if state + .deployments + .get(key.id()) + .is_some_and(|deployment| deployment.active_revision == key.revision()) + { + return Err(ConsoleError::ActiveDeployment(key.clone())); + } + state + .services + .get(key) + .map(ServiceRecord::view) + .ok_or_else(|| ConsoleError::ServiceNotFound(key.clone()))? + }; + + let runtime = self.runtime()?; + let runtime = runtime.as_ref().ok_or(ConsoleError::RuntimeNotRunning)?; + runtime.unregister(key)?; + self.state()?.services.remove(key); + + for path in [ + self.artifact_dir.join(format!("{key}.wasm")), + self.artifact_dir.join(format!("{key}.toml")), + ] { + if let Err(error) = fs::remove_file(&path) + && error.kind() != io::ErrorKind::NotFound + { + tracing::warn!(path = %path.display(), %error, "failed to remove unregistered artifact"); + } + } + self.push_event( + EventKind::Unregistered, + Some(key), + "component revision unregistered".to_owned(), + )?; + Ok(view) + } + fn start(&self, key: &ServiceKey, init_config: Vec) -> Result { let runtime = self.runtime()?; let runtime = runtime.as_ref().ok_or(ConsoleError::RuntimeNotRunning)?; @@ -1221,7 +1271,7 @@ pub fn app(console: Arc, allowed_origins: Vec) -> Router { .max(console.max_wit_package_bytes()) .saturating_add(1024 * 1024); let mut cors = CorsLayer::new() - .allow_methods([Method::GET, Method::POST]) + .allow_methods([Method::GET, Method::POST, Method::DELETE]) .allow_headers([header::CONTENT_TYPE]); if !allowed_origins.is_empty() { cors = cors.allow_origin(allowed_origins); @@ -1234,6 +1284,10 @@ pub fn app(console: Arc, allowed_origins: Vec) -> Router { .route("/api/v1/runtime/stop", post(stop_runtime)) .route("/api/v1/runtime/restart", post(restart_runtime)) .route("/api/v1/services", get(list_services).post(register)) + .route( + "/api/v1/services/{id}/{revision}", + delete(unregister_service), + ) .route("/api/v1/deployments", get(list_deployments)) .route("/api/v1/deployments/{id}", get(get_deployment)) .route( @@ -1573,6 +1627,16 @@ async fn stop_service( )) } +async fn unregister_service( + State(console): State>, + AxumPath((id, revision)): AxumPath<(String, String)>, +) -> Result, ApiError> { + let key = api_key(id, revision)?; + Ok(Json( + run_blocking(console, move |console| console.unregister(&key)).await?, + )) +} + async fn restart_service( State(console): State>, AxumPath((id, revision)): AxumPath<(String, String)>, @@ -1766,6 +1830,7 @@ impl IntoResponse for GatewayError { | RuntimeError::ComponentError { .. }, ) => (StatusCode::BAD_GATEWAY, "service_error", error.to_string()), ConsoleError::ArtifactTooLarge { .. } + | ConsoleError::ActiveDeployment(_) | ConsoleError::Package(_) | ConsoleError::WitRegistry(_) | ConsoleError::Storage { .. } @@ -1818,6 +1883,7 @@ impl IntoResponse for ApiError { ConsoleError::InvalidRequest(_) => (StatusCode::BAD_REQUEST, "invalid_request"), ConsoleError::ServiceNotFound(_) => (StatusCode::NOT_FOUND, "service_not_found"), ConsoleError::DeploymentNotFound(_) => (StatusCode::NOT_FOUND, "deployment_not_found"), + ConsoleError::ActiveDeployment(_) => (StatusCode::CONFLICT, "active_deployment"), ConsoleError::RuntimeNotRunning => (StatusCode::CONFLICT, "runtime_not_running"), ConsoleError::ArtifactTooLarge { .. } | ConsoleError::Package(PackageError::ComponentTooLarge { .. }) diff --git a/crates/wasmeld-console/src/persistence.rs b/crates/wasmeld-console/src/persistence.rs index 391f031..a577e61 100644 --- a/crates/wasmeld-console/src/persistence.rs +++ b/crates/wasmeld-console/src/persistence.rs @@ -5,6 +5,7 @@ //! history, and service-scoped Host KV entries. use std::{ + collections::BTreeSet, path::Path, sync::Arc, time::{SystemTime, UNIX_EPOCH}, @@ -130,6 +131,15 @@ impl Persistence { let mut db = self.db.lock().await; let mut tx = db.transaction().await?; + let service_keys = services + .iter() + .map(|service| service.service_key.clone()) + .collect::>(); + for stored in StoredService::all().exec(&mut tx).await? { + if !service_keys.contains(&stored.service_key) { + StoredService::delete_by_service_key(&mut tx, &stored.service_key).await?; + } + } for service in services { StoredService::upsert_by_service_key(service.service_key) .manifest_toml(service.manifest_toml) @@ -159,6 +169,15 @@ impl Persistence { .await?; } + let deployment_ids = deployments + .iter() + .map(|deployment| deployment.service_id.clone()) + .collect::>(); + for stored in StoredDeployment::all().exec(&mut tx).await? { + if !deployment_ids.contains(&stored.service_id) { + StoredDeployment::delete_by_service_id(&mut tx, &stored.service_id).await?; + } + } for deployment in deployments { StoredDeployment::upsert_by_service_id(deployment.service_id) .active_revision(deployment.active_revision) diff --git a/crates/wasmeld-console/tests/api.rs b/crates/wasmeld-console/tests/api.rs index b3311bf..cce781d 100644 --- a/crates/wasmeld-console/tests/api.rs +++ b/crates/wasmeld-console/tests/api.rs @@ -261,6 +261,74 @@ async fn routes_raw_gateway_calls_through_the_active_deployment() { 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"); @@ -888,6 +956,14 @@ fn empty_post(uri: &str) -> Request { .unwrap() } +fn delete_request(uri: &str) -> Request { + Request::builder() + .method("DELETE") + .uri(uri) + .body(Body::empty()) + .unwrap() +} + fn get_request(uri: &str) -> Request { Request::builder().uri(uri).body(Body::empty()).unwrap() }