use std::path::{Path, PathBuf};
use axum::extract::{Path as AxumPath, Query, State};
use axum::response::{IntoResponse, Response};
use serde::Deserialize;
use serde_json::{Value, json};
use crate::routes::memory_rpc;
use crate::routes::verdict::{ActionVerdict, first_line, validate_id};
use crate::routes::{ACTION_TIMEOUT, MEMORY_SERVICE, SEARCH_SERVICE_ID};
use crate::server::AppState;
pub(crate) async fn delete_palace_on_socket(socket: &Path, id: &str, force: bool) -> ActionVerdict {
if let Err(reason) = validate_id(id) {
return ActionVerdict::Invalid {
id: id.to_string(),
reason,
};
}
let payload = match memory_rpc::call_tool(
socket,
"palace_delete",
json!({ "palace_id": id, "force": force }),
id,
)
.await
{
Ok(payload) => payload,
Err(verdict) => return verdict,
};
match payload.get("deleted").and_then(Value::as_str) {
Some(deleted) if deleted == id => ActionVerdict::Succeeded {
id: id.to_string(),
detail: json!({ "deleted": deleted }),
},
_ => ActionVerdict::Refused {
id: id.to_string(),
reason: format!(
"{MEMORY_SERVICE} answered palace_delete without confirming the palace was deleted"
),
detail: payload,
},
}
}
pub(crate) async fn delete_index_on_socket(
socket: &Path,
id: &str,
delete_data: bool,
expected_root_path: Option<&str>,
) -> ActionVerdict {
if let Err(reason) = validate_id(id) {
return ActionVerdict::Invalid {
id: id.to_string(),
reason,
};
}
let mut params = json!({ "index_id": id, "delete_data": delete_data });
if let Some(root) = expected_root_path {
params["expected_root_path"] = Value::String(root.to_string());
}
let parsed = match crate::search_uds::call(
socket,
crate::search_uds::METHOD_INDEX_DELETE,
params,
ACTION_TIMEOUT,
)
.await
{
Ok(result) => result,
Err(crate::search_uds::SearchRpcError::Refused { code, message }) => {
return ActionVerdict::Refused {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} refused the delete (code {code}): {}",
first_line(&message)
),
detail: json!({ "code": code }),
};
}
Err(e) => {
return ActionVerdict::Unreachable {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} did not answer the delete: {}",
e.message()
),
};
}
};
if parsed.get("removed").and_then(Value::as_bool) != Some(true) {
let quiesced = parsed.get("quiesced").and_then(Value::as_bool);
return ActionVerdict::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 ActionVerdict::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 ActionVerdict::Refused {
id: id.to_string(),
reason: format!(
"{SEARCH_SERVICE_ID} deregistered '{id}' but did not delete its on-disk data"
),
detail: parsed,
};
}
ActionVerdict::Succeeded {
id: id.to_string(),
detail: parsed,
}
}
#[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 ActionVerdict::Invalid { id, reason }.into_response();
}
let socket: PathBuf = match trusty_common::daemon_socket_path(MEMORY_SERVICE) {
Ok(p) => p,
Err(e) => {
return ActionVerdict::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, ActionVerdict::Succeeded { .. }) {
refresh_metrics(&state, MEMORY_SERVICE, state.memory_metrics_cache()).await;
}
verdict.into_response()
}
pub(crate) fn purge_data_by_default() -> bool {
true
}
#[derive(Debug, Deserialize)]
pub struct IndexDeleteParams {
#[serde(default = "purge_data_by_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 ActionVerdict::Invalid { id, reason }.into_response();
}
let socket: PathBuf = match state.search_socket_path() {
Ok(p) => p,
Err(reason) => return ActionVerdict::Unreachable { id, reason }.into_response(),
};
let verdict = delete_index_on_socket(&socket, &id, params.delete_data, None).await;
if matches!(verdict, ActionVerdict::Succeeded { .. }) {
refresh_metrics(&state, SEARCH_SERVICE_ID, state.search_metrics_cache()).await;
}
verdict.into_response()
}
pub(crate) 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)]
pub(crate) mod tests {
use super::*;
use crate::routes::verdict::MAX_ID_LEN;
use axum::body::Body;
use axum::http::Request;
use axum::http::StatusCode;
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()
}
pub(crate) fn stub_search_socket<F>(dir: &std::path::Path, respond: F) -> PathBuf
where
F: Fn(&Value) -> Value + Send + Sync + 'static,
{
let socket = dir.join("sockets").join("search.sock");
let listener = trusty_common::uds::bind_hardened(&socket).expect("bind");
let respond = std::sync::Arc::new(respond);
tokio::spawn(async move {
loop {
let Ok((mut conn, _)) = listener.accept().await else {
return;
};
let respond = std::sync::Arc::clone(&respond);
tokio::spawn(async move {
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
let mut raw = Vec::new();
let _ = conn.read_to_end(&mut raw).await;
let request: Value = serde_json::from_slice(&raw).unwrap_or(Value::Null);
let reply = respond(&request).to_string();
let _ = conn.write_all(reply.as_bytes()).await;
let _ = conn.write_all(b"\n").await;
let _ = conn.flush().await;
});
}
});
socket
}
pub(crate) fn always_result(result: Value) -> impl Fn(&Value) -> Value + Send + Sync + 'static {
move |_| json!({ "jsonrpc": "2.0", "id": 1, "result": result.clone() })
}
fn always_error(
code: i64,
message: &'static str,
) -> impl Fn(&Value) -> Value + Send + Sync + 'static {
move |_| json!({ "jsonrpc": "2.0", "id": 1, "error": { "code": code, "message": message } })
}
fn assert_failure(verdict: &ActionVerdict, must_mention: &str) {
assert!(
!matches!(verdict, ActionVerdict::Succeeded { .. }),
"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!(
ActionVerdict::Succeeded {
id: id.clone(),
detail: Value::Null
}
.status(),
StatusCode::OK
);
assert_eq!(
ActionVerdict::Refused {
id: id.clone(),
reason: String::new(),
detail: Value::Null
}
.status(),
StatusCode::CONFLICT
);
assert_eq!(
ActionVerdict::Unreachable {
id: id.clone(),
reason: String::new()
}
.status(),
StatusCode::SERVICE_UNAVAILABLE
);
assert_eq!(
ActionVerdict::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, ActionVerdict::Succeeded { 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, ActionVerdict::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, ActionVerdict::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 tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(
tmp.path(),
always_result(
json!({"id":"scratch","removed":true,"data_deleted":false,"quiesced":true}),
),
);
let verdict = delete_index_on_socket(&socket, "scratch", false, None).await;
assert!(
matches!(&verdict, ActionVerdict::Succeeded { id, .. } if id == "scratch"),
"a confirmed delete must read as success: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_sends_the_id_and_the_data_choice_as_params() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(tmp.path(), |request: &Value| {
let params = &request["params"];
let ok = params["index_id"] == json!("scratch") && params["delete_data"] == json!(true);
json!({
"jsonrpc": "2.0",
"id": 1,
"result": {
"id": "scratch",
"removed": ok,
"data_deleted": ok,
"quiesced": true,
"method": request["method"].clone(),
},
})
});
let verdict = delete_index_on_socket(&socket, "scratch", true, None).await;
assert!(
matches!(&verdict, ActionVerdict::Succeeded { detail, .. }
if detail["method"] == json!("search.index.delete")),
"the params and the method name must both reach the daemon: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_a_skipped_delete_as_a_failure() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(
tmp.path(),
always_result(
json!({"id":"scratch","removed":false,"data_deleted":false,"quiesced":true}),
),
);
let verdict = delete_index_on_socket(&socket, "scratch", false, None).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 tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(tmp.path(), always_result(json!({})));
let verdict = delete_index_on_socket(&socket, "scratch", false, None).await;
assert_failure(&verdict, "skipped the delete");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_rejects_a_confirmation_for_another_id() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(
tmp.path(),
always_result(
json!({"id":"someone-else","removed":true,"data_deleted":false,"quiesced":true}),
),
);
let verdict = delete_index_on_socket(&socket, "scratch", false, None).await;
assert_failure(&verdict, "not for 'scratch'");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_names_an_abandoned_teardown() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(
tmp.path(),
always_result(
json!({"id":"scratch","removed":false,"data_deleted":false,"quiesced":false}),
),
);
let verdict = delete_index_on_socket(&socket, "scratch", false, None).await;
assert_failure(&verdict, "never quiesced");
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_reports_a_daemon_error_frame_as_a_failure() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(
tmp.path(),
always_error(-32602, "delete_data must be a boolean"),
);
let verdict = delete_index_on_socket(&socket, "scratch", false, None).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 tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(
tmp.path(),
always_result(
json!({"id":"scratch","removed":true,"data_deleted":false,"quiesced":true}),
),
);
let verdict = delete_index_on_socket(&socket, "scratch", true, None).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 tmp = tempfile::TempDir::new().expect("tempdir");
let verdict =
delete_index_on_socket(&tmp.path().join("absent.sock"), "scratch", false, None).await;
assert!(
matches!(verdict, ActionVerdict::Unreachable { .. }),
"a dead daemon must read as unreachable: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_forwards_the_expected_root_path() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_search_socket(tmp.path(), |request: &Value| {
json!({
"jsonrpc": "2.0",
"id": 1,
"result": {
"id": "scratch",
"removed": true,
"data_deleted": false,
"quiesced": true,
"sent": request["params"].clone(),
},
})
});
let pinned = delete_index_on_socket(&socket, "scratch", false, Some("/gone/scratch")).await;
assert!(
matches!(&pinned, ActionVerdict::Succeeded { detail, .. }
if detail["sent"]["expected_root_path"] == json!("/gone/scratch")),
"the expectation must reach the daemon verbatim: {pinned:?}"
);
let unpinned = delete_index_on_socket(&socket, "scratch", false, None).await;
assert!(
matches!(&unpinned, ActionVerdict::Succeeded { detail, .. }
if detail["sent"].get("expected_root_path").is_none()),
"a caller with no expectation must send no field: {unpinned:?}"
);
}
async fn delete_through_router(uri: &str) -> (StatusCode, Value) {
let tmp = tempfile::TempDir::new().expect("tempdir");
let router =
build_router(AppState::new(vec![]).with_search_socket(tmp.path().join("absent.sock")));
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_a_dead_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}"
);
}
async fn params_seen_by_the_daemon(uri: &str) -> Value {
let tmp = tempfile::TempDir::new().expect("tempdir");
let seen = std::sync::Arc::new(std::sync::Mutex::new(Value::Null));
let recorder = std::sync::Arc::clone(&seen);
let socket = stub_search_socket(tmp.path(), move |request: &Value| {
if let Ok(mut slot) = recorder.lock() {
*slot = request["params"].clone();
}
json!({
"jsonrpc": "2.0",
"id": 1,
"result": { "id": request["params"]["index_id"].clone(), "removed": true,
"data_deleted": true, "quiesced": true },
})
});
let router = build_router(AppState::new(vec![]).with_search_socket(socket));
let req = Request::builder()
.method("DELETE")
.uri(uri)
.body(Body::empty())
.expect("request");
let resp = router.oneshot(req).await.expect("response");
assert_eq!(
resp.status(),
StatusCode::OK,
"the stub confirms the delete"
);
seen.lock().expect("lock").clone()
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_route_purges_by_default() {
let params = params_seen_by_the_daemon("/api/console/search/indexes/scratch").await;
assert_eq!(
params["delete_data"],
json!(true),
"an unqualified console delete must take the on-disk data too: {params}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn index_delete_route_honours_the_deregister_only_opt_out() {
let params =
params_seen_by_the_daemon("/api/console/search/indexes/scratch?delete_data=false")
.await;
assert_eq!(
params["delete_data"],
json!(false),
"the opt-out must reach the daemon: {params}"
);
}
#[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"
);
}
}
}