- 1
//! A session ledger's lock belongs to its handle alone. The server reopens - 2
//! ledgers while other turns spawn tool workers, and a child being spawned - 3
//! holds duplicates of the server's descriptors until it execs. A lock - 4
//! released only by closing its descriptor stayed with such a child, so a - 5
//! reopen failed with `SessionError::Locked` while nothing held the ledger. - 6
- 7
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 8
- 9
use std::path::PathBuf; - 10
use std::sync::Arc; - 11
use std::sync::atomic::{AtomicBool, Ordering}; - 12
- 13
use serde_json::json; - 14
use vak_session::SessionLog; - 15
use vak_session::types::SessionError; - 16
use vak_tools::ToolContext; - 17
- 18
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 19
async fn reopening_a_ledger_while_workers_spawn_never_finds_it_locked() { - 20
let dir = tempfile::tempdir().unwrap(); - 21
let ledger = dir.path().join("session.jsonl"); - 22
std::fs::write(&ledger, "").unwrap(); - 23
std::fs::write(dir.path().join("note.txt"), "hello").unwrap(); - 24
- 25
let stop = Arc::new(AtomicBool::new(false)); - 26
let reopening = std::thread::spawn({ - 27
let stop = stop.clone(); - 28
move || { - 29
let (mut reopens, mut locked) = (0u32, 0u32); - 30
while !stop.load(Ordering::Relaxed) { - 31
match SessionLog::open(ledger.clone()) { - 32
Ok(log) => drop(log), - 33
Err(SessionError::Locked(_)) => locked += 1, - 34
Err(error) => panic!("reopen failed: {error}"), - 35
} - 36
reopens += 1; - 37
} - 38
(reopens, locked) - 39
} - 40
}); - 41
- 42
let worker = PathBuf::from(env!("CARGO_BIN_EXE_vak-tool-worker")); - 43
let read = vak_tools::brokered_default_tools(worker) - 44
.into_iter() - 45
.find(|tool| tool.name() == "read") - 46
.expect("read is a built-in tool"); - 47
let ctx = ToolContext::new(dir.path().to_path_buf()); - 48
let mut lanes = tokio::task::JoinSet::new(); - 49
for _ in 0..4 { - 50
let (read, ctx) = (read.clone(), ctx.clone()); - 51
lanes.spawn(async move { - 52
for _ in 0..25 { - 53
let output = read.execute(&json!({"path": "note.txt"}), &ctx).await; - 54
assert!(!output.is_error, "{}", output.content); - 55
} - 56
}); - 57
} - 58
while let Some(lane) = lanes.join_next().await { - 59
lane.unwrap(); - 60
} - 61
stop.store(true, Ordering::Relaxed); - 62
- 63
let (reopens, locked) = reopening.join().unwrap(); - 64
assert!(reopens > 0, "the ledger was never reopened"); - 65
assert_eq!( - 66
locked, 0, - 67
"{locked} of {reopens} reopens found the ledger locked while workers spawned" - 68
); - 69
} - 70
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.