use std::path::{Path, PathBuf};
use std::time::Duration;
use axum::extract::{Path as AxumPath, Query, State};
use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use serde::Deserialize;
use serde_json::{Value, json};
use crate::server::AppState;
const DELETE_TIMEOUT: Duration = Duration::from_secs(30);
const MAX_ID_LEN: usize = 128;
const SEARCH_SERVICE_ID: &str = "trusty-search";
const MEMORY_SERVICE: &str = "trusty-memory";
#[derive(Debug)]
pub(crate) enum DeleteVerdict {
Deleted { id: String, detail: Value },
Refused {
id: String,
reason: String,
detail: Value,
},
Unreachable { id: String, reason: String },
Invalid { id: String, reason: String },
}
impl DeleteVerdict {
fn status(&self) -> StatusCode {
match self {
Self::Deleted { .. } => StatusCode::OK,
Self::Refused { .. } => StatusCode::CONFLICT,
Self::Unreachable { .. } => StatusCode::SERVICE_UNAVAILABLE,
Self::Invalid { .. } => StatusCode::BAD_REQUEST,
}
}
}
impl IntoResponse for DeleteVerdict {
fn into_response(self) -> Response {
let status = self.status();
let body = match self {
Self::Deleted { id, detail } => json!({ "ok": true, "id": id, "detail": detail }),
Self::Refused { id, reason, detail } => {
json!({ "ok": false, "id": id, "error": reason, "detail": detail })
}
Self::Unreachable { id, reason } | Self::Invalid { id, reason } => {
json!({ "ok": false, "id": id, "error": reason })
}
};
(status, axum::Json(body)).into_response()
}
}
fn validate_id(id: &str) -> Result<(), String> {
if id.is_empty() {
return Err("the resource id is empty".to_string());
}
if id.len() > MAX_ID_LEN {
return Err(format!(
"the resource id is longer than the {MAX_ID_LEN}-byte limit"
));
}
if id.contains("..") {
return Err("the resource id contains '..'".to_string());
}
if let Some(bad) = id
.chars()
.find(|c| !(c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-')))
{
return Err(format!(
"the resource id contains {bad:?}; only letters, digits, '.', '_' and '-' are accepted"
));
}
Ok(())
}
pub(crate) async fn delete_palace_on_socket(socket: &Path, id: &str, force: bool) -> DeleteVerdict {
if let Err(reason) = validate_id(id) {
return DeleteVerdict::Invalid {
id: id.to_string(),
reason,
};
}
let request = json!({
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "palace_delete",
"arguments": { "palace_id": id, "force": force },
},
});
let sent =
trusty_common::uds::send_framed_request::<_, trusty_common::uds::server::RpcResponse>(
socket,
&request,
DELETE_TIMEOUT,
)
.await;
let response = match sent {
Ok(r) => r,
Err(e) => {
return DeleteVerdict::Unreachable {
id: id.to_string(),
reason: format!("{MEMORY_SERVICE} did not answer palace_delete: {e}"),
};
}
};
if let Some(error) = response.error {
return DeleteVerdict::Refused {
id: id.to_string(),
reason: format!(
"{MEMORY_SERVICE} refused palace_delete (code {}): {}",
error.code, error.message
),
detail: json!({ "code": error.code }),
};
}
let result = response.result.unwrap_or(Value::Null);
match confirmed_palace_id(&result) {
Some(deleted) if deleted == id => DeleteVerdict::Deleted {
id: id.to_string(),
detail: json!({ "deleted": deleted }),
},
_ => DeleteVerdict::Refused {
id: id.to_string(),
reason: format!(
"{MEMORY_SERVICE} answered palace_delete without confirming the palace was deleted"
),
detail: result,
},
}
}
fn confirmed_palace_id(result: &Value) -> Option<String> {
let text = result
.get("content")?
.as_array()?
.first()?
.get("text")?
.as_str()?;
serde_json::from_str::<Value>(text)
.ok()?
.get("deleted")?
.as_str()
.map(str::to_string)
}
pub(crate) async fn delete_index_on_daemon(
client: &reqwest::Client,
base_url: &str,
id: &str,
delete_data: bool,
) -> DeleteVerdict {
if let Err(reason) = validate_id(id) {
return DeleteVerdict::Invalid {
id: id.to_string(),
reason,
};
}
let base = crate::proxy::routes::normalize_base_url(base_url);
if !crate::proxy::routes::is_local_upstream(&base) {
return DeleteVerdict::Unreachable {
id: id.to_string(),
reason: format!("the resolved {SEARCH_SERVICE_ID} address '{base}' is not loopback"),
};
}
let url = format!(
"{}/indexes/{id}?delete_data={delete_data}",
base.trim_end_matches('/')
);
let sent = client.delete(&url).timeout(DELETE_TIMEOUT).send().await;
let response = match sent {
Ok(r) => r,
Err(e) => {
return DeleteVerdict::Unreachable {
id: id.to_string(),
reason: format!("{SEARCH_SERVICE_ID} did not answer the delete: {e}"),
};
}
};
let status = response.status();
let body = response.text().await.unwrap_or_default();
let parsed: Value = serde_json::from_str(&body).unwrap_or(Value::Null);
if !status.is_success() {
return DeleteVerdict::Refused {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} refused the delete with HTTP {}: {}",
status.as_u16(),
first_line(&body)
),
detail: parsed,
};
}
if parsed.get("removed").and_then(Value::as_bool) != Some(true) {
let quiesced = parsed.get("quiesced").and_then(Value::as_bool);
return DeleteVerdict::Refused {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} skipped the delete: no registration for '{id}' was removed{}",
match quiesced {
Some(false) =>
" (in-flight writers never quiesced, so the teardown was abandoned)",
_ => "",
}
),
detail: parsed,
};
}
match parsed.get("id").and_then(Value::as_str) {
Some(echoed) if echoed != id => {
return DeleteVerdict::Refused {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} confirmed a delete for '{echoed}', not for '{id}'"
),
detail: parsed,
};
}
_ => {}
}
if delete_data && parsed.get("data_deleted").and_then(Value::as_bool) != Some(true) {
return DeleteVerdict::Refused {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} deregistered '{id}' but did not delete its on-disk data"
),
detail: parsed,
};
}
DeleteVerdict::Deleted {
id: id.to_string(),
detail: parsed,
}
}
fn first_line(body: &str) -> String {
const MAX: usize = 300;
let line = body.lines().next().unwrap_or("").trim();
if line.len() <= MAX {
return line.to_string();
}
let cut = line
.char_indices()
.map(|(i, _)| i)
.take_while(|i| *i <= MAX)
.last()
.unwrap_or(0);
format!("{}…", &line[..cut])
}
#[derive(Debug, Default, Deserialize)]
pub struct PalaceDeleteParams {
#[serde(default)]
force: bool,
}
pub async fn delete_palace_handler(
State(state): State<AppState>,
AxumPath(id): AxumPath<String>,
Query(params): Query<PalaceDeleteParams>,
) -> Response {
if let Err(reason) = validate_id(&id) {
return DeleteVerdict::Invalid { id, reason }.into_response();
}
let socket: PathBuf = match trusty_common::daemon_socket_path(MEMORY_SERVICE) {
Ok(p) => p,
Err(e) => {
return DeleteVerdict::Unreachable {
id,
reason: format!("could not resolve the {MEMORY_SERVICE} socket path: {e:#}"),
}
.into_response();
}
};
let verdict = delete_palace_on_socket(&socket, &id, params.force).await;
if matches!(verdict, DeleteVerdict::Deleted { .. }) {
refresh_metrics(&state, MEMORY_SERVICE, state.memory_metrics_cache()).await;
}
verdict.into_response()
}
#[derive(Debug, Default, Deserialize)]
pub struct IndexDeleteParams {
#[serde(default)]
delete_data: bool,
}
pub async fn delete_index_handler(
State(state): State<AppState>,
AxumPath(id): AxumPath<String>,
Query(params): Query<IndexDeleteParams>,
) -> Response {
if let Err(reason) = validate_id(&id) {
return DeleteVerdict::Invalid { id, reason }.into_response();
}
let base_url = match state.poller_cache().snapshot().await {
Some(snap) => snap.url_map().get(SEARCH_SERVICE_ID).cloned(),
None => None,
};
let Some(base_url) = base_url else {
return DeleteVerdict::Unreachable {
id,
reason: format!(
"{SEARCH_SERVICE_ID} is not reachable: the console has no live address for it"
),
}
.into_response();
};
let client = state.http_client();
let verdict = delete_index_on_daemon(&client, &base_url, &id, params.delete_data).await;
if matches!(verdict, DeleteVerdict::Deleted { .. }) {
refresh_metrics(&state, SEARCH_SERVICE_ID, state.search_metrics_cache()).await;
}
verdict.into_response()
}
async fn refresh_metrics(
state: &AppState,
service_id: &str,
cache: &crate::metrics_poller::MetricsCache,
) {
let Some(handle) = state.mcp_handles().get(service_id).cloned() else {
tracing::warn!("delete: no MCP handle registered for {service_id}; roster not refreshed");
return;
};
crate::metrics_poller::poll_once(&handle, cache).await;
}
#[cfg(test)]
mod tests {
use super::*;
use axum::body::Body;
use axum::http::Request;
use http_body_util::BodyExt as _;
use tower::ServiceExt as _;
use crate::server::build_router;
fn stub_memory_daemon(dir: &std::path::Path, reply: impl Into<String>) -> PathBuf {
let socket = dir.join("sockets").join("memory.sock");
let reply = reply.into();
let listener = trusty_common::uds::bind_hardened(&socket).expect("bind");
tokio::spawn(async move {
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
let Ok((mut conn, _)) = listener.accept().await else {
return;
};
let mut sink = Vec::new();
let _ = conn.read_to_end(&mut sink).await;
let _ = conn.write_all(reply.as_bytes()).await;
let _ = conn.write_all(b"\n").await;
let _ = conn.flush().await;
});
socket
}
fn tools_call_reply(payload: &str) -> String {
json!({
"jsonrpc": "2.0",
"id": 1,
"result": { "content": [{ "type": "text", "text": payload }] },
})
.to_string()
}
async fn stub_search_daemon(status: StatusCode, body: Value) -> String {
let app = axum::Router::new().route(
"/indexes/{id}",
axum::routing::delete(move || {
let body = body.clone();
async move { (status, axum::Json(body)) }
}),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind");
let addr = listener.local_addr().expect("addr");
tokio::spawn(async move {
let _ = axum::serve(listener, app).await;
});
format!("http://{addr}")
}
fn client() -> reqwest::Client {
(*AppState::new(vec![]).http_client()).clone()
}
fn assert_failure(verdict: &DeleteVerdict, must_mention: &str) {
assert!(
!matches!(verdict, DeleteVerdict::Deleted { .. }),
"expected a failure verdict, got {verdict:?}"
);
let text = format!("{verdict:?}");
assert!(
text.contains(must_mention),
"the failure must carry the daemon's own words ({must_mention:?}): {text}"
);
}
#[test]
fn validate_id_accepts_ordinary_ids() {
for id in ["trusty-tools", "my_palace", "index.v2", "A1", "a-b_c.d-9"] {
assert!(validate_id(id).is_ok(), "{id} must be accepted");
}
}
#[test]
fn validate_id_rejects_traversal() {
for id in ["..", "../etc", "a/../b", "....//"] {
assert!(validate_id(id).is_err(), "{id:?} must be rejected");
}
}
#[test]
fn validate_id_rejects_separators_and_control_bytes() {
for id in [
"",
"a/b",
"a\\b",
"a?b",
"a#b",
"a b",
"a\nb",
"a\0b",
"a;rm -rf /",
"a%2fb",
] {
assert!(validate_id(id).is_err(), "{id:?} must be rejected");
}
let too_long = "a".repeat(MAX_ID_LEN + 1);
assert!(
validate_id(&too_long).is_err(),
"an oversized id is refused"
);
}
#[test]
fn verdict_status_codes_separate_the_four_arms() {
let id = "x".to_string();
assert_eq!(
DeleteVerdict::Deleted {
id: id.clone(),
detail: Value::Null
}
.status(),
StatusCode::OK
);
assert_eq!(
DeleteVerdict::Refused {
id: id.clone(),
reason: String::new(),
detail: Value::Null
}
.status(),
StatusCode::CONFLICT
);
assert_eq!(
DeleteVerdict::Unreachable {
id: id.clone(),
reason: String::new()
}
.status(),
StatusCode::SERVICE_UNAVAILABLE
);
assert_eq!(
DeleteVerdict::Invalid {
id,
reason: String::new()
}
.status(),
StatusCode::BAD_REQUEST
);
}
#[tokio::test(flavor = "multi_thread")]
async fn palace_delete_confirms_a_real_delete() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(tmp.path(), tools_call_reply(r#"{"deleted":"scratch"}"#));
let verdict = delete_palace_on_socket(&socket, "scratch", true).await;
assert!(
matches!(&verdict, DeleteVerdict::Deleted { id, .. } if id == "scratch"),
"a confirmed delete must read as success: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn palace_delete_reports_a_daemon_refusal_as_a_failure() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(
tmp.path(),
r#"{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"Palace 'scratch' still has 4 drawers; pass force=true"}}"#,
);
let verdict = delete_palace_on_socket(&socket, "scratch", false).await;
assert_failure(&verdict, "still has 4 drawers");
}
#[tokio::test(flavor = "multi_thread")]
async fn palace_delete_reports_an_unconfirmed_answer_as_a_failure() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(
tmp.path(),
tools_call_reply(r#"{"status":"noop","reason":"not loaded"}"#),
);
let verdict = delete_palace_on_socket(&socket, "scratch", true).await;
assert_failure(&verdict, "without confirming");
}
#[tokio::test(flavor = "multi_thread")]
async fn palace_delete_rejects_a_confirmation_for_another_id() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(
tmp.path(),
tools_call_reply(r#"{"deleted":"someone-else"}"#),
);
let verdict = delete_palace_on_socket(&socket, "scratch", true).await;
assert_failure(&verdict, "without confirming");
}
#[tokio::test(flavor = "multi_thread")]
async fn palace_delete_reports_a_dead_socket_as_unreachable() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let verdict =
delete_palace_on_socket(&tmp.path().join("absent.sock"), "scratch", false).await;
assert!(
matches!(verdict, DeleteVerdict::Unreachable { .. }),
"a dead socket must read as unreachable: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn palace_delete_refuses_a_bad_id_without_dialling() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let verdict = delete_palace_on_socket(&tmp.path().join("absent.sock"), "../x", false).await;
assert!(
matches!(verdict, DeleteVerdict::Invalid { .. }),
"a traversal id must be refused at the console: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_confirms_a_real_delete() {
let base = stub_search_daemon(
StatusCode::OK,
json!({"id":"scratch","removed":true,"data_deleted":false,"quiesced":true}),
)
.await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", false).await;
assert!(
matches!(&verdict, DeleteVerdict::Deleted { id, .. } if id == "scratch"),
"a confirmed delete must read as success: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_a_skipped_delete_as_a_failure() {
let base = stub_search_daemon(
StatusCode::OK,
json!({"id":"scratch","removed":false,"data_deleted":false,"quiesced":true}),
)
.await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", false).await;
assert_failure(&verdict, "skipped the delete");
let response = verdict.into_response();
assert_eq!(response.status(), StatusCode::CONFLICT);
let bytes = response
.into_body()
.collect()
.await
.expect("body")
.to_bytes();
let body: Value = serde_json::from_slice(&bytes).expect("json");
assert_eq!(
body["ok"],
json!(false),
"a no-op must never render ok:true"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_an_empty_body_as_a_failure() {
let base = stub_search_daemon(StatusCode::OK, json!({})).await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", false).await;
assert_failure(&verdict, "skipped the delete");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_rejects_a_confirmation_for_another_id() {
let base = stub_search_daemon(
StatusCode::OK,
json!({"id":"someone-else","removed":true,"data_deleted":false,"quiesced":true}),
)
.await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", false).await;
assert_failure(&verdict, "not for 'scratch'");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_names_an_abandoned_teardown() {
let base = stub_search_daemon(
StatusCode::OK,
json!({"id":"scratch","removed":false,"data_deleted":false,"quiesced":false}),
)
.await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", false).await;
assert_failure(&verdict, "never quiesced");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_a_daemon_error_status_as_a_failure() {
let base = stub_search_daemon(
StatusCode::BAD_REQUEST,
json!({"error":"delete_data must be a boolean"}),
)
.await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", false).await;
assert_failure(&verdict, "delete_data must be a boolean");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_undeleted_data_as_a_failure() {
let base = stub_search_daemon(
StatusCode::OK,
json!({"id":"scratch","removed":true,"data_deleted":false,"quiesced":true}),
)
.await;
let verdict = delete_index_on_daemon(&client(), &base, "scratch", true).await;
assert_failure(&verdict, "did not delete its on-disk data");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_a_dead_daemon_as_unreachable() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind");
let addr = listener.local_addr().expect("addr");
drop(listener);
let verdict =
delete_index_on_daemon(&client(), &format!("http://{addr}"), "scratch", false).await;
assert!(
matches!(verdict, DeleteVerdict::Unreachable { .. }),
"a dead daemon must read as unreachable: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_does_not_follow_a_redirect() {
let app = axum::Router::new().route(
"/indexes/{id}",
axum::routing::delete(|| async {
(
StatusCode::TEMPORARY_REDIRECT,
[(
axum::http::header::LOCATION,
"http://169.254.169.254/indexes/scratch",
)],
)
}),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind");
let addr = listener.local_addr().expect("addr");
tokio::spawn(async move {
let _ = axum::serve(listener, app).await;
});
let verdict =
delete_index_on_daemon(&client(), &format!("http://{addr}"), "scratch", false).await;
assert_failure(&verdict, "HTTP 307");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_refuses_a_non_loopback_upstream() {
let verdict =
delete_index_on_daemon(&client(), "http://10.0.0.9:7878", "scratch", false).await;
assert!(
matches!(&verdict, DeleteVerdict::Unreachable { reason, .. } if reason.contains("not loopback")),
"a remote upstream must be refused: {verdict:?}"
);
}
async fn delete_through_router(uri: &str) -> (StatusCode, Value) {
let router = build_router(AppState::new(vec![]));
let req = Request::builder()
.method("DELETE")
.uri(uri)
.body(Body::empty())
.expect("request");
let resp = router.oneshot(req).await.expect("response");
let status = resp.status();
let bytes = resp.into_body().collect().await.expect("body").to_bytes();
let body = serde_json::from_slice(&bytes).unwrap_or(Value::Null);
(status, body)
}
#[tokio::test]
async fn palace_route_rejects_a_traversal_id() {
let (status, body) = delete_through_router("/api/console/memory/palaces/..%2Fetc").await;
assert_eq!(status, StatusCode::BAD_REQUEST, "body: {body}");
assert_eq!(body["ok"], json!(false));
}
#[tokio::test]
async fn index_route_reports_an_unresolved_daemon_as_unreachable() {
let (status, body) = delete_through_router("/api/console/search/indexes/scratch").await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "body: {body}");
assert_eq!(body["ok"], json!(false));
assert!(
body["error"]
.as_str()
.unwrap_or_default()
.contains("trusty-search"),
"the error must name the daemon: {body}"
);
}
#[tokio::test]
async fn index_route_rejects_a_bad_id() {
let (status, body) = delete_through_router("/api/console/search/indexes/a%2Fb").await;
assert_eq!(
status,
StatusCode::BAD_REQUEST,
"validation must run before daemon resolution, so this cannot be 503: {body}"
);
assert_eq!(body["ok"], json!(false));
assert!(
body["error"].as_str().unwrap_or_default().contains('/'),
"the error must name the offending character: {body}"
);
}
#[tokio::test]
async fn delete_routes_reject_a_cross_origin_caller() {
for uri in [
"/api/console/memory/palaces/scratch",
"/api/console/search/indexes/scratch",
] {
let router = build_router(AppState::new(vec![]));
let req = Request::builder()
.method("DELETE")
.uri(uri)
.header("origin", "https://evil.example")
.body(Body::empty())
.expect("request");
let resp = router.oneshot(req).await.expect("response");
assert_eq!(
resp.status(),
StatusCode::FORBIDDEN,
"{uri} must refuse a cross-origin delete"
);
}
}
}