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() }