#![allow(clippy::expect_used, clippy::unwrap_used)]
mod support;
use std::error::Error;
use std::num::NonZeroU64;
use std::sync::mpsc;
use std::time::Duration;
use frame_conv::workflow::{
WorkflowCallError, WorkflowClient, WorkflowConversationId, WorkflowItem, WorkflowObservation,
WorkflowObserveError, WorkflowPhase, WorkflowProgressKind, WorkflowTerminal,
};
use serde::{Deserialize, Serialize};
use support::embedded::EmbeddedAion;
const DECIDE_AWL: &str = r"//! Decide a topic on a contributed ruling.
workflow decide
input topic: String
signal decision: Decision
outcome agreed: type Agreement, route success
outcome refused: type Refusal, route failure
type Decision { agree: Bool, note: String }
type Agreement { topic: String, note: String }
type Refusal { topic: String, note: String }
step await_decision
wait decision -> verdict
outcome yes: when verdict.agree,
route agreed(topic: topic, note: verdict.note)
outcome no: otherwise,
route refused(topic: topic, note: verdict.note)
";
#[derive(Debug, Serialize)]
struct DecideInput {
topic: String,
}
#[derive(Debug, Serialize)]
struct Decision {
agree: bool,
note: String,
}
#[derive(Debug, Deserialize, PartialEq)]
struct Agreement {
topic: String,
note: String,
}
#[derive(Debug, Deserialize, PartialEq)]
struct RoutedOutcome<P> {
outcome: String,
payload: P,
}
const OBSERVE: Duration = Duration::from_secs(12);
struct BoundedObserver {
items: mpsc::Receiver<Result<WorkflowItem, WorkflowObserveError>>,
}
impl BoundedObserver {
fn spawn(mut observation: WorkflowObservation) -> Self {
let (sender, items) = 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;
}
}
});
Self { items }
}
fn next(&self) -> Result<WorkflowItem, Box<dyn Error>> {
Ok(self.items.recv_timeout(OBSERVE)??)
}
fn collect_through_terminal(&self) -> Result<Vec<WorkflowItem>, Box<dyn Error>> {
let mut collected = Vec::new();
loop {
let item = self.next()?;
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 item_label(item: &WorkflowItem) -> String {
match item {
WorkflowItem::Progress(progress) => match &progress.kind {
WorkflowProgressKind::Opened { .. } => "opened".to_owned(),
WorkflowProgressKind::ContributionReceived { name, .. } => {
format!("contribution:{name}")
}
other => format!("{other:?}")
.split_whitespace()
.next()
.unwrap()
.to_lowercase(),
},
WorkflowItem::Terminal(terminal) => format!("terminal:{terminal:?}")
.split_whitespace()
.next()
.unwrap()
.to_lowercase(),
}
}
fn one() -> NonZeroU64 {
NonZeroU64::new(1).unwrap()
}
#[test]
fn open_contribute_observe_complete_with_contiguous_sequences() -> Result<(), Box<dyn Error>> {
let substrate = EmbeddedAion::start("workflow-complete", &[DECIDE_AWL])?;
let client = WorkflowClient::from_substrate_client(substrate.client())?;
let conversation = client.open(
"decide",
&DecideInput {
topic: "budget".to_owned(),
},
"open-key-1",
)?;
let identity = conversation.identity();
assert_eq!(
WorkflowConversationId::from_bytes(identity.conversation.to_bytes()),
identity.conversation,
"the identity round-trips through caller-persistable bytes"
);
let observer = BoundedObserver::spawn(conversation.observe_from(one()));
let first = observer.next()?;
match &first {
WorkflowItem::Progress(progress) => match &progress.kind {
WorkflowProgressKind::Opened {
workflow_kind,
run,
continued_from,
} => {
assert_eq!(workflow_kind, "decide");
assert_eq!(*run, identity.run);
assert!(continued_from.is_none());
}
other => return Err(format!("first item must be Opened, got {other:?}").into()),
},
WorkflowItem::Terminal(other) => {
return Err(format!("first item must be progress, got {other:?}").into());
}
}
conversation.contribute(
"decision",
&Decision {
agree: true,
note: "proceed".to_owned(),
},
)?;
let mut items = vec![first];
items.extend(observer.collect_through_terminal()?);
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,
"observed sequences are strictly increasing and gap-free"
);
assert!(
items.iter().any(|item| matches!(
item,
WorkflowItem::Progress(progress)
if matches!(&progress.kind, WorkflowProgressKind::ContributionReceived { name, .. } if name == "decision")
)),
"the contribution surfaces as typed progress"
);
let WorkflowItem::Terminal(WorkflowTerminal::Completed { result, .. }) =
items.last().expect("nonempty")
else {
return Err(format!("the run must complete: {items:?}").into());
};
let routed: RoutedOutcome<Agreement> = result.decode()?;
assert_eq!(
routed,
RoutedOutcome {
outcome: "agreed".to_owned(),
payload: Agreement {
topic: "budget".to_owned(),
note: "proceed".to_owned(),
},
},
"the terminal result is the outcome-tagged envelope (characterized)"
);
let replayer = BoundedObserver::spawn(client.conversation(identity).observe_from(one()));
let replay = replayer.collect_through_terminal()?;
let live_shape: Vec<(u64, String)> = items
.iter()
.map(|item| (item_seq(item), item_label(item)))
.collect();
let replay_shape: Vec<(u64, String)> = replay
.iter()
.map(|item| (item_seq(item), item_label(item)))
.collect();
assert_eq!(
live_shape, replay_shape,
"replay equals the live observation"
);
let terminal_seq = item_seq(items.last().expect("nonempty"));
let tail = BoundedObserver::spawn(
client
.conversation(identity)
.observe_from(NonZeroU64::new(terminal_seq).expect("terminal seq is nonzero")),
);
let tail_first = tail.next()?;
assert_eq!(item_seq(&tail_first), terminal_seq);
assert!(matches!(tail_first, WorkflowItem::Terminal(_)));
Ok(())
}
#[test]
fn open_idempotency_is_per_client_instance_and_conflicts_are_typed() -> Result<(), Box<dyn Error>> {
let substrate = EmbeddedAion::start("workflow-idempotency", &[DECIDE_AWL])?;
let client = WorkflowClient::from_substrate_client(substrate.client())?;
let input = DecideInput {
topic: "release".to_owned(),
};
let first = client.open("decide", &input, "shared-key")?;
let replayed = client.open("decide", &input, "shared-key")?;
assert_eq!(
first.identity(),
replayed.identity(),
"same client, same key, same request replays the same identity"
);
let conflict = client.open(
"decide",
&DecideInput {
topic: "different".to_owned(),
},
"shared-key",
);
assert!(
matches!(conflict, Err(WorkflowCallError::IdempotencyConflict { .. })),
"same key, different request: {conflict:?}"
);
let empty = client.open("decide", &input, "");
assert!(
matches!(empty, Err(WorkflowCallError::InvalidArgument { .. })),
"an empty key is refused typed: {empty:?}"
);
let second_client = WorkflowClient::from_substrate_client(substrate.client())?;
let double = second_client.open("decide", &input, "shared-key")?;
assert_ne!(
first.identity(),
double.identity(),
"a second client instance mints a SECOND conversation — the honest weaker contract"
);
Ok(())
}
#[test]
fn reopen_cancelled_run_and_typed_refusals() -> Result<(), Box<dyn Error>> {
let substrate = EmbeddedAion::start("workflow-reopen", &[DECIDE_AWL])?;
let client = WorkflowClient::from_substrate_client(substrate.client())?;
let conversation = client.open(
"decide",
&DecideInput {
topic: "merge".to_owned(),
},
"reopen-key",
)?;
let identity = conversation.identity();
let inspection = conversation.inspect()?;
assert_eq!(inspection.phase, WorkflowPhase::Running);
assert!(inspection.recorded_events >= 1);
let live_refusal = client.reopen(identity.conversation, Some(identity.run));
assert!(
matches!(live_refusal, Err(WorkflowCallError::InvalidState { .. })),
"a live run is observed, never restarted: {live_refusal:?}"
);
let raw = substrate.client();
let workflow_id =
aion_core::WorkflowId::new(uuid::Uuid::from_bytes(identity.conversation.to_bytes()));
let run_id = aion_core::RunId::new(uuid::Uuid::from_bytes(identity.run.to_bytes()));
substrate
.runtime
.block_on(raw.cancel(&workflow_id, Some(&run_id), "operator cancel"))?;
let observer = BoundedObserver::spawn(conversation.observe_from(one()));
let items = observer.collect_through_terminal()?;
let WorkflowItem::Terminal(WorkflowTerminal::Cancelled { reason, .. }) =
items.last().expect("nonempty")
else {
return Err(format!("the cancel must surface as the Cancelled terminal: {items:?}").into());
};
assert_eq!(reason, "operator cancel");
assert_eq!(conversation.inspect()?.phase, WorkflowPhase::Cancelled);
let reopened = client.reopen(identity.conversation, Some(identity.run))?;
assert_eq!(reopened.run, identity.run, "reopen continues the SAME run");
assert_eq!(reopened.phase, WorkflowPhase::Running);
assert_eq!(conversation.inspect()?.phase, WorkflowPhase::Running);
let cancelled_seq = item_seq(items.last().expect("nonempty"));
let post = BoundedObserver::spawn(
conversation.observe_from(NonZeroU64::new(cancelled_seq + 1).expect("nonzero")),
);
conversation.contribute(
"decision",
&Decision {
agree: true,
note: "after reopen".to_owned(),
},
)?;
let post_items = post.collect_through_terminal()?;
assert!(
post_items.iter().any(|item| matches!(
item,
WorkflowItem::Progress(progress)
if matches!(&progress.kind, WorkflowProgressKind::Reopened { run, .. } if *run == identity.run)
)),
"the reopen projects typed Reopened progress: {post_items:?}"
);
let WorkflowItem::Terminal(WorkflowTerminal::Completed { result, .. }) =
post_items.last().expect("nonempty")
else {
return Err(format!("the reopened run must complete: {post_items:?}").into());
};
let routed: RoutedOutcome<Agreement> = result.decode()?;
assert_eq!(routed.payload.note, "after reopen");
let absent = client.reopen(WorkflowConversationId::from_bytes([0x5A; 16]), None);
assert!(
matches!(absent, Err(WorkflowCallError::Unknown { .. })),
"an absent conversation is the typed Unknown: {absent:?}"
);
Ok(())
}