use std::collections::HashMap;
use aion_core::{ContentType, Payload};
use aion_package::{ExtractionLimits, Package};
use super::AdmissionReason;
use crate::EngineBuilder;
use crate::engine::api::Engine;
type TestResult<T = ()> = Result<T, Box<dyn std::error::Error>>;
const V1: &str = r"//! Contract-drift fixture, first deploy.
workflow contract_drift
input amount: Int
outcome completed: type Result, route success
type Result { approved: Bool }
worker payments
action charge(amount: Int) -> Result
step run
charge(amount: amount) -> result
route completed(approved: result.approved)
";
const V2: &str = r"//! Contract-drift fixture, second deploy.
workflow contract_drift
input amount: Int
input currency: String
outcome completed: type Result, route success
type Result { approved: Bool }
worker payments
action charge(amount: Int, currency: String) -> Result
step run
charge(amount: amount, currency: currency) -> result
route completed(approved: result.approved)
";
const HELD_V1: &str = r"//! Contract-drift fixture whose run parks on a signal.
workflow contract_drift_held
input amount: Int
signal go: Ruling
outcome completed: type Result, route success
type Result { approved: Bool }
type Ruling { proceed: Bool }
worker payments
action charge(amount: Int) -> Result
step hold
wait go -> decision
step run
charge(amount: amount) -> result
route completed(approved: result.approved)
";
const HELD_V2: &str = r"//! Contract-drift fixture, parked-run workflow, second deploy.
workflow contract_drift_held
input amount: Int
input currency: String
signal go: Ruling
outcome completed: type Result, route success
type Result { approved: Bool }
type Ruling { proceed: Bool }
worker payments
action charge(amount: Int, currency: String) -> Result
step hold
wait go -> decision
step run
charge(amount: amount, currency: currency) -> result
route completed(approved: result.approved)
";
async fn engine() -> TestResult<Engine> {
Ok(EngineBuilder::new()
.stop_drain_timeout(std::time::Duration::from_secs(5))
.store(aion_store::InMemoryStore::default())
.in_memory_visibility()
.build()
.await?)
}
async fn deploy(engine: &Engine, source: &str) -> TestResult<String> {
let root = tempfile::tempdir()?;
let prepared =
aion_awl_package::compile_and_assemble_awl(source, root.path(), "contract_drift.awl")?;
let package = Package::load_from_bytes(prepared.archive, ExtractionLimits::unbounded())?;
let version = package.content_hash().to_string();
engine.load_package(package).await?;
Ok(version)
}
#[tokio::test]
async fn a_superseded_version_with_no_live_run_binds_no_worker() -> TestResult {
let engine = engine().await?;
let stale = deploy(&engine, V1).await?;
let current = deploy(&engine, V2).await?;
let admission = engine.worker_contracts_for_admission("payments")?;
assert_eq!(
admission.required.len(),
1,
"only the routed version is reachable: {admission:?}"
);
assert_eq!(
admission.required[0].contract.package_version.to_string(),
current
);
assert_eq!(admission.required[0].reason, AdmissionReason::RouteActive);
assert_eq!(
admission
.unreachable
.iter()
.map(|contract| contract.package_version.to_string())
.collect::<Vec<_>>(),
vec![stale],
"the superseded version must be reported as reachable by nothing"
);
Ok(())
}
#[tokio::test]
async fn a_superseded_version_with_a_live_run_still_binds_every_worker() -> TestResult {
let engine = engine().await?;
let held = deploy(&engine, HELD_V1).await?;
let handle = engine
.start_workflow(
"contract_drift_held",
Payload::new(ContentType::Json, br#"{"amount":1}"#.to_vec()),
HashMap::new(),
"default".to_owned(),
)
.await?;
assert_eq!(handle.loaded_version().to_string(), held);
let successor = deploy(&engine, HELD_V2).await?;
assert_ne!(successor, held);
let admission = engine.worker_contracts_for_admission("payments")?;
let reasons = admission
.required
.iter()
.map(|required| {
(
required.contract.package_version.to_string(),
required.reason,
)
})
.collect::<HashMap<_, _>>();
assert_eq!(
reasons.get(&held),
Some(&AdmissionReason::LiveWorkflow),
"a superseded version a live run is pinned to still binds every worker: {admission:?}"
);
assert_eq!(
reasons.get(&successor),
Some(&AdmissionReason::RouteActive),
"the new route binds every worker: {admission:?}"
);
assert!(
admission.unreachable.is_empty(),
"one version routes, the other holds a live run: {admission:?}"
);
Ok(())
}
#[tokio::test]
async fn a_rollback_moves_the_demanded_version_with_the_route() -> TestResult {
let engine = engine().await?;
let first = deploy(&engine, V1).await?;
let second = deploy(&engine, V2).await?;
let rolled_back = first.parse()?;
engine
.route_workflow_version("contract_drift", &rolled_back)
.await?;
let admission = engine.worker_contracts_for_admission("payments")?;
assert_eq!(admission.required.len(), 1, "{admission:?}");
assert_eq!(
admission.required[0].contract.package_version.to_string(),
first,
"the rolled-back-to version is the one new starts resolve"
);
assert_eq!(
admission
.unreachable
.iter()
.map(|contract| contract.package_version.to_string())
.collect::<Vec<_>>(),
vec![second],
"the rolled-away-from version binds nobody once nothing runs on it"
);
Ok(())
}
#[tokio::test]
async fn an_undeclared_queue_demands_nothing() -> TestResult {
let engine = engine().await?;
drop(deploy(&engine, V1).await?);
let admission = engine.worker_contracts_for_admission("no_such_queue")?;
assert!(admission.required.is_empty(), "{admission:?}");
assert!(admission.unreachable.is_empty(), "{admission:?}");
Ok(())
}