use futures::executor::block_on;
use pipecrab_core::{DataFrame, Direction, SystemFrame, Transcript};
use pipecrab_runtime::{Received, link};
#[test]
fn interrupt_preempts_backed_up_data() {
block_on(async {
let (out, mut inb) = link(16);
for i in 0..8 {
out.send_data(Transcript::user_final(i.to_string()).into())
.await
.unwrap();
}
out.send_system(Direction::Down, SystemFrame::Interrupt)
.await
.unwrap();
let r = inb.recv().await.unwrap();
assert!(
matches!(r, Received::Sys(Direction::Down, SystemFrame::Interrupt)),
"interrupt must jump the backlog, got {r:?}",
);
});
}
#[test]
fn fatal_error_propagates_upstream_ahead_of_data() {
block_on(async {
let (out, mut inb) = link(16);
for i in 0..8 {
out.send_data(Transcript::user_final(i.to_string()).into())
.await
.unwrap();
}
out.send_system(
Direction::Up,
SystemFrame::Error {
message: "inference exploded".into(),
fatal: true,
},
)
.await
.unwrap();
match inb.recv().await.unwrap() {
Received::Sys(Direction::Up, SystemFrame::Error { message, .. }) => {
assert_eq!(message, "inference exploded".into());
}
other => panic!("expected Sys(Up, Error), got {other:?}"),
}
});
}
#[test]
fn data_lane_is_fifo() {
block_on(async {
let (out, mut inb) = link(16);
for i in 0..4 {
out.send_data(Transcript::user_final(i.to_string()).into())
.await
.unwrap();
}
for i in 0..4 {
match inb.recv().await.unwrap() {
Received::Data(DataFrame::Transcript(s)) => {
assert_eq!(s.text, i.to_string().into())
}
other => panic!("expected Data(Transcript({i})), got {other:?}"),
}
}
});
}
#[test]
fn data_lane_is_always_downstream() {
block_on(async {
let (out, mut inb) = link(16);
out.send_data(Transcript::user_final("a").into())
.await
.unwrap();
out.send_data(Transcript::user_final("b").into())
.await
.unwrap();
assert!(matches!(inb.recv().await.unwrap(), Received::Data(_)));
assert!(matches!(inb.recv().await.unwrap(), Received::Data(_)));
});
}
#[test]
fn dropping_the_outbound_closes_both_lanes() {
block_on(async {
let (out, mut inb) = link(16);
out.send_data(Transcript::user_final("buffered").into())
.await
.unwrap();
drop(out);
match inb.recv().await.unwrap() {
Received::Data(DataFrame::Transcript(s)) => {
assert_eq!(s.text, "buffered".into())
}
other => panic!("buffered frame must survive sender drop, got {other:?}"),
}
assert!(
inb.recv().await.is_none(),
"closed lanes must signal shutdown via None"
);
});
}