use std::path::PathBuf;
use std::rc::Rc;
use std::sync::Arc;
use std::time::Duration;
use bash_interop::rig::{Answer, Driving, ExitStatus, Failure, Layout, Message, Reacting, Rig, Run, Shell};
use tokio::sync::Notify;
use crate::{ENTRY, beginning, behind, provisioned, report};
use bash_interop::scratch::{Scripts, bash, sourcing};
struct Answering {
steps: PathBuf,
}
struct Soak {
steps: PathBuf,
heard: Vec<Message>,
answered: usize,
}
impl AsRef<[Message]> for Soak {
fn as_ref(&self) -> &[Message] {
&self.heard
}
}
impl Rig for Answering {
type Reaction = Soak;
fn bash(&self, _at: &Layout) -> String {
r#"
NOTE() { BC_INSTR SOAK say NOTE "$@"; }
"#
.to_string()
}
async fn joined(&self, _at: &Layout, _shell: Arc<Shell>) -> Result<Soak, Failure> {
Ok(Soak {
steps: self.steps.clone(),
heard: Vec::new(),
answered: 0,
})
}
}
impl Driving for Answering {}
impl Reacting for Soak {
type Kept = Self;
async fn hear(&mut self, said: Message) -> Result<(), Failure> {
self.heard.push(said);
Ok(())
}
async fn answer(&mut self, asked: Message) -> Result<Answer, Failure> {
let step: usize = asked
.words
.last()
.and_then(|word| word.parse().ok())
.unwrap_or(0);
self.answered += 1;
self.heard.push(asked);
Ok(match step % 7 {
0 => Answer::status(0),
1 => Answer::of(
"declare",
["-g".to_string(), format!("mark_{step}=set")],
),
2 => Answer::of("eval", [format!("NOTE eval {step}")]),
3 => Answer::of(
"NOTE",
["call".to_string(), step.to_string()],
),
4 => {
let step_bash = self.steps.join(format!("step.{step}.bash"));
return sourcing(
&step_bash,
&format!("NOTE source {step}"),
);
}
5 => {
std::thread::sleep(Duration::from_millis(2));
Answer::status(0)
}
6 => Answer::of(
"printf",
["%s".to_string(), "x".repeat(100_000)],
),
_ => Answer::status(3),
})
}
async fn finish(self) -> Result<Self, Failure> {
Ok(self)
}
}
#[tokio::test]
async fn a_session_survives_every_way_of_answering() {
let scripts = Scripts::of(&[
(
ENTRY,
r#"
declare -i i=0
while (( i < 56 )); do
BC_INSTR SOAK say REC tick "$i"
if (( i % 7 == 6 )); then
got="$(BC_INSTR SOAK ask step "$i")"
BC_INSTR SOAK say REC big "${#got}"
else
BC_INSTR SOAK ask step "$i" || BC_INSTR SOAK say REC refused "$i"
fi
(( i += 1 ))
done
wide="$(printf 'W%.0s' {1..9000})"
BC_INSTR SOAK say REC wide "$wide"
bash "${BASH_SOURCE[0]%/*}/other.bash"
BC_INSTR SOAK say REC marks ${!mark_@}
"#,
),
(
"other.bash",
r#"
BC_INSTR SOAK ask step 4
BC_INSTR SOAK say REC other done
"#,
),
]);
let answering = Answering {
steps: scripts.dir().to_path_buf(),
};
let ran = answering
.run(
&bash(scripts.at(ENTRY)),
provisioned(&answering),
)
.await
.and_then(Run::whole)
.unwrap_or_else(|error| panic!("{error}"));
assert_eq!(
ran.subject,
ExitStatus::Code(0),
"{}",
report(&ran.shells)
);
assert_eq!(
ran.shells
.iter()
.map(|at| at.kept.answered)
.collect::<Vec<_>>(),
[48, 1, 1, 1, 1, 1, 1, 1, 1, 1],
"each shell's own questions, counted where they were answered — the \
eight command substitutions are shells of their own, then the child"
);
let said = behind(&ran.shells, "REC");
assert_eq!(beginning(&said, "tick"), 56);
assert_eq!(beginning(&said, "refused"), 0);
assert_eq!(
beginning(&said, "big"),
8,
"{}",
report(&ran.shells)
);
assert!(
said.iter()
.filter(|words| words[0] == "big")
.all(|words| words[1] == "100000")
);
assert_eq!(
beginning(&said, "other"),
1,
"the second shell got its answer too"
);
assert!(
said.iter()
.any(|words| words.iter().any(|word| word.len() == 9000)),
"the wide message arrived as exactly what was written"
);
let notes = behind(&ran.shells, "NOTE");
for form in ["eval", "call", "source"] {
assert!(
beginning(¬es, form) > 0,
"no answer arrived by {form}{}",
report(&ran.shells)
);
}
let marks = said
.iter()
.find(|words| words.first().is_some_and(|first| first == "marks"))
.expect("the marks message");
for name in ["mark_1", "mark_50"] {
assert!(
marks.iter().any(|word| word == name),
"`declare -g` reached the subject's own scope, but {name} is missing from {marks:?}"
);
}
}
struct Gated {
open: Rc<Notify>,
}
struct Gate {
open: Rc<Notify>,
}
impl Rig for Gated {
type Reaction = Gate;
fn bash(&self, _at: &Layout) -> String {
String::new()
}
async fn joined(&self, _at: &Layout, _shell: Arc<Shell>) -> Result<Gate, Failure> {
Ok(Gate {
open: Rc::clone(&self.open),
})
}
}
impl Driving for Gated {}
impl Reacting for Gate {
type Kept = ();
async fn hear(&mut self, _said: Message) -> Result<(), Failure> {
self.open.notify_one();
Ok(())
}
async fn answer(&mut self, _asked: Message) -> Result<Answer, Failure> {
self.open.notified().await;
Ok(Answer::of("echo", ["opened"]))
}
async fn finish(self) -> Result<(), Failure> {
Ok(())
}
}
#[tokio::test]
async fn an_answer_may_wait_on_another_shells_word() {
let scripts = Scripts::of(&[(
ENTRY,
r#"
got="$(BC_INSTR GATE ask open-please)" &
sleep 0.2
BC_INSTR GATE say REC opening
wait
"#,
)]);
let gated = Gated {
open: Rc::new(Notify::new()),
};
let argv = bash(scripts.at(ENTRY));
let ran = tokio::time::timeout(
Duration::from_secs(10),
gated.run(&argv, provisioned(&gated)),
)
.await
.expect("served concurrently, or this would never return")
.unwrap();
assert_eq!(ran.subject, ExitStatus::Code(0));
assert_eq!(
ran.shells.len(),
2,
"the asker in its subshell, and the script"
);
assert!(ran.failed.is_none());
}
impl crate::Joins for Answering {
fn joining(&self, at: &Layout) -> String {
format!(
"BC_JOIN SOAK {}\n",
bash_strings::emit_scalar(at.text())
)
}
}
impl crate::Joins for Gated {
fn joining(&self, at: &Layout) -> String {
format!(
"BC_JOIN GATE {}\n",
bash_strings::emit_scalar(at.text())
)
}
}