mod common;
use common::{fixture_with_dropped_deliberation, fixture_with_dropped_stream};
use magi::graph::Runner;
use magi::run::RunStatus;
#[tokio::test]
async fn a_dropped_stream_is_resumed_once_and_the_candidate_recovers() {
let _home = common::home_lock().await;
let fx = fixture_with_dropped_stream(_home, &["impl-B"]);
let mut runner = Runner::start(&fx.repo, "create note.txt".to_owned(), fx.config.clone())
.await
.expect("start");
runner.execute().await.expect("execute");
let state = &runner.state;
let events: Vec<&str> = state
.events
.iter()
.filter(|e| e.node == "implement")
.map(|e| e.message.as_str())
.collect();
assert!(
events
.iter()
.any(|m| m.contains("resuming the conversation")),
"{events:?}"
);
let art = state.dir().join("artifacts");
assert!(
art.join("impl-B-resume.out").exists(),
"the resumed call must have run"
);
let b = state
.candidates
.iter()
.find(|c| c.label == 'B')
.expect("candidate B");
assert!(!b.empty, "{b:?}");
assert!(b.failed.is_none(), "{b:?}");
assert!(b.commits > 0, "{b:?}");
}
#[tokio::test]
async fn a_dropped_stream_with_no_session_left_is_not_resumed_into_a_blank_prompt() {
let _home = common::home_lock().await;
let mut fx = fixture_with_dropped_stream(_home, &["impl-B"]);
fx.config.graph.sessions = false;
let mut runner = Runner::start(&fx.repo, "create note.txt".to_owned(), fx.config.clone())
.await
.expect("start");
runner.execute().await.expect("execute");
let state = &runner.state;
let events: Vec<&str> = state
.events
.iter()
.filter(|e| e.node == "implement")
.map(|e| e.message.as_str())
.collect();
assert!(
events
.iter()
.any(|m| m.contains("no session left to resume")),
"{events:?}"
);
assert!(
!events
.iter()
.any(|m| m.contains("resuming the conversation")),
"{events:?}"
);
let art = state.dir().join("artifacts");
assert!(
!art.join("impl-B-resume.out").exists(),
"there is nothing to resume into, so no retry call should have run"
);
let b = state
.candidates
.iter()
.find(|c| c.label == 'B')
.expect("candidate B");
assert!(b.empty, "{b:?}");
assert_eq!(b.commits, 0, "{b:?}");
}
#[tokio::test]
async fn a_dropped_deliberation_turn_is_skipped_not_read_as_the_judges_position() {
let _home = common::home_lock().await;
let fx = fixture_with_dropped_deliberation(_home, &["judge-1"]);
let mut runner = Runner::start(&fx.repo, "create note.txt".to_owned(), fx.config.clone())
.await
.expect("start");
runner.execute().await.expect("execute");
let state = &runner.state;
assert_eq!(state.deliberation.len(), 1);
let round = &state.deliberation[0];
assert_eq!(round.turns.len(), 2, "{:?}", round.turns);
assert!(
round.turns.iter().all(|t| t.judge != 1),
"{:?}",
round.turns
);
for t in &round.turns {
assert!(
!t.body.contains("conversation_id") && !t.body.contains("subscriber fell behind"),
"a judge's position must never be the CLI's raw error JSON: {:?}",
t.body
);
}
let events: Vec<&str> = state
.events
.iter()
.filter(|e| e.node == "deliberate")
.map(|e| e.message.as_str())
.collect();
assert!(
events
.iter()
.any(|m| m.contains("judge 1 skipped") && m.contains("dropped the stream")),
"{events:?}"
);
assert_eq!(
state.status,
RunStatus::Ready,
"{}",
magi::report::run(state)
);
}