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.
This commit is contained in:
Maofeng
2026-07-29 20:26:06 +08:00
parent 778fbd2a06
commit d0d2e06f68
3 changed files with 163 additions and 2 deletions
+68 -2
View File
@@ -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<Self, ConsoleError> {
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<ServiceView, ConsoleError> {
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<u8>) -> Result<ServiceView, ConsoleError> {
let runtime = self.runtime()?;
let runtime = runtime.as_ref().ok_or(ConsoleError::RuntimeNotRunning)?;
@@ -1221,7 +1271,7 @@ pub fn app(console: Arc<Console>, allowed_origins: Vec<HeaderValue>) -> 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<Console>, allowed_origins: Vec<HeaderValue>) -> 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<Arc<Console>>,
AxumPath((id, revision)): AxumPath<(String, String)>,
) -> Result<Json<ServiceView>, 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<Arc<Console>>,
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 { .. })
+19
View File
@@ -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::<BTreeSet<_>>();
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::<BTreeSet<_>>();
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)
+76
View File
@@ -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<Body> {
.unwrap()
}
fn delete_request(uri: &str) -> Request<Body> {
Request::builder()
.method("DELETE")
.uri(uri)
.body(Body::empty())
.unwrap()
}
fn get_request(uri: &str) -> Request<Body> {
Request::builder().uri(uri).body(Body::empty()).unwrap()
}