use super::*;
use crate::session::dispatch::transport::NUDGE_TEXT;
use crate::session::dispatch::turn_end::{DELIVERY_BLOCKED, TURN_END_WITHOUT_COMPLETE};
use onlyne_proto::new_task_id;
use serde_json::Value;
use tempfile::{TempDir, tempdir};
fn plugin_beat(agent: &str, delivery: &str, recovery: &str) -> Value {
serde_json::json!({
"version": { "generation": 1, "seq": 1005 },
"generation_live": true,
"isolate_after": 1,
"terminate_after": 3,
"mismatch_count": 0,
"agent": agent,
"delivery": delivery,
"resource": "attached",
"recovery": recovery,
})
}
fn staged_state(dir: &TempDir, task: &str) -> DispatchState {
let store = ClientStore::open(dir.path().join("client.db")).expect("client store");
let state = DispatchState::new(
"planner",
dir.path(),
Vec::new(),
8,
Arc::new(crate::backend::fake::FakeBackend::new()),
store,
);
state.inner.lock().bridge.track_live(SessionRef {
task_id: task.to_string(),
backend: "fake".into(),
backend_ref: Value::Null,
generation: 1,
});
state
}
fn seeded_ready(state: &DispatchState, task: &str) {
let inner = state.inner.lock();
feed_created(&inner.bridge, &inner.store, task).expect("seed the session row");
feed_ready(&inner.bridge, &inner.store, task).expect("ready");
}
async fn beat(state: &DispatchState, task: &str, spelling: &str, seq: u64) {
on_plugin_report(
state,
None,
Report::Heartbeat {
task_id: task.to_string(),
session_id: String::new(),
generation: 1,
seq,
observed: plugin_beat(spelling, "none", "none"),
projection: None,
cluster_ref: None,
},
)
.await
.expect("the beat is handled");
}
fn opened_task(state: &DispatchState, task: &str) {
let inner = state.inner.lock();
inner
.store
.open_task(&Causality::root(task.to_string()), "root")
.expect("open the task record");
}
fn serving_slot(state: &DispatchState, task: &str, msg_id: &str) {
let mut inner = state.inner.lock();
inner.sessions.insert(
task.to_string(),
SessionSlot {
keeps_idle: false,
family: None,
idle_since: None,
suspended: false,
opened_at: Instant::now(),
command: Vec::new(),
resume_handle: None,
session: SessionRef {
task_id: task.to_string(),
backend: "fake".into(),
backend_ref: Value::Null,
generation: 1,
},
task_id: Some(task.to_string()),
ready: true,
payload: None,
msg_id: Some(msg_id.to_string()),
origin: Some(Principal::role("reviewer")),
causality: Causality::root(task.to_string()),
dropped_at: None,
last_beat: None,
read_only: false,
tools_token: String::new(),
delivered_roles: BTreeSet::new(),
},
);
}
fn verdict(state: &DispatchState, task: &str) -> (TaskState, Option<String>) {
let inner = state.inner.lock();
let record = inner
.store
.task(task)
.expect("read the task record")
.expect("opened task record");
(
record.task_state,
inner.store.out_head(task).expect("read the head column"),
)
}
#[tokio::test]
async fn a_beat_below_the_rows_watermark_still_records_a_new_turn() {
let dir = tempdir().expect("tempdir");
let task = new_task_id();
let state = staged_state(&dir, &task);
seeded_ready(&state, &task);
{
let inner = state.inner.lock();
let row = inner
.store
.get_session(&task)
.expect("read the row")
.expect("the seeded row");
inner
.store
.bump_session_version(
&task,
row.generation.max(0) as u64,
row.seq.max(0) as u64 + 2000,
)
.expect("the client's own feeds moved the watermark");
}
beat(&state, &task, "running", 1001).await;
let inner = state.inner.lock();
let row = inner.store.get_session(&task).expect("read the row");
let observed = stored_observation(&inner.store, row.as_ref());
assert_eq!(
observed.agent,
AgentPhase::Running,
"a new beat whose sequence sits below the client's own watermark must still move the row"
);
}
#[tokio::test]
async fn a_turn_ending_over_an_open_task_nudges_once_and_then_settles() {
let dir = tempdir().expect("tempdir");
let task = new_task_id();
let state = staged_state(&dir, &task);
seeded_ready(&state, &task);
opened_task(&state, &task);
serving_slot(&state, &task, "msg-open");
let (client_side, test_side) = tokio::io::duplex(4096);
let serving = AdapterIo::new(client_side, Duration::from_secs(5), Duration::from_secs(5));
let (_plugin, mut inbound) =
AdapterIo::new_with_inbound(test_side, Duration::from_secs(5), Duration::from_secs(5));
{
let mut inner = state.inner.lock();
inner
.transports
.insert(task.clone(), (serving, vec![Capability::Inject]));
}
let queued = || -> Vec<ClientOp> {
let inner = state.inner.lock();
inner
.store
.due_intents(chrono::Utc::now(), 64)
.expect("read the outbound queue")
.iter()
.map(crate::runtime::intent::op_for_intent)
.collect::<Result<Vec<_>>>()
.expect("decode the queued ops")
};
beat(&state, &task, "running", 1005).await;
beat(&state, &task, "idle", 1006).await;
match inbound
.recv()
.await
.expect("the ending hands its plugin the one sentence")
.msg
{
AdapterMsg::Host(HostOp::Nudge {
task_id: nudged,
text,
}) => {
assert_eq!(nudged, task, "the sentence names the task it is about");
assert_eq!(text, NUDGE_TEXT, "the sentence is 3c's, verbatim");
}
other => panic!("a nudge is what an ending sends: {other:?}"),
}
assert_eq!(
verdict(&state, &task).0,
TaskState::Pending,
"a nudge is not a verdict: the delivery stays open"
);
beat(&state, &task, "running", 1007).await;
beat(&state, &task, "idle", 1008).await;
assert_eq!(
verdict(&state, &task).0,
TaskState::Blocked,
"the second turn that ends without a completion settles it blocked"
);
while let Ok(frame) = inbound.try_recv() {
assert!(
!matches!(frame.msg, AdapterMsg::Host(HostOp::Nudge { .. })),
"one nudge per delivery, and no second try: {:?}",
frame.msg
);
}
let rule: Vec<String> = queued()
.iter()
.filter_map(|op| match op {
ClientOp::PublishEvent(args) => Some(args.class.clone()),
_ => None,
})
.collect();
assert_eq!(
rule,
vec![
TURN_END_WITHOUT_COMPLETE.to_string(),
TURN_END_WITHOUT_COMPLETE.to_string(),
DELIVERY_BLOCKED.to_string(),
],
"every step of the rule is published, in order: {rule:?}"
);
let nudges: Vec<bool> = queued()
.iter()
.filter_map(|op| match op {
ClientOp::PublishEvent(args) if args.class == TURN_END_WITHOUT_COMPLETE => args
.payload
.get("nudge")
.and_then(serde_json::Value::as_bool),
_ => None,
})
.collect();
assert_eq!(nudges, vec![true, false], "one nudge, then the settlement");
let why: Vec<String> = queued()
.iter()
.filter_map(|op| match op {
ClientOp::PublishEvent(args) if args.class == DELIVERY_BLOCKED => args
.payload
.get("reason")
.and_then(serde_json::Value::as_str)
.map(str::to_string),
_ => None,
})
.collect();
assert_eq!(why.len(), 1, "the settlement publishes one reason");
assert!(
!why[0].is_empty(),
"a settlement that names no reason leaves a hook with the fact and none of the cause"
);
}
async fn settled_with_attached_runtime(
state: &DispatchState,
task: &str,
) -> (
AdapterIo,
tokio::sync::mpsc::Receiver<onlyne_adapter::IncomingFrame>,
) {
seeded_ready(state, task);
opened_task(state, task);
serving_slot(state, task, "msg-open");
{
let inner = state.inner.lock();
crate::reconcile::feed_resource_attached(&inner.bridge, &inner.store, task)
.expect("the resource is attached");
}
let (client_side, test_side) = tokio::io::duplex(4096);
let serving = AdapterIo::new(client_side, Duration::from_secs(5), Duration::from_secs(5));
let (plugin, inbound) =
AdapterIo::new_with_inbound(test_side, Duration::from_secs(5), Duration::from_secs(5));
{
let mut inner = state.inner.lock();
inner
.transports
.insert(task.to_string(), (serving, vec![Capability::Recycle]));
}
(plugin, inbound)
}
#[tokio::test]
async fn a_settled_non_kept_session_asks_its_runtime_to_leave() {
let dir = tempdir().expect("tempdir");
let task = new_task_id();
let state = staged_state(&dir, &task);
let (_plugin, mut inbound) = settled_with_attached_runtime(&state, &task).await;
on_out(
&state,
&task,
Outcome::Blocked,
Some("waiting on the outside".to_string()),
None,
SettleAuthority::ClientOwned,
)
.await
.expect("the ending settles");
match tokio::time::timeout(Duration::from_secs(5), inbound.recv())
.await
.expect("a leave request reaches the runtime")
.expect("the connection stays open")
.msg
{
AdapterMsg::Host(HostOp::Recycle(args)) => {
assert_eq!(args.task_id, task, "the request names the session it frees");
assert_eq!(args.outcome, None, "the leave request settles nothing");
}
other => panic!("a settled non-kept session owes its runtime a leave request: {other:?}"),
}
let inner = state.inner.lock();
assert!(
inner
.sessions
.get(&task)
.is_some_and(|slot| slot.task_id.is_none()),
"the session waits for its runtime to leave"
);
}
#[tokio::test]
async fn a_kept_session_between_deliveries_is_not_asked_to_leave() {
let dir = tempdir().expect("tempdir");
let task = new_task_id();
let state = staged_state(&dir, &task);
let (_plugin, mut inbound) = settled_with_attached_runtime(&state, &task).await;
{
let mut inner = state.inner.lock();
inner
.sessions
.get_mut(&task)
.expect("the staged slot")
.keeps_idle = true;
}
on_out(
&state,
&task,
Outcome::Done,
Some("finished".to_string()),
None,
SettleAuthority::ClientOwned,
)
.await
.expect("the ending settles");
assert!(
inbound.try_recv().is_err(),
"a kept session between deliveries is not asked to leave"
);
let inner = state.inner.lock();
assert!(
inner.sessions.contains_key(&task),
"the kept session stays for its family's next delivery"
);
}
#[tokio::test]
async fn a_close_reaches_a_settled_session_that_serves_no_task() {
let dir = tempdir().expect("tempdir");
let task = new_task_id();
let state = staged_state(&dir, &task);
let (_plugin, mut inbound) = settled_with_attached_runtime(&state, &task).await;
on_out(
&state,
&task,
Outcome::Done,
Some("finished".to_string()),
None,
SettleAuthority::ClientOwned,
)
.await
.expect("the ending settles");
let _leave = tokio::time::timeout(Duration::from_secs(5), inbound.recv())
.await
.expect("the leave request leaves first")
.expect("the connection stays open");
{
let mut inner = state.inner.lock();
super::super::retire::release_locked(
&mut inner,
&task,
Some(crate::backend::CloseReason::Cancelled),
)
.expect("the operator's close");
assert!(
!inner.sessions.contains_key(&task),
"the closed session leaves this client's books"
);
}
let inner = state.inner.lock();
let row = inner
.store
.get_session(&task)
.expect("read the row")
.expect("the session row");
assert_eq!(row.agent_state, "gone", "the close feeds the agent's exit");
assert_eq!(
row.resource_state, "closed",
"the close reaches the resource the settled session left behind"
);
}