#![allow(clippy::expect_used, clippy::unwrap_used)]
mod support;
use std::error::Error;
use std::fs;
use std::io::{BufRead, BufReader, Write};
use std::num::NonZeroU64;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::mpsc;
use std::time::Duration;
use aion::{ActivityDispatch, ActivityDispatcher};
use frame_conv::workflow::{
WorkflowClient, WorkflowConversationId, WorkflowIdentity, WorkflowItem, WorkflowProgressKind,
WorkflowRunId, WorkflowTerminal,
};
use serde::Deserialize;
use support::embedded::EmbeddedAion;
const TWO_STEP_AWL: &str = r"//! Two effectful steps split by a durable gate.
workflow two_steps
input job: String
signal proceed: Go
outcome done: type Done, route success
type Go { note: String }
type Ack { receipt: String }
type Done { receipt: String }
worker effects
action step_a(job: String) -> Ack
action step_b(job: String) -> Done
step first
step_a(job: job) -> a
step gate
wait proceed -> go
step second
step_b(job: job) -> b
b |> route done
";
#[derive(Debug, Deserialize, PartialEq)]
struct Done {
receipt: String,
}
#[derive(Debug, Deserialize, PartialEq)]
struct RoutedOutcome<P> {
outcome: String,
payload: P,
}
const OBSERVE: Duration = Duration::from_secs(20);
struct Ledger {
path: PathBuf,
}
impl Ledger {
fn begin(&self, dispatch: &ActivityDispatch) {
self.append(&format!(
"begin {} {} attempt={}\n",
dispatch.name, dispatch.activity_id, dispatch.attempt
));
}
fn effect_once(&self, dispatch: &ActivityDispatch) {
let line = format!("effect {} {}\n", dispatch.name, dispatch.activity_id);
let existing = fs::read_to_string(&self.path).unwrap_or_default();
if !existing.contains(&line) {
self.append(&line);
}
}
fn append(&self, line: &str) {
let mut file = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&self.path)
.expect("ledger opens");
file.write_all(line.as_bytes()).expect("ledger appends");
file.sync_all().expect("ledger reaches disk");
}
}
struct LedgerDispatcher {
ledger: Ledger,
wedge: Option<String>,
}
impl ActivityDispatcher for LedgerDispatcher {
fn dispatch(&self, request: ActivityDispatch) -> Result<String, String> {
self.ledger.begin(&request);
if self.wedge.as_deref() == Some(request.name.as_str()) {
loop {
std::thread::park();
}
}
self.ledger.effect_once(&request);
match request.name.as_str() {
"step_a" => Ok(r#"{"receipt":"a-done"}"#.to_owned()),
"step_b" => Ok(r#"{"receipt":"b-done"}"#.to_owned()),
other => Err(format!("unknown activity {other}")),
}
}
}
fn ledger_lines(dir: &Path) -> Vec<String> {
fs::read_to_string(dir.join("effects.log"))
.unwrap_or_default()
.lines()
.map(ToOwned::to_owned)
.collect()
}
fn dispatcher(dir: &Path, wedge: Option<&str>) -> LedgerDispatcher {
LedgerDispatcher {
ledger: Ledger {
path: dir.join("effects.log"),
},
wedge: wedge.map(ToOwned::to_owned),
}
}
fn run_child(window: &str, dir: &Path) -> Result<(), Box<dyn Error>> {
let wedge = (window == "mid-step").then_some("step_b");
let substrate = EmbeddedAion::start_at(dir.to_path_buf(), &[TWO_STEP_AWL], |builder| {
builder.activity_dispatcher(std::sync::Arc::new(dispatcher(dir, wedge)))
})?;
let client = WorkflowClient::from_substrate_client(substrate.client())?;
let conversation = client.open("two_steps", &serde_json::json!({"job": "j1"}), "crash-open")?;
let identity = conversation.identity();
fs::write(
dir.join("identity"),
[
identity.conversation.to_bytes().as_slice(),
identity.run.to_bytes().as_slice(),
]
.concat(),
)?;
match window {
"between-steps" => {
let mut observation = conversation.observe_from(NonZeroU64::new(1).expect("nonzero"));
loop {
if let WorkflowItem::Progress(progress) = observation.next_item()?
&& matches!(progress.kind, WorkflowProgressKind::StepCompleted { .. })
{
break;
}
}
}
"mid-step" => {
conversation.contribute("proceed", &serde_json::json!({"note": "go"}))?;
let (sender, receiver) = mpsc::channel();
let dir_owned = dir.to_path_buf();
std::thread::spawn(move || {
loop {
if ledger_lines(&dir_owned)
.iter()
.any(|line| line.starts_with("begin step_b"))
{
let _ = sender.send(());
break;
}
std::thread::yield_now();
}
});
receiver.recv_timeout(OBSERVE)?;
}
other => return Err(format!("unknown window {other}").into()),
}
println!("READY {window}");
std::io::stdout().flush()?;
loop {
std::thread::park();
}
}
fn run_and_kill_child(window: &str, dir: &Path) -> Result<(), Box<dyn Error>> {
let executable = std::env::current_exe()?;
let mut child = Command::new(executable)
.args([
"--exact",
"kill_nine_between_steps_and_mid_step",
"--nocapture",
])
.env("F3B_CRASH_WINDOW", window)
.env("F3B_CRASH_DIR", dir)
.stdout(Stdio::piped())
.spawn()?;
let stdout = child.stdout.take().expect("child stdout was piped");
let mut ready = false;
for line in BufReader::new(stdout).lines() {
if line?.contains(&format!("READY {window}")) {
ready = true;
break;
}
}
if !ready {
let _ = child.kill();
return Err(format!("{window}: child exited before the READY witness").into());
}
child.kill()?;
let status = child.wait()?;
if status.success() {
return Err(format!("{window}: SIGKILLed child reported success").into());
}
Ok(())
}
fn held_identity(dir: &Path) -> Result<WorkflowIdentity, Box<dyn Error>> {
let bytes = fs::read(dir.join("identity"))?;
if bytes.len() != 32 {
return Err(format!("identity file holds {} bytes, want 32", bytes.len()).into());
}
let mut conversation = [0_u8; 16];
let mut run = [0_u8; 16];
conversation.copy_from_slice(&bytes[..16]);
run.copy_from_slice(&bytes[16..]);
Ok(WorkflowIdentity {
conversation: WorkflowConversationId::from_bytes(conversation),
run: WorkflowRunId::from_bytes(run),
})
}
fn collect_through_terminal(
client: &WorkflowClient,
identity: WorkflowIdentity,
) -> Result<Vec<WorkflowItem>, Box<dyn Error>> {
let mut observation = client
.conversation(identity)
.observe_from(NonZeroU64::new(1).expect("nonzero"));
let (sender, receiver) = mpsc::channel();
std::thread::spawn(move || {
loop {
let item = observation.next_item();
let last = matches!(&item, Ok(WorkflowItem::Terminal(_)) | Err(_));
if sender.send(item).is_err() || last {
break;
}
}
});
let mut collected = Vec::new();
loop {
let item = receiver.recv_timeout(OBSERVE)??;
let done = matches!(item, WorkflowItem::Terminal(_));
collected.push(item);
if done {
return Ok(collected);
}
}
}
fn item_seq(item: &WorkflowItem) -> u64 {
match item {
WorkflowItem::Progress(progress) => progress.seq,
WorkflowItem::Terminal(
WorkflowTerminal::Completed { seq, .. }
| WorkflowTerminal::Failed { seq, .. }
| WorkflowTerminal::Cancelled { seq, .. }
| WorkflowTerminal::TimedOut { seq, .. }
| WorkflowTerminal::ContinuedAsNew { seq, .. },
) => *seq,
}
}
fn assert_no_foreign_durable_state(dir: &Path) -> Result<(), Box<dyn Error>> {
for entry in fs::read_dir(dir)? {
let name = entry?.file_name().to_string_lossy().into_owned();
let sanctioned = name.starts_with("aion.db")
|| name.starts_with("workflow-")
|| name == "effects.log"
|| name == "identity";
if !sanctioned {
return Err(format!("unsanctioned durable file in the store dir: {name}").into());
}
}
Ok(())
}
#[test]
fn kill_nine_between_steps_and_mid_step() -> Result<(), Box<dyn Error>> {
if let Ok(window) = std::env::var("F3B_CRASH_WINDOW") {
let dir = PathBuf::from(std::env::var("F3B_CRASH_DIR")?);
return run_child(&window, &dir);
}
for window in ["between-steps", "mid-step"] {
let dir = support::store_dir(&format!("workflow-crash-{window}"))?;
run_and_kill_child(window, &dir)?;
let killed_lines = ledger_lines(&dir);
assert!(
killed_lines
.iter()
.any(|line| line.starts_with("effect step_a")),
"{window}: step A's effect committed before the kill: {killed_lines:?}"
);
let substrate = EmbeddedAion::start_at(dir.clone(), &[TWO_STEP_AWL], |builder| {
builder.activity_dispatcher(std::sync::Arc::new(dispatcher(&dir, None)))
})?;
let client = WorkflowClient::from_substrate_client(substrate.client())?;
let identity = held_identity(&dir)?;
if window == "between-steps" {
client
.conversation(identity)
.contribute("proceed", &serde_json::json!({"note": "resume"}))?;
}
let items = collect_through_terminal(&client, identity)?;
let seqs: Vec<u64> = items.iter().map(item_seq).collect();
let expected: Vec<u64> = (seqs[0]..seqs[0] + seqs.len() as u64).collect();
assert_eq!(seqs, expected, "{window}: gap-free, duplicate-free");
let step_completions = items
.iter()
.filter(|item| {
matches!(
item,
WorkflowItem::Progress(progress)
if matches!(progress.kind, WorkflowProgressKind::StepCompleted { .. })
)
})
.count();
assert_eq!(
step_completions, 2,
"{window}: exactly one committed completion per step — replay, not re-execution"
);
let WorkflowItem::Terminal(WorkflowTerminal::Completed { result, .. }) =
items.last().expect("nonempty")
else {
return Err(format!("{window}: the run must complete: {items:?}").into());
};
let routed: RoutedOutcome<Done> = result.decode()?;
assert_eq!(routed.outcome, "done");
assert_eq!(routed.payload.receipt, "b-done");
let final_lines = ledger_lines(&dir);
let effect_a = final_lines
.iter()
.filter(|line| line.starts_with("effect step_a"))
.count();
let effect_b = final_lines
.iter()
.filter(|line| line.starts_with("effect step_b"))
.count();
let begin_b = final_lines
.iter()
.filter(|line| line.starts_with("begin step_b"))
.count();
assert_eq!(
effect_a, 1,
"{window}: committed step A occurred exactly once across both processes"
);
assert_eq!(
effect_b, 1,
"{window}: step B's effect committed exactly once"
);
match window {
"mid-step" => assert!(
begin_b >= 2,
"mid-step: the wedged begin plus the re-execution: {final_lines:?}"
),
_ => assert_eq!(
begin_b, 1,
"between-steps: step B was never admitted before the kill"
),
}
assert_no_foreign_durable_state(&dir)?;
}
Ok(())
}
#[test]
fn completion_and_last_disconnect_race_leaves_terminal_observable() -> Result<(), Box<dyn Error>> {
let dir = support::store_dir("workflow-race")?;
let substrate = EmbeddedAion::start_at(dir.clone(), &[TWO_STEP_AWL], |builder| {
builder.activity_dispatcher(std::sync::Arc::new(dispatcher(&dir, None)))
})?;
let client = WorkflowClient::from_substrate_client(substrate.client())?;
let conversation = client.open(
"two_steps",
&serde_json::json!({"job": "race"}),
"race-open",
)?;
let identity = conversation.identity();
let racer = client
.conversation(identity)
.observe_from(NonZeroU64::new(1).expect("nonzero"));
conversation.contribute("proceed", &serde_json::json!({"note": "go"}))?;
drop(racer);
let items = collect_through_terminal(&client, identity)?;
let WorkflowItem::Terminal(WorkflowTerminal::Completed { .. }) =
items.last().expect("nonempty")
else {
return Err(format!("the terminal survives the race: {items:?}").into());
};
Ok(())
}