use std::fs;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::mpsc;
use std::time::Duration;
use tempfile::tempdir;
use super::{Frame, Lane};
use crate::boundary::answer::answer;
use crate::boundary::dispatch::{Caller, Deps, dispatch};
use crate::boundary::reply::Reply;
use crate::boundary::{Action, Query};
use crate::cli_outbound::{Chunk, Cli, ExitInfo, Streamed};
use crate::login::runs::Runs;
use crate::login::{LoginRun, LoginView};
use crate::ui_state::UiState;
fn workspace() -> PathBuf {
crate::test_support::world::fixture_workspace()
}
fn deps_with(root: &Path, bz: &Path) -> Deps {
let world = crate::test_support::world::world_under(root);
fs::create_dir_all(world.yog_state_root()).expect("state root");
Deps {
litany: Cli::new("/definitely/not/a/litany-xyz"),
bl: Cli::new("/definitely/not/a/bl-xyz"),
state_root: world.yog_state_root(),
yog_binary: root.join("yog"),
world,
home: root.join("home"),
yog_data_root: root.join("data/yog"),
balls_state_root: root.join("state/balls"),
snapshot: Arc::new(crate::boundary::tests::snapshot(
&workspace(),
"alba",
vec![],
vec![],
)),
caller: Caller {
logins: Runs::of(Cli::new(bz)),
..Caller::default()
},
}
}
fn script(dir: &Path, name: &str, body: &str) -> PathBuf {
let path = dir.join(name);
crate::test_support::write_exec(&path, body);
path
}
fn ask(deps: &Deps, provider: &str) -> LoginView {
let ui = UiState::open(PathBuf::from("/nonexistent/ui.json"));
let asked = Query::LoginTail {
workspace: crate::naming::leaf(&workspace()),
provider: provider.to_owned(),
};
match answer(&asked, deps, &ui, 0) {
Ok(Reply::Login(view)) => view,
other => panic!("the sign-in lane answers a standing: {other:?}"),
}
}
fn settle(deps: &Deps, provider: &str) -> LoginView {
for _ in 0..600 {
let view = ask(deps, provider);
if view.outcome.is_some() {
return view;
}
std::thread::sleep(Duration::from_millis(5));
}
panic!("the engine never settled {provider}'s run");
}
fn wired(state_root: &Path) -> (LoginRun, mpsc::Sender<Chunk>) {
let (tx, rx) = mpsc::channel();
(
LoginRun::from_streamed(Streamed::from_rx(rx), vec!["bz".to_owned()], state_root),
tx,
)
}
fn lane(runs: &Runs, provider: &str) -> Lane {
Lane::holding(
runs.clone(),
workspace(),
provider.to_owned(),
2,
Duration::ZERO,
)
}
#[test]
fn the_act_starts_the_run_answers_its_standing_and_the_read_says_the_rest() {
let dir = tempdir().expect("tmp");
let bz = script(
dir.path(),
"bz",
"#!/bin/sh\nprintf 'open https://x/auth\\n' 1>&2\nexit 0\n",
);
let deps = deps_with(dir.path(), &bz);
let mut ui = UiState::open(PathBuf::from("/nonexistent/ui.json"));
assert_eq!(ask(&deps, "openai"), LoginView::default());
let act = Action::Login {
workspace: crate::naming::leaf(&workspace()),
provider: "openai".to_owned(),
};
let Ok(Reply::Login(receipt)) = dispatch(&deps, &mut ui, "100", &act) else {
panic!("the sign-in act answers a standing");
};
assert_eq!(receipt.outcome, None, "the act never waits for the run");
let settled = settle(&deps, "openai");
assert_eq!(settled.outcome, Some(0));
assert_eq!(
settled
.lines
.iter()
.map(|l| l.text.clone())
.collect::<Vec<_>>(),
["open https://x/auth".to_owned()]
);
}
#[test]
fn a_lane_on_a_pair_with_no_run_opens_on_the_emptiness_and_then_ends() {
let runs = Runs::default();
let mut lane = lane(&runs, "openai");
assert_eq!(lane.next(), Some(Reply::Login(LoginView::default())));
assert_eq!(lane.next(), None);
}
#[test]
fn a_lane_replays_from_the_start_appends_and_ends_on_the_settled_exit() {
let dir = tempdir().expect("tmp");
let runs = Runs::default();
let (run, tx) = wired(dir.path());
let serial = runs.seat(&workspace(), "openai", run, 0);
tx.send(Chunk::Stderr(b"open https://x/auth\n".to_vec()))
.expect("send");
runs.read_once(&workspace(), "openai", serial);
let mut lane = lane(&runs, "openai");
let Some(Reply::Login(first)) = lane.next() else {
panic!("a lane opens on the standing");
};
assert_eq!(first.lines.len(), 1);
assert_eq!(first.outcome, None);
assert!(matches!(lane.poll(), Frame::Waiting));
tx.send(Chunk::Stderr(b"waiting\n".to_vec())).expect("send");
runs.read_once(&workspace(), "openai", serial);
let Frame::Ready(second) = lane.poll() else {
panic!("an append is a frame");
};
assert_eq!(
second
.lines
.iter()
.map(|l| l.text.clone())
.collect::<Vec<_>>(),
["waiting".to_owned()],
"the append, not the whole buffer again"
);
tx.send(Chunk::Exited(ExitInfo::Code(78))).expect("send");
runs.read_once(&workspace(), "openai", serial);
let Some(Reply::Login(last)) = lane.next() else {
panic!("the exit is a frame");
};
assert_eq!(last.outcome, Some(78));
assert!(
last.fallback.is_some(),
"a non-zero exit carries the §8.3 fallback"
);
assert!(matches!(lane.poll(), Frame::Over));
assert_eq!(lane.next(), None);
}
#[test]
fn a_lane_ends_when_the_run_it_was_reading_is_gone() {
let dir = tempdir().expect("tmp");
let world = crate::test_support::world::world_under(dir.path());
let state = tempdir().expect("tmp");
let runs = Runs::of(Cli::new("/bin/sh"));
let (run, _tx) = wired(dir.path());
runs.seat(&workspace(), "openai", run, 0);
let mut lane = lane(&runs, "openai");
assert!(matches!(lane.poll(), Frame::Ready(_)), "it opened on a run");
runs.start(&world, &workspace(), "other", state.path(), "4000")
.expect("started");
assert!(matches!(lane.poll(), Frame::Over));
assert_eq!(lane.next(), None);
}
#[test]
fn a_quiet_lane_ends_on_its_hold_rather_than_holding_forever() {
let dir = tempdir().expect("tmp");
let runs = Runs::default();
let (run, _tx) = wired(dir.path());
runs.seat(&workspace(), "openai", run, 0);
let mut lane = lane(&runs, "openai");
assert!(lane.next().is_some(), "the opening frame");
assert_eq!(lane.next(), None);
}