feat(console): add deployment-routed gateway
- persist service-to-revision deployments in libSQL with an additive migration - restore active actors and atomically switch revisions - expose an isolated raw-byte gateway on a second listener with stable errors - cover routing, isolation, revision switching, and restart restoration
This commit is contained in:
@@ -15,7 +15,7 @@ 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_console::{Console, ConsoleConfig, app, gateway_app};
|
||||
use wasmeld_package::{
|
||||
module::{ModuleLock, sync_dependencies},
|
||||
wit_package::build_wit_package,
|
||||
@@ -89,6 +89,199 @@ async fn manages_a_resident_component_over_http() {
|
||||
assert!(events["events"].as_array().unwrap().len() >= 5);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn routes_raw_gateway_calls_through_the_active_deployment() {
|
||||
let artifact_dir = TempDir::new().expect("temporary artifact directory");
|
||||
let (management, gateway) = test_apps(artifact_dir.path()).await;
|
||||
let component = fs::read(echo_component()).expect("echo component should be readable");
|
||||
|
||||
let response = gateway
|
||||
.clone()
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/echo/invoke",
|
||||
"application/octet-stream",
|
||||
b"not deployed",
|
||||
))
|
||||
.await
|
||||
.expect("gateway request should complete");
|
||||
assert_eq!(response.status(), StatusCode::NOT_FOUND);
|
||||
assert_eq!(
|
||||
response_json(response).await["error"]["code"],
|
||||
"service_not_found"
|
||||
);
|
||||
|
||||
let response = gateway
|
||||
.clone()
|
||||
.oneshot(get_request("/api/v1/services"))
|
||||
.await
|
||||
.expect("gateway route isolation request should complete");
|
||||
assert_eq!(response.status(), StatusCode::NOT_FOUND);
|
||||
|
||||
for revision in ["0.1.0", "0.2.0"] {
|
||||
let response = management
|
||||
.clone()
|
||||
.oneshot(package_request("echo", revision, &component))
|
||||
.await
|
||||
.expect("register request should complete");
|
||||
assert_eq!(response.status(), StatusCode::CREATED);
|
||||
}
|
||||
|
||||
let response = management
|
||||
.clone()
|
||||
.oneshot(json_request(
|
||||
"/api/v1/deployments/echo/activate",
|
||||
json!({ "revision": "0.1.0" }),
|
||||
))
|
||||
.await
|
||||
.expect("activation request should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let deployment = response_json(response).await;
|
||||
assert_eq!(deployment["active_revision"], "0.1.0");
|
||||
assert_eq!(deployment["status"], "running");
|
||||
|
||||
let response = gateway
|
||||
.clone()
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/echo/invoke",
|
||||
"application/octet-stream",
|
||||
b"public bytes",
|
||||
))
|
||||
.await
|
||||
.expect("gateway invocation should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
assert_eq!(
|
||||
response.headers()["content-type"],
|
||||
"application/octet-stream"
|
||||
);
|
||||
assert_eq!(response.headers()["x-wasmeld-revision"], "0.1.0");
|
||||
assert_eq!(
|
||||
to_bytes(response.into_body(), usize::MAX).await.unwrap(),
|
||||
b"public bytes".as_slice()
|
||||
);
|
||||
|
||||
let response = management
|
||||
.clone()
|
||||
.oneshot(json_request(
|
||||
"/api/v1/deployments/echo/activate",
|
||||
json!({ "revision": "0.2.0" }),
|
||||
))
|
||||
.await
|
||||
.expect("switch request should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
|
||||
let response = gateway
|
||||
.clone()
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/echo/invoke",
|
||||
"application/octet-stream",
|
||||
b"switched",
|
||||
))
|
||||
.await
|
||||
.expect("gateway invocation should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
assert_eq!(response.headers()["x-wasmeld-revision"], "0.2.0");
|
||||
|
||||
let response = gateway
|
||||
.clone()
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/echo/invoke",
|
||||
"application/json",
|
||||
b"{}",
|
||||
))
|
||||
.await
|
||||
.expect("unsupported media request should complete");
|
||||
assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
|
||||
assert_eq!(
|
||||
response_json(response).await["error"]["code"],
|
||||
"unsupported_media_type"
|
||||
);
|
||||
|
||||
let response = management
|
||||
.clone()
|
||||
.oneshot(empty_post("/api/v1/services/echo/0.2.0/stop"))
|
||||
.await
|
||||
.expect("active actor stop request should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let response = gateway
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/echo/invoke",
|
||||
"application/octet-stream",
|
||||
b"stopped",
|
||||
))
|
||||
.await
|
||||
.expect("stopped deployment request should complete");
|
||||
assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
|
||||
assert_eq!(
|
||||
response_json(response).await["error"]["code"],
|
||||
"service_unavailable"
|
||||
);
|
||||
|
||||
let response = management
|
||||
.oneshot(get_request("/api/v1/deployments"))
|
||||
.await
|
||||
.expect("deployment list request should complete");
|
||||
let deployments = response_json(response).await;
|
||||
assert_eq!(deployments["deployments"].as_array().unwrap().len(), 1);
|
||||
assert_eq!(deployments["deployments"][0]["active_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");
|
||||
let component = fs::read(echo_component()).expect("echo component should be readable");
|
||||
|
||||
{
|
||||
let (management, gateway) = test_apps(artifact_dir.path()).await;
|
||||
let response = management
|
||||
.clone()
|
||||
.oneshot(package_request("restored", "1.0.0", &component))
|
||||
.await
|
||||
.expect("register request should complete");
|
||||
assert_eq!(response.status(), StatusCode::CREATED);
|
||||
let response = management
|
||||
.oneshot(json_request(
|
||||
"/api/v1/deployments/restored/activate",
|
||||
json!({ "revision": "1.0.0" }),
|
||||
))
|
||||
.await
|
||||
.expect("activation request should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
|
||||
let response = gateway
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/restored/invoke",
|
||||
"application/octet-stream",
|
||||
b"before restart",
|
||||
))
|
||||
.await
|
||||
.expect("gateway invocation should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
}
|
||||
|
||||
let (management, gateway) = test_apps(artifact_dir.path()).await;
|
||||
let response = management
|
||||
.oneshot(get_request("/api/v1/deployments/restored"))
|
||||
.await
|
||||
.expect("restored deployment request should complete");
|
||||
let deployment = response_json(response).await;
|
||||
assert_eq!(deployment["active_revision"], "1.0.0");
|
||||
assert_eq!(deployment["status"], "running");
|
||||
|
||||
let response = gateway
|
||||
.oneshot(raw_request(
|
||||
"/v1/services/restored/invoke",
|
||||
"application/octet-stream",
|
||||
b"after restart",
|
||||
))
|
||||
.await
|
||||
.expect("restored gateway invocation should complete");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
assert_eq!(
|
||||
to_bytes(response.into_body(), usize::MAX).await.unwrap(),
|
||||
b"after restart".as_slice()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reloads_manifests_without_restoring_actor_memory() {
|
||||
let artifact_dir = TempDir::new().expect("temporary artifact directory");
|
||||
@@ -473,6 +666,10 @@ async fn client_fetches_transitive_wit_dependencies_from_the_registry() {
|
||||
}
|
||||
|
||||
async fn test_app(artifact_dir: &Path) -> Router {
|
||||
test_apps(artifact_dir).await.0
|
||||
}
|
||||
|
||||
async fn test_apps(artifact_dir: &Path) -> (Router, Router) {
|
||||
let console = Console::new(ConsoleConfig {
|
||||
artifact_dir: artifact_dir.to_path_buf(),
|
||||
wit_registry_dir: artifact_dir.join("wit-packages"),
|
||||
@@ -481,7 +678,8 @@ async fn test_app(artifact_dir: &Path) -> Router {
|
||||
})
|
||||
.await
|
||||
.expect("console should start");
|
||||
app(Arc::new(console), Vec::new())
|
||||
let console = Arc::new(console);
|
||||
(app(Arc::clone(&console), Vec::new()), gateway_app(console))
|
||||
}
|
||||
|
||||
fn package_request(id: &str, revision: &str, component: &[u8]) -> Request<Body> {
|
||||
@@ -552,6 +750,15 @@ fn json_request(uri: &str, value: Value) -> Request<Body> {
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn raw_request(uri: &str, content_type: &str, body: &[u8]) -> Request<Body> {
|
||||
Request::builder()
|
||||
.method("POST")
|
||||
.uri(uri)
|
||||
.header(header::CONTENT_TYPE, content_type)
|
||||
.body(Body::from(body.to_vec()))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn empty_post(uri: &str) -> Request<Body> {
|
||||
Request::builder()
|
||||
.method("POST")
|
||||
|
||||
Reference in New Issue
Block a user