use std::path::{Path, PathBuf};
use axum::extract::{Path as AxumPath, State};
use axum::response::{IntoResponse, Response};
use serde_json::{Value, json};
use crate::routes::MEMORY_SERVICE;
use crate::routes::memory_rpc;
use crate::routes::verdict::{ActionVerdict, validate_id};
use crate::server::AppState;
pub(crate) async fn compact_palace_on_socket(socket: &Path, id: &str) -> 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_compact", json!({ "palace": id }), id).await {
Ok(payload) => payload,
Err(verdict) => return verdict,
};
match payload.get("palace").and_then(Value::as_str) {
Some(compacted) if compacted == id => ActionVerdict::Succeeded {
id: id.to_string(),
detail: payload,
},
_ => ActionVerdict::Refused {
id: id.to_string(),
reason: format!(
"{MEMORY_SERVICE} answered palace_compact without confirming it compacted '{id}'"
),
detail: payload,
},
}
}
pub async fn compact_palace_handler(
State(state): State<AppState>,
AxumPath(id): AxumPath<String>,
) -> 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 = compact_palace_on_socket(&socket, &id).await;
if verdict.succeeded() {
crate::routes::deletes::refresh_metrics(
&state,
MEMORY_SERVICE,
state.memory_metrics_cache(),
)
.await;
}
verdict.into_response()
}
#[cfg(test)]
mod tests {
use super::*;
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: &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 post_through_router(uri: &str, body: Value) -> (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("POST")
.uri(uri)
.header("content-type", "application/json")
.body(Body::from(body.to_string()))
.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 parsed = serde_json::from_slice(&bytes).unwrap_or(Value::Null);
(status, parsed)
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_confirms_a_real_compaction() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(
tmp.path(),
tools_call_reply(
r#"{"palace":"scratch","total_checked":120,"orphans_removed":7,"index_size_before":120,"index_size_after":113}"#,
),
);
let verdict = compact_palace_on_socket(&socket, "scratch").await;
assert!(
matches!(&verdict, ActionVerdict::Succeeded { id, .. } if id == "scratch"),
"a confirmed compaction must read as success: {verdict:?}"
);
let response = verdict.into_response();
assert_eq!(response.status(), StatusCode::OK);
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!(true));
assert_eq!(
body["detail"]["orphans_removed"],
json!(7),
"the operator sees what was reclaimed: {body}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_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"}"#));
let verdict = compact_palace_on_socket(&socket, "scratch").await;
assert!(
matches!(&verdict, ActionVerdict::Refused { reason, .. } if reason.contains("without confirming")),
"an unconfirmed answer must read as a failure: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_rejects_a_confirmation_for_another_palace() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_memory_daemon(
tmp.path(),
tools_call_reply(r#"{"palace":"someone-else","orphans_removed":3}"#),
);
let verdict = compact_palace_on_socket(&socket, "scratch").await;
assert!(
matches!(verdict, ActionVerdict::Refused { .. }),
"a confirmation for another palace is not one for this one: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_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' is not open"}}"#,
);
let verdict = compact_palace_on_socket(&socket, "scratch").await;
assert!(
matches!(&verdict, ActionVerdict::Refused { reason, .. } if reason.contains("is not open")),
"the refusal must carry the daemon's words: {verdict:?}"
);
assert_eq!(verdict.status(), StatusCode::CONFLICT);
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_reports_a_dead_socket_as_unreachable() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let verdict = compact_palace_on_socket(&tmp.path().join("absent.sock"), "scratch").await;
assert!(
matches!(verdict, ActionVerdict::Unreachable { .. }),
"a dead socket must read as unreachable: {verdict:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn compact_refuses_a_bad_id_without_dialling() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let verdict = compact_palace_on_socket(&tmp.path().join("absent.sock"), "../x").await;
assert!(
matches!(verdict, ActionVerdict::Invalid { .. }),
"a traversal id must be refused at the console: {verdict:?}"
);
}
#[tokio::test]
async fn compact_route_rejects_a_traversal_id() {
let (status, body) =
post_through_router("/api/console/memory/palaces/..%2Fetc/compact", json!({})).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "body: {body}");
assert_eq!(body["ok"], json!(false));
}
#[tokio::test]
async fn cleanup_routes_reject_a_cross_origin_caller() {
for uri in ["/api/console/memory/palaces/scratch/compact"] {
let router = build_router(AppState::new(vec![]));
let req = Request::builder()
.method("POST")
.uri(uri)
.header("origin", "https://evil.example")
.header("content-type", "application/json")
.body(Body::from(json!({}).to_string()))
.expect("request");
let resp = router.oneshot(req).await.expect("response");
assert_eq!(
resp.status(),
StatusCode::FORBIDDEN,
"{uri} must refuse a cross-origin cleanup"
);
}
}
}