use std::sync::Arc;
use aion_core::{Event, Payload, WorkflowId};
use aion_package::ContentHash;
use aion_store::EventStore;
use crate::EngineError;
use crate::loader::{ActivityServing, LoadedWorkflow, PinnedWorkflow, WorkflowCatalog};
use crate::registry::Registry;
pub(in crate::lifecycle) fn resolve_contract(
catalog: &WorkflowCatalog,
workflow_type: &str,
loaded_version: Option<&ContentHash>,
) -> Result<PinnedWorkflow, EngineError> {
let pinned = match loaded_version {
Some(version) => catalog.resolve_exact(workflow_type, version)?,
None => catalog.resolve_routed(workflow_type)?,
}
.ok_or_else(|| EngineError::WorkflowNotFound {
workflow_type: workflow_type.to_owned(),
})?;
let loaded = pinned.workflow();
let contract = loaded
.contract()
.map_err(|source| EngineError::ContractIdentity {
workflow_type: workflow_type.to_owned(),
source,
})?;
if catalog.activity_serving() == ActivityServing::QueueRouted
&& !contract.unscoped_activities.is_empty()
{
let mut activities = contract.unscoped_activities.clone();
activities.sort();
return Err(EngineError::NoQueueDeclaration {
workflow_type: workflow_type.to_owned(),
version: loaded.version().clone(),
activities: activities.join(","),
});
}
Ok(pinned)
}
pub(super) fn admit_declared_input(
loaded: &LoadedWorkflow,
input: &Payload,
) -> Result<(), EngineError> {
let workflow_type = loaded.workflow_type();
let contract = loaded
.contract()
.map_err(|source| EngineError::ContractIdentity {
workflow_type: workflow_type.to_owned(),
source,
})?;
let schema = contract.entry_input_schema(workflow_type);
if aion_package::declares_nothing(schema) {
return Ok(());
}
let refuse = |reason: String| EngineError::StartInputRefused {
workflow_type: workflow_type.to_owned(),
version: loaded.version().clone(),
reason,
};
let value = input
.to_json()
.map_err(|error| refuse(format!("the input is not decodable JSON: {error}")))?;
match aion_package::admit_value(schema, &value) {
Ok(()) => Ok(()),
Err(aion_package::AdmissionError::UnusableSchema { reason }) => {
tracing::warn!(
workflow_type,
version = %loaded.version(),
%reason,
"package declares an input schema that is not valid JSON Schema; the start input could not be admitted against it and was allowed through"
);
Ok(())
}
Err(error @ aion_package::AdmissionError::Mismatch { .. }) => {
Err(refuse(error.to_string()))
}
}
}
pub(in crate::lifecycle) async fn workflow_identity(
store: &Arc<dyn EventStore>,
registry: &Registry,
requested: Option<WorkflowId>,
) -> Result<(WorkflowId, u64), EngineError> {
let Some(workflow_id) = requested else {
return Ok((WorkflowId::new_v4(), 0));
};
if let Some(incumbent) = registry.sole_handle(&workflow_id)? {
return Err(EngineError::WorkflowIdAlreadyLive {
workflow_id: workflow_id.to_string(),
holder_run_id: incumbent.run_id().to_string(),
holder_pid: incumbent.pid(),
});
}
let initial_head = store
.read_history(&workflow_id)
.await?
.iter()
.map(Event::seq)
.max()
.unwrap_or_default();
Ok((workflow_id, initial_head))
}