use super::*;
use crate::kv_integration::ExactStateRecordAdmission;
use std::sync::atomic::Ordering;
#[test]
fn l1_is_visible_while_l3_spill_is_blocked() {
let radix = Mutex::new(UnifiedRadixCache::new());
let blobs = Mutex::new(CacheBlobStore::new(4));
let (root, tier) = test_l3("l1-before-l3");
let budget = StorageBudget::new();
let spill_gate = (Mutex::new((false, false)), std::sync::Condvar::new());
std::thread::scope(|scope| {
let spill_gate_ref = &spill_gate;
let before_l3_spill = move || {
let (state, ready) = spill_gate_ref;
let mut state = state.lock().unwrap();
state.0 = true;
ready.notify_one();
while !state.1 {
state = ready.wait(state).unwrap();
}
};
let worker_radix = &radix;
let worker_blobs = &blobs;
let worker_tier = &tier;
let worker_budget = &budget;
let worker = scope.spawn(move || {
store_exact_radix_record_with_codec(
worker_radix,
worker_blobs,
1,
limits(0, 0),
None,
DurableRecordTarget {
l3: Some(worker_tier),
cachegen_enabled: false,
before_l3_spill: Some(&before_l3_spill),
},
pending("first", &[1, 2], b"first-exact-state", worker_budget),
)
});
let (state, ready) = &spill_gate;
let mut state = state.lock().unwrap();
while !state.0 {
state = ready.wait(state).unwrap();
}
drop(state);
let visible_page_id = radix
.lock()
.unwrap()
.lookup_recurrent("model", &[1, 2])
.map(|lookup| lookup.value.page_id);
let mut state = spill_gate.0.lock().unwrap();
state.1 = true;
ready.notify_one();
drop(state);
worker.join().unwrap().unwrap();
assert_eq!(visible_page_id.as_deref(), Some("first"));
});
assert!(
tier.locate_longest("model", &[1, 2], 2).unwrap().is_some(),
"released spill must still persist the durable entry"
);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn l3_refusal_preserves_l1_record() {
let radix = Mutex::new(UnifiedRadixCache::new());
let blobs = Mutex::new(CacheBlobStore::new(4));
let root = std::env::temp_dir()
.join("skippy-server-l3-tests")
.join(format!("refusal-preserves-l1-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let tier = L3Tier::open(&root, 4, "blake3:test-tier".to_string(), 4096).unwrap();
let budget = StorageBudget::new();
store_exact_radix_record(
&radix,
&blobs,
1,
limits(0, 0),
None,
Some(&tier),
pending("first", &[1, 2], b"first-exact-state", &budget),
)
.unwrap();
assert_eq!(
radix
.lock()
.unwrap()
.lookup_recurrent("model", &[1, 2])
.expect("L3 refusal must leave the L1 entry usable")
.value
.page_id,
"first"
);
assert!(
tier.locate_longest("model", &[1, 2], 2).unwrap().is_none(),
"refused durable entry must not be published in L3"
);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn resident_only_exact_record_does_not_spill_to_l3() {
let radix = Mutex::new(UnifiedRadixCache::new());
let blobs = Mutex::new(CacheBlobStore::new(4));
let (root, tier) = test_l3("resident-only");
let budget = StorageBudget::new();
let mut record = pending("full-prompt", &[1, 2], b"full-prompt-state", &budget);
record.write_through_l3 = false;
store_exact_radix_record(&radix, &blobs, 1, limits(0, 0), None, Some(&tier), record).unwrap();
let radix_entry = radix
.lock()
.unwrap()
.lookup_recurrent("model", &[1, 2])
.expect("resident-only exact state must remain available in L1")
.value
.clone();
assert_eq!(radix_entry.page_id, "full-prompt");
assert!(
!radix_entry.l3_promotion_eligible,
"a later L1 hit must not promote an off-checkpoint state into L3"
);
assert!(
tier.locate_longest("model", &[1, 2], 2).unwrap().is_none(),
"resident-only exact state must not create a durable manifest"
);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn blocked_l3_spill_does_not_head_of_line_block_later_l1_records() {
struct SpillPauseGuard(Arc<std::sync::atomic::AtomicBool>);
impl Drop for SpillPauseGuard {
fn drop(&mut self) {
self.0.store(false, Ordering::Release);
}
}
let root = std::env::temp_dir()
.join("skippy-server-l3-manager-tests")
.join(format!("l1-while-l3-blocked-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let manager = L3CacheManager::acquire(&root, StoreLimits::new(1_000_000, 0)).unwrap();
let mut config = enabled_auto_config("future/model");
config.kv_cache.as_mut().unwrap().payload = StageKvCachePayload::FullState;
let kv = KvStageIntegration::from_loaded_model_with_l3_manager(
&config,
Some(ModelStateKind::Dense),
None,
Some(manager),
None,
)
.unwrap()
.expect("disk-backed exact cache should be enabled");
let budget = StorageBudget::new();
let wait = std::time::Duration::from_secs(30);
kv.l3_spill_worker_pause.store(true, Ordering::Release);
let spill_pause_guard = SpillPauseGuard(kv.l3_spill_worker_pause.clone());
assert_eq!(
kv.enqueue_exact_state_record(pending("first", &[1, 2], b"first-exact-state", &budget,)),
ExactStateRecordAdmission::Queued,
);
kv.wait_for_exact_state_recording(wait)
.expect("first L1 record should publish");
let deadline = std::time::Instant::now() + wait;
while kv.l3_spill_worker_received.load(Ordering::Acquire) == 0 {
assert!(
std::time::Instant::now() < deadline,
"durable worker did not receive the first spill"
);
std::thread::sleep(std::time::Duration::from_millis(2));
}
assert_eq!(
kv.enqueue_exact_state_record(pending(
"second",
&[1, 2, 3],
b"second-exact-state",
&budget,
)),
ExactStateRecordAdmission::Queued,
);
kv.wait_for_exact_state_recording(wait)
.expect("second L1 record must not wait for the first L3 spill");
let (both_visible, promotion_eligible) = {
let mut radix = kv.radix.lock().unwrap();
let both_visible = radix.recurrent_exact("model", &[1, 2]).is_some()
&& radix.recurrent_exact("model", &[1, 2, 3]).is_some();
let promotion_eligible = radix
.peek_recurrent("model", &[1, 2, 3])
.is_some_and(|entry| entry.value.l3_promotion_eligible);
(both_visible, promotion_eligible)
};
drop(spill_pause_guard);
assert!(
both_visible,
"both L1 records must be visible while the first durable spill is blocked"
);
assert!(
promotion_eligible,
"an async spill refusal must leave the L1 entry eligible for later promotion"
);
drop(kv);
let _ = std::fs::remove_dir_all(root);
}