use aion::DeployedWorkerContract;
use aion_package::{ActionBodyContract, ContentHash};
use super::declared_body::DeclaredBodyLookup;
use super::declared_body_ambiguity::DeclaringVersion;
#[must_use]
pub fn select_declared_body(
contracts: &[DeployedWorkerContract],
action: &str,
run_version: Option<&ContentHash>,
) -> DeclaredBodyLookup {
match run_version {
Some(version) => body_of_version(contracts, action, version),
None => body_across_queue(contracts, action),
}
}
fn body_of_version(
contracts: &[DeployedWorkerContract],
action: &str,
version: &ContentHash,
) -> DeclaredBodyLookup {
for deployed in contracts {
if &deployed.package_version != version {
continue;
}
for declared in &deployed.contract.actions {
if declared.name == action
&& let Some(body) = &declared.body
{
return DeclaredBodyLookup::Declared(body.clone());
}
}
}
DeclaredBodyLookup::None
}
fn body_across_queue(contracts: &[DeployedWorkerContract], action: &str) -> DeclaredBodyLookup {
let mut bodies: Vec<ActionBodyContract> = Vec::new();
let mut declaring: Vec<DeclaringVersion> = Vec::new();
for deployed in contracts {
for declared in &deployed.contract.actions {
if declared.name == action
&& let Some(body) = &declared.body
{
let index = bodies
.iter()
.position(|seen| seen == body)
.unwrap_or_else(|| {
bodies.push(body.clone());
bodies.len() - 1
});
declaring.push(DeclaringVersion {
content_hash: deployed.package_version.to_string(),
workflow_types: deployed.workflow_types.clone(),
route_active: deployed.route_active,
body: index,
});
break;
}
}
}
match bodies.len() {
0 => DeclaredBodyLookup::None,
1 => match bodies.pop() {
Some(body) => DeclaredBodyLookup::Declared(body),
None => DeclaredBodyLookup::None,
},
_ => DeclaredBodyLookup::Ambiguous { declaring },
}
}
#[cfg(test)]
mod tests {
use aion::DeployedWorkerContract;
use aion_package::{ActionBodyContract, ActionContract, ContentHash, WorkerContract};
use super::super::declared_body::DeclaredBodyLookup;
use super::select_declared_body;
fn version(byte: u8) -> ContentHash {
ContentHash::from_bytes([byte; 32])
}
fn contract(
hash: u8,
workflow_type: &str,
route_active: bool,
action: &str,
command: Option<&str>,
) -> DeployedWorkerContract {
DeployedWorkerContract {
package_version: version(hash),
contract: WorkerContract {
task_queue: "local".to_owned(),
actions: vec![ActionContract {
name: action.to_owned(),
input_schema: serde_json::json!({"type": "object"}),
output_schema: serde_json::json!({"type": "object"}),
node: None,
timeout: None,
retry: None,
advisory: false,
agent: false,
body: command.map(|command| ActionBodyContract::Run {
command: command.to_owned(),
}),
}],
},
workflow_types: vec![workflow_type.to_owned()],
route_active,
}
}
fn command_of(lookup: &DeclaredBodyLookup) -> Option<&str> {
match lookup {
DeclaredBodyLookup::Declared(ActionBodyContract::Run { command }) => Some(command),
_ => None,
}
}
#[test]
fn an_in_flight_run_keeps_its_own_versions_body_after_a_redeploy() {
let retained = [
contract(0x0a, "sweep", false, "collect", Some("echo old")),
contract(0x0b, "sweep", true, "collect", Some("echo new")),
];
let old_run = select_declared_body(&retained, "collect", Some(&version(0x0a)));
assert_eq!(
command_of(&old_run),
Some("echo old"),
"the run pinned to the superseded version must keep its own body: {old_run:?}"
);
let new_run = select_declared_body(&retained, "collect", Some(&version(0x0b)));
assert_eq!(
command_of(&new_run),
Some("echo new"),
"a run on the route-active version gets the new body: {new_run:?}"
);
let blind = select_declared_body(&retained, "collect", None);
assert!(
matches!(blind, DeclaredBodyLookup::Ambiguous { .. }),
"without the run's version the queue-wide reading must still refuse: {blind:?}"
);
}
#[test]
fn a_version_without_a_body_delegates_rather_than_borrowing_one() {
let retained = [
contract(0x0a, "sweep", false, "collect", None),
contract(0x0b, "sweep", true, "collect", Some("echo new")),
];
let lookup = select_declared_body(&retained, "collect", Some(&version(0x0a)));
assert!(
matches!(lookup, DeclaredBodyLookup::None),
"a run whose package declares no body must delegate, not borrow: {lookup:?}"
);
}
#[test]
fn a_version_absent_from_the_queue_declares_nothing_here() {
let retained = [contract(0x0b, "sweep", true, "collect", Some("echo new"))];
let lookup = select_declared_body(&retained, "collect", Some(&version(0x0c)));
assert!(
matches!(lookup, DeclaredBodyLookup::None),
"an unrelated version cannot pick up this queue's body: {lookup:?}"
);
}
#[test]
fn identical_bodies_across_versions_still_collapse_to_one() {
let retained = [
contract(0x0a, "sweep", false, "collect", Some("echo same")),
contract(0x0b, "sweep", true, "collect", Some("echo same")),
];
let lookup = select_declared_body(&retained, "collect", None);
assert_eq!(command_of(&lookup), Some("echo same"), "{lookup:?}");
}
#[test]
fn a_bodiless_action_is_none_with_or_without_the_run_version() {
let retained = [contract(0x0a, "sweep", true, "collect", None)];
assert!(matches!(
select_declared_body(&retained, "collect", None),
DeclaredBodyLookup::None
));
assert!(matches!(
select_declared_body(&retained, "collect", Some(&version(0x0a))),
DeclaredBodyLookup::None
));
}
#[test]
fn the_blind_refusal_still_names_every_declaring_version() -> Result<(), String> {
let retained = [
contract(0x0a, "sweep", false, "collect", Some("echo old")),
contract(0x0b, "sweep", true, "collect", Some("echo new")),
];
let lookup = select_declared_body(&retained, "collect", None);
let DeclaredBodyLookup::Ambiguous { declaring } = lookup else {
return Err(format!(
"different bodies must refuse when read blind: {lookup:?}"
));
};
assert_eq!(declaring.len(), 2, "{declaring:?}");
assert_eq!(declaring[0].content_hash, version(0x0a).to_string());
assert!(!declaring[0].route_active);
assert!(declaring[1].route_active);
assert_ne!(
declaring[0].body, declaring[1].body,
"the two versions must be recorded as carrying DIFFERENT bodies"
);
Ok(())
}
}