use super::*;
use crate::backend::fake::FakeBackend;
use onlyne_proto::new_task_id;
use tempfile::tempdir;
fn state_with_slot(task: &str) -> (DispatchState, Arc<FakeBackend>) {
let dir = tempdir().unwrap();
let store = ClientStore::open(dir.path().join("client.db")).expect("client store");
let backend = Arc::new(FakeBackend::new());
let state = DispatchState::new("planner", dir.path(), Vec::new(), 8, backend.clone(), store);
dispatch(
&state,
&new_envelope(
MsgKind::Task,
Principal::role("sender"),
Principal::role("planner"),
Body::text("repair the failing widget"),
Some(Causality::root(task.to_string())),
)
.expect("task envelope"),
)
.expect("the delivery takes a session");
(state, backend)
}
fn connection() -> (AdapterIo, tokio::io::DuplexStream) {
let (stream, peer) = tokio::io::duplex(1024);
(
AdapterIo::new(stream, Duration::from_secs(5), Duration::from_secs(5)),
peer,
)
}
fn completion(task: &str) -> onlyne_proto::Report {
onlyne_proto::Report::Complete {
task_id: task.to_string(),
outcome: onlyne_proto::Outcome::Done,
head: Some("done".to_string()),
details: None,
files: Vec::new(),
reply_to: None,
cluster_ref: None,
}
}
fn deliver(state: &DispatchState, io: &AdapterIo, role: &str) {
let mut envelope = new_envelope(
MsgKind::Task,
Principal::role("builder"),
Principal::role(role),
Body::text("carry it on"),
Some(Causality::root(new_task_id())),
)
.expect("a task envelope the protocol accepts");
state
.stamp_tools_send(io, &mut envelope)
.expect("the session's own record stamps the frame");
state
.plugin_send(io, &envelope)
.expect("the client carries the delivery");
}
#[tokio::test]
async fn a_session_owes_its_downstream_edges_and_never_the_role_it_answers() {
let task = new_task_id();
let (state, _backend) = state_with_slot(&task);
{
let mut inner = state.inner.lock();
inner.required_targets = vec!["sender".to_string()];
}
let token = state
.tools_token(&task)
.expect("the session minted its token");
let (io, _peer) = connection();
state.bind_tools_mount(&token, io.clone()).expect("bind");
assert!(
state
.completion_refusal(Some(&io), &completion(&task))
.is_none(),
"the role that handed this session its task is not owed a second delivery",
);
{
let mut inner = state.inner.lock();
inner.required_targets = vec!["sender".to_string(), "downstream".to_string()];
}
let refusal = state
.completion_refusal(Some(&io), &completion(&task))
.expect("a downstream edge is still owed");
assert_eq!(
refusal
.error
.expect("the refusal carries its error")
.message,
"relay guard: missing handoff to: downstream (this session delivered to: none)",
);
deliver(&state, &io, "downstream");
assert!(
state
.completion_refusal(Some(&io), &completion(&task))
.is_none(),
"the downstream delivery settles the obligation",
);
}
#[tokio::test]
async fn a_self_addressed_entry_owes_nothing() {
let task = new_task_id();
let (state, _backend) = state_with_slot(&task);
{
let mut inner = state.inner.lock();
let key = inner
.sessions
.keys()
.next()
.cloned()
.expect("the staged session");
inner
.sessions
.get_mut(&key)
.expect("the slot the task opened")
.origin = Some(Principal::role("planner"));
inner.required_targets = vec!["planner".to_string()];
}
let token = state
.tools_token(&task)
.expect("the session minted its token");
let (io, _peer) = connection();
state.bind_tools_mount(&token, io.clone()).expect("bind");
assert!(
state
.completion_refusal(Some(&io), &completion(&task))
.is_none(),
"a one-role cluster's self-addressed list is not an obligation it could discharge",
);
}