use std::path::{Path, PathBuf};
use std::time::Duration;
use serde_json::{json, Value};
use tokio::sync::oneshot;
use trusty_common::uds::server::{RpcResponse, CODE_METHOD_NOT_FOUND};
use trusty_common::uds::{send_framed_request_capped, send_framed_stream_request_capped};
use super::{
build_router, dispatcher_method_count, serve_with_shutdown, socket_path, FOLDED_METHODS,
MAX_FRAME_BYTES, METHOD_HEALTH, STREAM_METHODS,
};
use crate::AppState;
const CALL_TIMEOUT: Duration = Duration::from_secs(20);
fn test_state() -> AppState {
trusty_common::memory_core::retrieval::seed_shared_embedder_with_mock();
let tmp = tempfile::tempdir().expect("tempdir");
let root = tmp.path().to_path_buf();
std::mem::forget(tmp);
unsafe {
std::env::set_var("TRUSTY_SKIP_PALACE_ENFORCEMENT", "1");
}
let state = AppState::new(root);
state.set_ready();
state
}
fn frame(id: i64, method: &str, params: Value) -> Value {
json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
}
struct Daemon {
socket: PathBuf,
stop: Option<oneshot::Sender<()>>,
joined: Option<tokio::task::JoinHandle<()>>,
}
impl Daemon {
async fn start(state: AppState) -> Self {
let tmp = tempfile::tempdir().expect("tempdir");
let socket = tmp.path().join("sockets").join("trusty-memory.sock");
std::mem::forget(tmp);
let (stop, shutdown) = oneshot::channel::<()>();
let serve_socket = socket.clone();
let joined = tokio::spawn(async move {
let _ = serve_with_shutdown(state, &serve_socket, async {
let _ = shutdown.await;
})
.await;
});
for _ in 0..200 {
if trusty_common::uds::socket_is_serving(&socket, Duration::from_millis(200)).await {
break;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
Self {
socket,
stop: Some(stop),
joined: Some(joined),
}
}
fn socket(&self) -> &Path {
&self.socket
}
async fn call(&self, method: &str, params: Value) -> RpcResponse {
send_framed_request_capped(
&self.socket,
&frame(1, method, params),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
.unwrap_or_else(|e| panic!("call {method}: {e}"))
}
async fn ok(&self, method: &str, params: Value) -> Value {
let response = self.call(method, params).await;
match (response.result, response.error) {
(Some(result), _) => result,
(None, Some(e)) => panic!("{method} failed: {} ({})", e.message, e.code),
(None, None) => panic!("{method} answered neither result nor error"),
}
}
async fn shutdown(mut self) {
if let Some(stop) = self.stop.take() {
let _ = stop.send(());
}
if let Some(joined) = self.joined.take() {
let _ = tokio::time::timeout(Duration::from_secs(10), joined).await;
}
}
}
#[tokio::test]
async fn rpc_router_registers_every_documented_method() {
let router = build_router(test_state());
let registered: Vec<&str> = router.method_names().collect();
let mut documented = FOLDED_METHODS.to_vec();
documented.sort_unstable();
assert_eq!(
registered, documented,
"FOLDED_METHODS must equal what build_router registers"
);
let streams: Vec<&str> = router.stream_names().collect();
let mut documented_streams = STREAM_METHODS.to_vec();
documented_streams.sort_unstable();
assert_eq!(
streams, documented_streams,
"STREAM_METHODS must equal what build_router registers"
);
}
#[tokio::test]
async fn rpc_folded_names_do_not_collide_with_dispatcher_names() {
let dispatcher = crate::transport::rpc::method_names();
for folded in FOLDED_METHODS.iter().chain(STREAM_METHODS) {
assert!(
!dispatcher.contains(folded),
"{folded} is registered AND routed by the dispatcher; the folded \
registration would silently shadow it"
);
}
}
#[test]
fn rpc_reports_the_dispatcher_surface_size() {
assert!(
dispatcher_method_count() >= 30,
"the dispatcher routes the whole tool surface; got {}",
dispatcher_method_count()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_dispatcher_method_answers_through_the_fallback() {
let daemon = Daemon::start(test_state()).await;
let result = daemon.ok("palace_list", json!({})).await;
assert!(
result["palaces"]
.as_array()
.expect("palace_list returns {palaces: []}")
.is_empty(),
"a fresh state lists zero palaces, got {result}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_reports_method_not_found_for_an_unknown_method() {
let daemon = Daemon::start(test_state()).await;
let response = daemon.call("definitely_not_a_method", json!({})).await;
let error = response.error.expect("an unknown method must be refused");
assert_eq!(error.code, CODE_METHOD_NOT_FOUND);
assert!(
error.message.contains("definitely_not_a_method"),
"the refusal must name what was asked for: {}",
error.message
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_status_answers_with_no_params() {
let daemon = Daemon::start(test_state()).await;
let bare = send_framed_request_capped::<_, RpcResponse>(
daemon.socket(),
&json!({"jsonrpc": "2.0", "id": 1, "method": "memory.status"}),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
.expect("status with absent params");
let result = bare.result.expect("status must answer");
assert!(result["version"].is_string());
assert_eq!(result["palace_count"], 0);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_health_answers_over_a_real_socket() {
let daemon = Daemon::start(test_state()).await;
let result = daemon.ok(METHOD_HEALTH, json!({})).await;
assert_eq!(result["status"], "ok");
assert!(result["version"].is_string());
assert!(
result.get("addr").is_none(),
"the TCP address field retired with the listener: {result}"
);
assert!(
result["socket"]
.as_str()
.is_some_and(|s| s.ends_with(".sock")),
"health must name the socket it serves: {result}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_config_reports_whether_a_key_is_set_without_echoing_it() {
let daemon = Daemon::start(test_state()).await;
let result = daemon.ok("memory.config", json!({})).await;
assert!(result["openrouter_configured"].is_boolean());
assert!(result["model"].is_string());
assert!(
result.get("openrouter_api_key").is_none(),
"the key must never cross the wire: {result}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_drawer_create_list_and_delete_round_trip() {
let state = test_state();
let palace = seed_palace(&state, "alpha");
let daemon = Daemon::start(state).await;
let created = daemon
.ok(
"memory.drawer_create",
json!({"palace_id": palace, "content": "a folded drawer for the round trip"}),
)
.await;
let drawer_id = created["id"]
.as_str()
.expect("create returns an id")
.to_string();
let listed = daemon
.ok("memory.drawers_list", json!({"palace_id": palace}))
.await;
assert!(
listed.to_string().contains(&drawer_id),
"the created drawer must be listed: {listed}"
);
let deleted = daemon
.ok(
"memory.drawer_delete",
json!({"palace_id": palace, "drawer_id": drawer_id}),
)
.await;
assert_eq!(
deleted["deleted"], true,
"delete answers a body, not the former 204: {deleted}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_drawer_create_attributes_the_caller_it_was_given() {
let state = test_state();
let palace = seed_palace(&state, "attributed");
let daemon = Daemon::start(state).await;
daemon
.ok(
"memory.drawer_create",
json!({
"palace_id": palace,
"content": "attributed to a caller that named itself",
"client": "trusty-console",
"workstream": "feat-6286",
}),
)
.await;
let listed = daemon
.ok("memory.drawers_list", json!({"palace_id": palace}))
.await
.to_string();
assert!(
listed.contains("feat-6286"),
"the caller's workstream must reach the drawer's tags: {listed}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_drawer_delete_reports_not_found_for_an_unknown_id() {
let state = test_state();
let palace = seed_palace(&state, "missing-drawer");
let daemon = Daemon::start(state).await;
let response = daemon
.call(
"memory.drawer_delete",
json!({"palace_id": palace, "drawer_id": uuid::Uuid::new_v4().to_string()}),
)
.await;
let error = response.error.expect("an absent drawer must be refused");
assert_eq!(error.code, crate::transport::CODE_NOT_FOUND);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_palace_get_reports_not_found_for_an_unknown_id() {
let daemon = Daemon::start(test_state()).await;
let response = daemon
.call("memory.palace_get", json!({"palace_id": "no-such-palace"}))
.await;
let error = response.error.expect("an absent palace must be refused");
assert_eq!(error.code, crate::transport::CODE_NOT_FOUND);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_palaces_list_reports_counts_per_palace() {
let state = test_state();
let written = seed_palace(&state, "roster-one");
let empty = seed_palace(&state, "roster-two");
let daemon = Daemon::start(state).await;
daemon
.ok(
"memory.drawer_create",
json!({
"palace_id": written,
"content": "a drawer the roster has to count",
"force": true,
}),
)
.await;
let result = daemon.ok("memory.palaces_list", json!({})).await;
let rows = result["palaces"]
.as_array()
.expect("palaces_list answers a palaces array");
assert_eq!(rows.len(), 2, "one row per palace: {result}");
let row = rows
.iter()
.find(|r| r["id"] == written.as_str())
.expect("the written palace is listed");
assert!(row["error"].is_null(), "a readable palace carries no error");
assert_eq!(
row["palace"]["drawer_count"], 1,
"the count must be a measurement, not a peeked zero: {row}"
);
let row = rows
.iter()
.find(|r| r["id"] == empty.as_str())
.expect("the empty palace is listed");
assert!(row["error"].is_null());
assert_eq!(row["palace"]["drawer_count"], 0);
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_palaces_list_reports_an_unreadable_palace_rather_than_dropping_it() {
let state = test_state();
let good = seed_palace(&state, "roster-good");
let broken = seed_palace_off_registry(&state, "roster-broken");
let dir = state.data_root.join(&broken);
std::fs::create_dir_all(dir.join("identity.txt")).expect("wedge the palace identity");
let daemon = Daemon::start(state).await;
let result = daemon.ok("memory.palaces_list", json!({})).await;
let rows = result["palaces"]
.as_array()
.expect("palaces_list answers a palaces array");
assert_eq!(
rows.len(),
2,
"an unreadable palace must still be a row: {result}"
);
let row = rows
.iter()
.find(|r| r["id"] == broken.as_str())
.expect("the wedged palace is listed");
assert!(
row["error"].is_string(),
"the row must carry why it could not be read: {row}"
);
assert!(
row["palace"].is_null(),
"a failed read must not become a row of zeros: {row}"
);
let row = rows
.iter()
.find(|r| r["id"] == good.as_str())
.expect("the healthy palace is unaffected");
assert!(row["error"].is_null());
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_kg_all_clamps_an_absurd_limit() {
let state = test_state();
let palace = seed_palace(&state, "kg-clamp");
let daemon = Daemon::start(state).await;
let result = daemon
.ok(
"memory.kg_all",
json!({"palace_id": palace, "limit": 1_000_000_000_u64}),
)
.await;
assert!(result.is_array(), "kg_all answers an array: {result}");
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_kg_graph_neighbors_refuses_a_bad_direction() {
let state = test_state();
let palace = seed_palace(&state, "kg-direction");
let daemon = Daemon::start(state).await;
let response = daemon
.call(
"memory.kg_graph_neighbors",
json!({"palace_id": palace, "node": "n", "direction": "sideways"}),
)
.await;
let error = response.error.expect("a bad direction must be refused");
assert_eq!(error.code, trusty_common::uds::server::CODE_INVALID_PARAMS);
assert!(
error.message.contains("sideways"),
"the refusal must name what was sent: {}",
error.message
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_remember_async_rejects_short_content() {
let daemon = Daemon::start(test_state()).await;
let response = daemon
.call("memory.remember_async", json!({"content": "too short"}))
.await;
let error = response.error.expect("short content must be refused");
assert_eq!(error.code, trusty_common::uds::server::CODE_INVALID_PARAMS);
assert!(
error.message.contains("too short"),
"the refusal must say why: {}",
error.message
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_remember_async_queues_and_persists() {
let state = test_state();
let palace = seed_palace(&state, "queued");
let daemon = Daemon::start(state).await;
let queued = daemon
.ok(
"memory.remember_async",
json!({
"palace": palace,
"content": "the fire and forget path really does persist this",
}),
)
.await;
assert_eq!(queued["status"], "queued");
let mut listed = String::new();
for _ in 0..100 {
tokio::time::sleep(Duration::from_millis(100)).await;
listed = daemon
.ok("memory.drawers_list", json!({"palace_id": palace}))
.await
.to_string();
if listed.contains("really does persist") {
break;
}
}
assert!(
listed.contains("really does persist"),
"the queued write must reach the palace: {listed}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_logs_tail_answers_a_bounded_page() {
let daemon = Daemon::start(test_state()).await;
let result = daemon.ok("memory.logs_tail", json!({"n": 5})).await;
assert!(
result["lines"].is_array(),
"logs_tail returns lines: {result}"
);
assert!(result["total"].is_number());
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_activity_refuses_an_unparseable_since() {
let daemon = Daemon::start(test_state()).await;
let ok = daemon.ok("memory.activity", json!({})).await;
assert!(ok["entries"].is_array(), "activity returns entries: {ok}");
let response = daemon
.call("memory.activity", json!({"since": "not-a-timestamp"}))
.await;
let error = response.error.expect("a bad timestamp must be refused");
assert_eq!(error.code, trusty_common::uds::server::CODE_INVALID_PARAMS);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_dream_status_answers_an_aggregate() {
let daemon = Daemon::start(test_state()).await;
let result = daemon.ok("memory.dream_status", json!({})).await;
assert!(
result.is_object(),
"dream_status returns an object: {result}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_chat_providers_answers_both_upstreams() {
let daemon = Daemon::start(test_state()).await;
let result = daemon.ok("memory.chat_providers", json!({})).await;
let names: Vec<&str> = result["providers"]
.as_array()
.expect("providers is an array")
.iter()
.filter_map(|p| p["name"].as_str())
.collect();
assert_eq!(names, vec!["ollama", "openrouter"]);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_messages_send_list_and_mark_read_round_trip() {
let state = test_state();
let inbox = seed_palace(&state, "inbox");
let daemon = Daemon::start(state).await;
let sent = daemon
.ok(
"memory.message_send",
json!({
"to_palace": inbox,
"purpose": "handoff",
"content": "the folded message path still delivers",
"from_palace": "sender",
}),
)
.await;
assert_eq!(sent["status"], "sent");
let drawer_id = sent["drawer_id"].as_str().expect("a drawer id").to_string();
let listed = daemon
.ok(
"memory.messages_list",
json!({"palace": inbox, "unread_only": true}),
)
.await;
assert_eq!(
listed.as_array().expect("an array").len(),
1,
"the sent message must be unread in the inbox: {listed}"
);
let first = daemon
.ok(
"memory.message_mark_read",
json!({"palace": inbox, "drawer_id": drawer_id.clone()}),
)
.await;
assert_eq!(first["flipped"], true, "the first ack flips the flag");
let second = daemon
.ok(
"memory.message_mark_read",
json!({"palace": inbox, "drawer_id": drawer_id}),
)
.await;
assert_eq!(
second["flipped"], false,
"a second ack is a success that flipped nothing, not an error"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn mcp_writes_carry_distinct_ws_tags_per_caller_over_rpc() {
let state = test_state();
let palace = seed_palace(&state, "cred-multi-caller");
let daemon = Daemon::start(state).await;
for (workstream, marker) in [("ws-alpha", "alpha"), ("ws-beta", "beta")] {
let answered = daemon
.ok(
"tools/call",
json!({
"name": "memory_remember",
"arguments": {
"palace": palace,
"text": format!(
"distinct-caller regression content marker {marker} with enough tokens"
),
"room": "General",
"workstream": workstream,
"force": true,
}
}),
)
.await;
let text = answered["content"][0]["text"].as_str().unwrap_or_default();
let inner: Value = serde_json::from_str(text).unwrap_or_default();
assert_eq!(
inner["status"], "stored",
"expected a stored (not skipped) envelope for {workstream}; got {inner:?}"
);
}
let listed = daemon
.ok(
"memory.drawers_list",
json!({ "palace_id": palace, "limit": 10 }),
)
.await;
let drawers = listed.as_array().expect("drawers array");
assert_eq!(drawers.len(), 2, "expected two drawers, got {drawers:?}");
for drawer in drawers {
let content = drawer["content"].as_str().unwrap_or_default();
let tags: Vec<&str> = drawer["tags"]
.as_array()
.expect("tags array")
.iter()
.filter_map(|t| t.as_str())
.collect();
let (own, other) = if content.contains("marker alpha") {
("ws-alpha", "ws-beta")
} else if content.contains("marker beta") {
("ws-beta", "ws-alpha")
} else {
panic!("drawer content did not match either marker: {content:?}");
};
assert!(
tags.contains(&format!("ws:{own}").as_str())
&& tags.contains(&format!("creator:workstream={own}").as_str()),
"the {own} drawer must carry its own ws tags; got {tags:?}"
);
assert!(
!tags.contains(&format!("ws:{other}").as_str())
&& !tags.contains(&format!("creator:workstream={other}").as_str()),
"the {own} drawer must NOT leak the {other} caller's ws tags \
(daemon-vs-caller mis-attribution); got {tags:?}"
);
}
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_chat_refuses_a_unary_call_naming_the_stream_requirement() {
let daemon = Daemon::start(test_state()).await;
let response = daemon
.call("memory.chat", json!({"message": "hello"}))
.await;
let error = response
.error
.expect("a unary call to a streaming method must be refused");
assert_eq!(
error.code,
trusty_common::uds::server::CODE_STREAM_REQUIRED,
"the refusal must be the stream-required code, got {}",
error.message
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_chat_reports_a_provider_failure_as_the_terminal_error_frame() {
let daemon = Daemon::start(test_state()).await;
let mut stream = send_framed_stream_request_capped::<_, Value>(
daemon.socket(),
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "memory.chat",
"stream": true,
"params": {"message": "hello"},
}),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
.expect("the stream must open even when the handler will fail");
let first = stream
.next_frame()
.await
.expect("a failed open is a terminal frame, never an empty end");
let error = first.expect_err("no provider is configured on a test machine");
assert!(
matches!(error, trusty_common::uds::UdsRpcError::Stream { .. }),
"expected the server's terminal error frame, got {error:?}"
);
assert!(
stream.next_frame().await.is_none(),
"the failure is reported once, then the stream is done"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_activity_stream_delivers_an_event_without_polling() {
let state = test_state();
let palace = seed_palace(&state, "stream-palace");
let daemon = Daemon::start(state).await;
let mut stream = send_framed_stream_request_capped::<_, Value>(
daemon.socket(),
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "memory.activity_stream",
"stream": true,
}),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
.expect("the stream opens");
daemon
.ok(
"memory.drawer_create",
json!({
"palace_id": palace,
"content": "an event the stream has to carry",
"force": true,
}),
)
.await;
let frame = tokio::time::timeout(CALL_TIMEOUT, stream.next_frame())
.await
.expect("an emitted event must reach an open stream promptly")
.expect("the stream carries the event rather than ending")
.expect("the frame is an item, not a terminal error");
assert_eq!(
frame["type"], "drawer_added",
"the frame body is the `type`-tagged DaemonEvent: {frame}"
);
assert_eq!(frame["palace_id"], palace.as_str());
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_activity_stream_does_not_replay_history() {
let state = test_state();
let palace = seed_palace(&state, "replay-palace");
let daemon = Daemon::start(state).await;
daemon
.ok(
"memory.drawer_create",
json!({
"palace_id": palace,
"content": "an event from before the stream opened",
"force": true,
}),
)
.await;
let mut stream = send_framed_stream_request_capped::<_, Value>(
daemon.socket(),
&json!({
"jsonrpc": "2.0",
"id": 1,
"method": "memory.activity_stream",
"stream": true,
}),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
.expect("the stream opens");
let quiet = tokio::time::timeout(Duration::from_millis(300), stream.next_frame()).await;
assert!(
quiet.is_err(),
"a freshly opened stream must not replay: {quiet:?}"
);
daemon
.ok(
"memory.drawer_create",
json!({
"palace_id": palace,
"content": "an event from after the stream opened",
"force": true,
}),
)
.await;
let frame = tokio::time::timeout(CALL_TIMEOUT, stream.next_frame())
.await
.expect("a later event still arrives")
.expect("the stream is open")
.expect("the frame is an item");
assert_eq!(frame["type"], "drawer_added");
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_serves_concurrent_connections() {
const CLIENTS: usize = 24;
let daemon = Daemon::start(test_state()).await;
let socket = daemon.socket().to_path_buf();
let calls = (0..CLIENTS).map(|i| {
let socket = socket.clone();
tokio::spawn(async move {
send_framed_request_capped::<_, RpcResponse>(
&socket,
&frame(i as i64, "memory.status", json!({})),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
})
});
let answers = futures::future::join_all(calls).await;
for (i, answer) in answers.into_iter().enumerate() {
let response = answer
.unwrap_or_else(|e| panic!("client {i} task panicked: {e}"))
.unwrap_or_else(|e| panic!("client {i} call failed: {e}"));
assert!(
response.result.is_some(),
"client {i} got an error: {:?}",
response.error
);
assert_eq!(response.id, json!(i as i64), "ids must not cross wires");
}
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_accepts_a_request_larger_than_the_shared_default() {
let daemon = Daemon::start(test_state()).await;
let padding = "x".repeat(12 * 1024 * 1024);
let response = send_framed_request_capped::<_, RpcResponse>(
daemon.socket(),
&frame(1, "memory.status", json!({"pad": padding})),
CALL_TIMEOUT,
MAX_FRAME_BYTES,
)
.await
.expect("a 12 MiB frame is inside this service's budget");
assert!(response.result.is_some(), "{:?}", response.error);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_refuses_a_request_past_its_own_budget() {
let daemon = Daemon::start(test_state()).await;
let padding = "x".repeat(MAX_FRAME_BYTES as usize);
let refused = send_framed_request_capped::<_, RpcResponse>(
daemon.socket(),
&frame(1, "memory.status", json!({"pad": padding})),
CALL_TIMEOUT,
MAX_FRAME_BYTES * 2,
)
.await;
assert!(
refused.is_err(),
"a frame past the server's budget must not be answered"
);
let after = daemon.ok("memory.status", json!({})).await;
assert!(
after["version"].is_string(),
"the accept loop must survive an oversized frame: {after}"
);
daemon.shutdown().await;
}
#[tokio::test(flavor = "multi_thread")]
async fn rpc_unlinks_its_socket_on_shutdown() {
let daemon = Daemon::start(test_state()).await;
let socket = daemon.socket().to_path_buf();
assert!(socket.exists(), "the socket must exist while serving");
daemon.shutdown().await;
assert!(
!socket.exists(),
"the socket file must be unlinked on shutdown: {}",
socket.display()
);
}
#[test]
fn socket_path_is_named_for_this_daemon() {
let path = socket_path().expect("the data directory must resolve");
assert!(
path.ends_with("trusty-memory.sock"),
"unexpected socket path: {}",
path.display()
);
}
#[test]
fn remove_if_present_deletes_a_stale_file_and_tolerates_an_absent_one() {
let tmp = tempfile::tempdir().expect("tempdir");
let stale = tmp.path().join("http_addr");
std::fs::write(&stale, "127.0.0.1:7070\n").expect("write");
super::remove_if_present(&stale);
assert!(!stale.exists(), "a stale discovery file must be removed");
super::remove_if_present(&stale);
}
fn seed_palace(state: &AppState, name: &str) -> String {
use trusty_common::memory_core::palace::{Palace, PalaceId};
let id = PalaceId::new(name);
let palace = Palace {
id: id.clone(),
name: name.to_string(),
description: None,
created_at: chrono::Utc::now(),
data_dir: state.data_root.join(name),
};
state
.registry
.create_palace(&state.data_root, palace)
.expect("create the test palace");
id.0
}
fn seed_palace_off_registry(state: &AppState, name: &str) -> String {
use trusty_common::memory_core::palace::{Palace, PalaceId};
use trusty_common::memory_core::PalaceRegistry;
let id = PalaceId::new(name);
let palace = Palace {
id: id.clone(),
name: name.to_string(),
description: None,
created_at: chrono::Utc::now(),
data_dir: state.data_root.join(name),
};
PalaceRegistry::new()
.create_palace(&state.data_root, palace)
.expect("create the test palace off-registry");
id.0
}