use super::*;
use std::time::Duration;
use trusty_common::memory_core::timeouts;
const TEST_BUDGET: Duration = Duration::from_millis(300);
fn hang_guard() -> Duration {
timeouts::write_lock_timeout() + timeouts::open_queue_timeout() + Duration::from_secs(30)
}
async fn within_hang_guard<T>(what: &str, call: impl std::future::Future<Output = T>) -> T {
let guard = hang_guard();
tokio::time::timeout(guard, call)
.await
.unwrap_or_else(|_| panic!("{what} never returned (hang guard {guard:?})"))
}
fn reported_wait(msg: &str) -> Duration {
let token = msg
.split("timed out after ")
.nth(1)
.and_then(|rest| rest.split_whitespace().next())
.unwrap_or_else(|| panic!("the error must report the wait it gave up after; got: {msg}"));
let unit_at = token
.find(char::is_alphabetic)
.unwrap_or_else(|| panic!("no unit on reported wait {token:?}"));
let (value, unit) = token.split_at(unit_at);
let value: f64 = value
.parse()
.unwrap_or_else(|_| panic!("non-numeric reported wait {token:?}"));
let per_unit = match unit {
"s" => 1.0,
"ms" => 1e-3,
"\u{b5}s" => 1e-6,
"ns" => 1e-9,
other => panic!("unknown unit {other:?} on reported wait {token:?}"),
};
Duration::from_secs_f64(value * per_unit)
}
fn assert_waited_within_budget(tool: &str, msg: &str) {
let waited = reported_wait(msg);
assert!(
waited <= TEST_BUDGET,
"{tool} must give up inside its {TEST_BUDGET:?} budget; it reported waiting {waited:?} \
(pre-fix the write-lock leg alone waited {:?} — issue #4002); error: {msg}",
timeouts::write_lock_timeout()
);
}
fn budgeted_state() -> (AppState, tempfile::TempDir) {
skip_palace_enforcement();
seed_embedder();
let tmp = tempfile::tempdir().expect("tempdir");
let state = AppState::new(tmp.path().to_path_buf()).with_write_op_budget(TEST_BUDGET);
state.set_ready();
(state, tmp)
}
#[tokio::test]
async fn memory_remember_gives_up_within_one_budget_not_the_leg_sum() {
let (state, _tmp) = budgeted_state();
let _ = dispatch_tool(&state, "palace_create", json!({"name": "budget"}))
.await
.expect("palace_create");
let pre_fix_worst_case = timeouts::write_lock_timeout() + timeouts::open_queue_timeout();
assert!(
pre_fix_worst_case > TEST_BUDGET,
"the composed pre-fix wait ({pre_fix_worst_case:?}) must exceed the injected budget \
({TEST_BUDGET:?}), or this test proves nothing (issue #4002)"
);
let write_lock = state.palace_write_lock("budget");
let _held = write_lock.lock().await;
let err = within_hang_guard(
"memory_remember",
handle_memory_remember(
&state,
json!({"palace": "budget", "text": "a sufficiently long fact to clear the content gate"}),
),
)
.await
.expect_err("a permanently held write mutex must surface an error, never block forever");
let msg = format!("{err:#}");
assert!(
msg.contains("memory_remember") && msg.contains("write-lock acquisition timed out"),
"the error must name the tool and the leg that expired; got: {msg}"
);
assert_waited_within_budget("memory_remember", &msg);
}
#[tokio::test]
async fn memory_note_gives_up_within_one_budget_not_the_leg_sum() {
let (state, _tmp) = budgeted_state();
let _ = dispatch_tool(&state, "palace_create", json!({"name": "budget"}))
.await
.expect("palace_create");
let write_lock = state.palace_write_lock("budget");
let _held = write_lock.lock().await;
let err = within_hang_guard(
"memory_note",
handle_memory_note(
&state,
json!({
"palace": "budget",
"content": "Masa prefers snake_case for every identifier in this workspace",
}),
),
)
.await
.expect_err("a permanently held write mutex must surface an error");
let msg = format!("{err:#}");
assert!(
msg.contains("memory_note"),
"the error must name the tool that gave up; got: {msg}"
);
assert_waited_within_budget("memory_note", &msg);
}
#[tokio::test]
async fn a_short_budget_does_not_disturb_an_uncontended_write() {
let (state, _tmp) = budgeted_state();
let _ = dispatch_tool(&state, "palace_create", json!({"name": "budget"}))
.await
.expect("palace_create");
let stored = handle_memory_remember(
&state,
json!({"palace": "budget", "text": "an uncontended write must still land on disk"}),
)
.await
.expect("an uncontended write must succeed under a short budget");
assert_eq!(stored["status"], "stored");
let listed = dispatch_tool(&state, "memory_list", json!({"palace": "budget"}))
.await
.expect("memory_list");
assert_eq!(
listed["drawers"].as_array().map(Vec::len),
Some(1),
"the budgeted write must be readable afterwards"
);
}
#[tokio::test]
async fn task_add_gives_up_within_one_budget_not_the_leg_sum() {
let (state, _tmp) = budgeted_state();
let _ = dispatch_tool(&state, "palace_create", json!({"name": "budget"}))
.await
.expect("palace_create");
let write_lock = state.palace_write_lock("budget");
let _held = write_lock.lock().await;
let err = within_hang_guard(
"task_add",
dispatch_tool(
&state,
"task_add",
json!({"palace": "budget", "content": "ship the joint budget"}),
),
)
.await
.expect_err("a permanently held write mutex must surface an error");
let msg = format!("{err:#}");
assert!(
msg.contains("task_add"),
"the error must name the tool that gave up; got: {msg}"
);
assert_waited_within_budget("task_add", &msg);
}
#[test]
fn an_exhausted_write_budget_gives_the_open_queue_leg_nothing() {
let spent = timeouts::OpBudget::start(Duration::ZERO);
assert_eq!(
spent.leg(timeouts::open_queue_timeout()),
Duration::ZERO,
"once the write-lock leg has spent the budget, the open-queue leg must not open \
a fresh {:?} window (issue #4002)",
timeouts::open_queue_timeout()
);
}
const TEST_PIPELINE_BUDGET: Duration = Duration::ZERO;
async fn pipeline_capped_state() -> (AppState, tempfile::TempDir) {
skip_palace_enforcement();
seed_embedder();
let tmp = tempfile::tempdir().expect("tempdir");
let state =
AppState::new(tmp.path().to_path_buf()).with_write_pipeline_budget(TEST_PIPELINE_BUDGET);
state.set_ready();
let _ = dispatch_tool(&state, "palace_create", json!({"name": "pipeline"}))
.await
.expect("palace_create");
(state, tmp)
}
#[tokio::test]
async fn memory_note_surfaces_the_pipeline_ceiling() {
let (state, _tmp) = pipeline_capped_state().await;
let err = within_hang_guard(
"memory_note",
handle_memory_note(
&state,
json!({"palace": "pipeline", "content": "a curated fact that clears the content gate"}),
),
)
.await
.expect_err("a write that cannot fit its pipeline ceiling must error, not stall");
let msg = format!("{err:#}");
assert!(
msg.contains("#6366") && msg.contains("write pipeline exceeded"),
"the error must name the pipeline ceiling and its issue; got: {msg}"
);
}
#[tokio::test]
async fn a_refused_write_leaves_the_palace_write_mutex_free() {
let (state, _tmp) = pipeline_capped_state().await;
let _ = handle_memory_note(
&state,
json!({"palace": "pipeline", "content": "a curated fact that clears the content gate"}),
)
.await;
let write_lock = state.palace_write_lock("pipeline");
assert!(
write_lock.try_lock().is_ok(),
"#6366: a refused write must not leave the palace write mutex held — \
that would stall every writer behind it, which is the reported defect"
);
}