use super::*;
#[test]
fn replay_is_bounded_by_event_count() {
let frame = ServerFrame::new(ServerMessage::Error {
code: "test".into(),
message: String::new(),
fatal: false,
});
let mut replay = VecDeque::from(vec![
ReplayEntry::new(frame.clone()).expect("measure frame");
REPLAY_CAPACITY
]);
let mut replay_bytes = serde_json::to_vec(&frame).expect("encode frame").len() * replay.len();
let (events, _) = broadcast::channel(1);
FRAME_SIZE_MEASUREMENTS.with(|count| count.set(0));
assert!(
publish_frame(&mut replay, &mut replay_bytes, &events, frame, true).expect("record event")
);
assert_eq!(replay.len(), REPLAY_CAPACITY);
assert_eq!(FRAME_SIZE_MEASUREMENTS.with(std::cell::Cell::get), 1);
}
#[test]
fn replay_is_bounded_by_encoded_bytes() {
let (events, _) = broadcast::channel(1);
let mut replay = VecDeque::new();
let mut replay_bytes = 0;
let large_message = "x".repeat(MAX_REPLAY_BYTES / 2);
let first = ServerFrame::new(ServerMessage::Error {
code: "first".into(),
message: large_message.clone(),
fatal: false,
});
let second = ServerFrame::new(ServerMessage::Error {
code: "second".into(),
message: large_message,
fatal: false,
});
assert!(
!publish_frame(&mut replay, &mut replay_bytes, &events, first, true)
.expect("record first frame")
);
assert!(
publish_frame(&mut replay, &mut replay_bytes, &events, second, true)
.expect("record second frame")
);
assert_eq!(replay.len(), 1);
assert!(replay_bytes <= MAX_REPLAY_BYTES);
}
#[test]
fn suppressed_frames_enter_replay_without_broadcasting() {
let mut replay = VecDeque::new();
let mut replay_bytes = 0;
let (events, mut receiver) = broadcast::channel(4);
let history = ServerFrame::new(ServerMessage::Error {
code: "history".into(),
message: "recorded only".into(),
fatal: false,
});
publish_frame(
&mut replay,
&mut replay_bytes,
&events,
history.clone(),
true,
)
.expect("record history");
assert_eq!(replay.back().map(|entry| &entry.frame), Some(&history));
assert!(matches!(
receiver.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
}
#[test]
fn transient_controls_are_broadcast_without_entering_replay() {
let (events, mut receiver) = broadcast::channel(5);
let mut replay = VecDeque::new();
let mut replay_bytes = 0;
let messages = [
EventMsg::SessionResumeRequested(mobius::protocol::SessionResumeRequestedEvent {
session_id: "target".into(),
context: Default::default(),
}),
EventMsg::Frontend(FrontendEvent::Preview {
symbol: None,
duration_ms: None,
started_at_ms: None,
id: "preview".into(),
title: "Preview".into(),
subtitle: String::new(),
page_id: "preview:latest".into(),
update: mobius::protocol::FrontendPreviewUpdate::Replace,
events: Vec::new(),
next: None,
}),
EventMsg::Frontend(FrontendEvent::Picker {
title: "Choose".into(),
options: Vec::new(),
}),
EventMsg::Frontend(FrontendEvent::Widget {
capability: "test".into(),
item: mobius::protocol::FrontendWidget {
id: "status".into(),
slot: mobius::protocol::FrontendSlot::Header,
text: "Current".into(),
tone: mobius::protocol::FrontendTone::Neutral,
symbol: None,
icon_only: false,
progress: None,
content: None,
action: None,
},
}),
EventMsg::Frontend(FrontendEvent::RemoveWidget {
capability: "test".into(),
id: "status".into(),
}),
];
for (index, msg) in messages.into_iter().enumerate() {
let frame = ServerFrame::new(ServerMessage::AgentEvent {
session_id: "source".into(),
record: RecordedEvent {
sequence: u64::try_from(index + 1).expect("sequence"),
recorded_at_ms: 1,
event: Event {
submission_id: Some("transient".into()),
msg,
},
stream_metrics: Vec::new(),
blocks: Vec::new(),
preview: None,
},
});
publish_frame(
&mut replay,
&mut replay_bytes,
&events,
frame.clone(),
false,
)
.expect("broadcast transient control");
assert_eq!(receiver.try_recv().expect("live transient control"), frame);
}
assert!(replay.is_empty());
assert_eq!(replay_bytes, 0);
}
#[test]
fn completed_step_compacts_only_its_progressive_replay_frames() {
let frame = |sequence, msg| {
ServerFrame::new(ServerMessage::AgentEvent {
session_id: "session".into(),
record: RecordedEvent {
sequence,
recorded_at_ms: 1,
event: Event {
submission_id: Some("submission".into()),
msg,
},
stream_metrics: Vec::new(),
blocks: Vec::new(),
preview: None,
},
})
};
let frames = [
frame(
1,
EventMsg::AssistantContentDelta(mobius::protocol::AssistantContentDeltaEvent {
session_id: "session".into(),
turn_id: "turn".into(),
model_step_id: "completed".into(),
delta: "answer".into(),
phase: mobius::protocol::ModelStepContentPhase::FinalAnswer,
}),
),
frame(
2,
EventMsg::AssistantContentDelta(mobius::protocol::AssistantContentDeltaEvent {
session_id: "session".into(),
turn_id: "turn".into(),
model_step_id: "completed".into(),
delta: "reasoning".into(),
phase: mobius::protocol::ModelStepContentPhase::Reasoning,
}),
),
frame(
3,
EventMsg::AssistantContentDelta(mobius::protocol::AssistantContentDeltaEvent {
session_id: "session".into(),
turn_id: "turn".into(),
model_step_id: "active".into(),
delta: "partial".into(),
phase: mobius::protocol::ModelStepContentPhase::FinalAnswer,
}),
),
];
let mut replay = frames
.into_iter()
.map(|frame| ReplayEntry::new(frame).expect("measure frame"))
.collect::<VecDeque<_>>();
let mut replay_bytes = replay
.iter()
.map(|frame| {
serde_json::to_vec(&frame.frame)
.expect("encode frame")
.len()
})
.sum();
FRAME_SIZE_MEASUREMENTS.with(|count| count.set(0));
compact_replay_deltas(&mut replay, &mut replay_bytes, "completed");
assert_eq!(FRAME_SIZE_MEASUREMENTS.with(std::cell::Cell::get), 0);
assert_eq!(replay.len(), 1);
assert_eq!(
replay
.front()
.and_then(|entry| event_sequence(&entry.frame)),
Some(3)
);
assert_eq!(
replay_bytes,
serde_json::to_vec(&replay.front().expect("remaining frame").frame)
.expect("encode remaining frame")
.len()
);
}
#[test]
fn replacement_startup_is_published_only_after_ready() {
let (events, mut receiver) = broadcast::channel(4);
let ready = ServerFrame::new(ServerMessage::Error {
code: "ready".into(),
message: String::new(),
fatal: false,
});
let startup = ServerFrame::new(ServerMessage::Error {
code: "startup".into(),
message: String::new(),
fatal: false,
});
publish_ready_and_pending(&events, ready, vec![startup]);
assert!(matches!(
receiver.try_recv().expect("ready frame").message,
ServerMessage::Error { code, .. } if code == "ready"
));
assert!(matches!(
receiver.try_recv().expect("startup frame").message,
ServerMessage::Error { code, .. } if code == "startup"
));
}
fn publish_frame(
replay: &mut VecDeque<ReplayEntry>,
replay_bytes: &mut usize,
events: &broadcast::Sender<ServerFrame>,
frame: ServerFrame,
suppress_broadcast: bool,
) -> Result<bool> {
Ok(record_and_publish(
replay,
replay_bytes,
events,
ReplayEntry::new(frame)?,
suppress_broadcast,
))
}