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_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"
);
}