use af_context::{InstanceId, RunId, SubjectId, TenantId};
use std::collections::{BTreeMap, BTreeSet};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkflowDefinition {
pub id: String,
pub name: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct WorkflowInstance {
pub id: String,
pub tenant_id: TenantId,
pub subject_id: SubjectId,
pub definition_id: String,
pub revision: u64,
pub execution_profile_id: String,
pub execution_profile_revision: u64,
pub lifecycle: LifecyclePolicy,
pub status: String,
pub state_version: i64,
pub event_sequence: i64,
pub control_epoch: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TriggerBinding {
pub id: String,
pub revision: u64,
pub source: String,
pub event_type: String,
pub instance_id: InstanceId,
pub predicate: Value,
pub ordering: OrderingPolicy,
pub starts_at: Option<DateTime<Utc>>,
pub expires_at: Option<DateTime<Utc>>,
#[serde(default = "default_gap_wait_ms")]
pub gap_wait_ms: u64,
#[serde(default = "default_gap_limit")]
pub gap_limit: u32,
}
impl TriggerBinding {
pub fn validate(&self) -> Result<(), ContractError> {
for (name, value) in [
("trigger binding id", self.id.as_str()),
("trigger source", self.source.as_str()),
("trigger event_type", self.event_type.as_str()),
("trigger instance_id", self.instance_id.as_str()),
] {
required(name, value)?;
}
if self.gap_wait_ms == 0 || self.gap_wait_ms > 3_600_000 {
return Err(ContractError::Invalid(
"trigger gap_wait_ms must be between 1 and 3600000".into(),
));
}
if self.gap_limit == 0 || self.gap_limit > 10_000 {
return Err(ContractError::Invalid(
"trigger gap_limit must be between 1 and 10000".into(),
));
}
if !self.predicate.is_object() {
return Err(ContractError::Invalid(
"trigger predicate must be a JSON object".into(),
));
}
if matches!((self.starts_at, self.expires_at), (Some(start), Some(end)) if end <= start) {
return Err(ContractError::Invalid(
"trigger expires_at must be after starts_at".into(),
));
}
Ok(())
}
}
const fn default_gap_wait_ms() -> u64 {
30_000
}
const fn default_gap_limit() -> u32 {
100
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct WorkflowEvent {
pub instance_id: InstanceId,
pub sequence: i64,
pub event_type: String,
pub payload: Value,
pub content_digest: String,
pub occurred_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionMode {
Simulation,
Backtest,
Paper,
Live,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DurabilityGrade {
Standard,
FundsGrade,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ExecutionProfileRevision {
pub id: String,
pub revision: u64,
pub content_digest: String,
pub mode: ExecutionMode,
pub durability_grade: DurabilityGrade,
pub trigger_provider: String,
pub data_provider: String,
pub clock_model: String,
pub action_provider: String,
#[serde(default)]
pub models: BTreeMap<String, Value>,
#[serde(default)]
pub environment: Value,
#[serde(default)]
pub policy_bundle: Value,
#[serde(default)]
pub connection_bindings: BTreeMap<String, String>,
}
impl ExecutionProfileRevision {
pub fn validate(&self) -> Result<(), ContractError> {
for (name, value) in [
("execution profile id", self.id.as_str()),
(
"execution profile content_digest",
self.content_digest.as_str(),
),
("trigger_provider", self.trigger_provider.as_str()),
("data_provider", self.data_provider.as_str()),
("clock_model", self.clock_model.as_str()),
("action_provider", self.action_provider.as_str()),
] {
required(name, value)?;
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub struct CapabilityPin {
pub id: String,
pub contract_version: String,
pub content_digest: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowRevision {
pub definition_id: String,
pub revision: u64,
pub content_digest: String,
pub kernel_abi_version: String,
pub dependency_set_digest: String,
#[serde(default)]
pub expression_versions: BTreeMap<String, String>,
#[serde(default)]
pub capabilities: Vec<CapabilityPin>,
#[serde(default)]
pub template_provenance: Value,
pub spec: crate::Spec,
}
impl WorkflowRevision {
pub fn validate(&self) -> Result<(), ContractError> {
required("definition_id", &self.definition_id)?;
required("content_digest", &self.content_digest)?;
required("kernel_abi_version", &self.kernel_abi_version)?;
required("dependency_set_digest", &self.dependency_set_digest)?;
self.spec
.validate_structure()
.map_err(|error| ContractError::Invalid(error.to_string()))?;
let unique = self
.capabilities
.iter()
.map(|pin| (&pin.id, &pin.contract_version))
.collect::<BTreeSet<_>>();
if unique.len() != self.capabilities.len() {
return Err(ContractError::Invalid(
"capability pins must be unique by id and contract version".into(),
));
}
for pin in &self.capabilities {
required("capability id", &pin.id)?;
required("capability contract_version", &pin.contract_version)?;
required("capability content_digest", &pin.content_digest)?;
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CapabilityKind {
Trigger,
Expression,
Guard,
Action,
Sink,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Effect {
Pure,
Read,
InternalWrite,
ExternalWrite,
Funds,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IdempotencyMode {
None,
Native,
ReconcileBeforeRetry,
NeverAutomaticRetry,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CapabilityLifecycle {
Installed,
Active,
Deprecated,
Disabled,
Unavailable,
EmergencyRevoked,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RetryPolicy {
pub max_attempts: u32,
pub timeout_ms: u64,
pub initial_backoff_ms: u64,
pub max_backoff_ms: u64,
}
impl RetryPolicy {
pub fn backoff_ms(&self, attempt: u32, jitter_seed: u64) -> u64 {
let factor = 1_u64
.checked_shl(attempt.saturating_sub(1).min(20))
.unwrap_or(u64::MAX);
let base = self
.initial_backoff_ms
.saturating_mul(factor)
.min(self.max_backoff_ms);
let jitter_ceiling = (base / 4).max(1);
base.saturating_add(jitter_seed % jitter_ceiling)
.min(self.max_backoff_ms)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CapabilityManifest {
pub id: String,
pub contract_version: String,
pub content_digest: String,
pub kind: CapabilityKind,
pub input_schema: Value,
pub output_schema: Value,
pub effect: Effect,
pub deterministic: bool,
pub idempotency_mode: IdempotencyMode,
pub retry: RetryPolicy,
#[serde(default)]
pub permissions: BTreeSet<String>,
#[serde(default)]
pub required_guards: BTreeSet<GuardKind>,
#[serde(default)]
pub taint_rules: Value,
#[serde(default)]
pub resource_cost: Value,
pub supports_simulation: bool,
pub supports_replay: bool,
pub supports_paper: bool,
pub supports_live: bool,
#[serde(default)]
pub clock_requirements: Value,
#[serde(default)]
pub data_requirements: Value,
pub supports_reconciliation: bool,
pub lifecycle: CapabilityLifecycle,
}
impl CapabilityManifest {
pub fn action(
id: impl Into<String>,
contract_version: impl Into<String>,
content_digest: impl Into<String>,
effect: Effect,
idempotency_mode: IdempotencyMode,
supports_reconciliation: bool,
) -> Self {
let required_guards = match effect {
Effect::Funds => BTreeSet::from([
GuardKind::Authorization,
GuardKind::Freshness,
GuardKind::Reservation,
]),
Effect::ExternalWrite => BTreeSet::from([GuardKind::Authorization]),
_ => BTreeSet::new(),
};
Self {
id: id.into(),
contract_version: contract_version.into(),
content_digest: content_digest.into(),
kind: CapabilityKind::Action,
input_schema: serde_json::json!({"type": "object"}),
output_schema: serde_json::json!({"type": "object"}),
effect,
deterministic: false,
idempotency_mode,
retry: RetryPolicy {
max_attempts: 1,
timeout_ms: 30_000,
initial_backoff_ms: 100,
max_backoff_ms: 5_000,
},
permissions: BTreeSet::new(),
required_guards,
taint_rules: Value::Null,
resource_cost: Value::Null,
supports_simulation: true,
supports_replay: false,
supports_paper: true,
supports_live: true,
clock_requirements: Value::Null,
data_requirements: Value::Null,
supports_reconciliation,
lifecycle: CapabilityLifecycle::Active,
}
}
pub fn validate(&self) -> Result<(), ContractError> {
required("capability id", &self.id)?;
required("contract_version", &self.contract_version)?;
required("content_digest", &self.content_digest)?;
if self.retry.max_attempts == 0
|| self.retry.timeout_ms == 0
|| self.retry.initial_backoff_ms == 0
|| self.retry.max_backoff_ms < self.retry.initial_backoff_ms
{
return Err(ContractError::Invalid(
"retry attempts, timeout and backoff bounds are invalid".into(),
));
}
if self.effect == Effect::Funds {
if self.kind != CapabilityKind::Action {
return Err(ContractError::Invalid(
"funds effect is only valid for action capabilities".into(),
));
}
if !self.supports_reconciliation && self.idempotency_mode != IdempotencyMode::Native {
return Err(ContractError::Invalid(
"funds actions require native idempotency or reconciliation".into(),
));
}
for guard in [
GuardKind::Authorization,
GuardKind::Freshness,
GuardKind::Reservation,
] {
if !self.required_guards.contains(&guard) {
return Err(ContractError::Invalid(format!(
"funds action must require {guard:?} guard"
)));
}
}
}
Ok(())
}
pub fn can_start_new_work(&self) -> bool {
matches!(
self.lifecycle,
CapabilityLifecycle::Installed
| CapabilityLifecycle::Active
| CapabilityLifecycle::Deprecated
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum GuardKind {
Authorization,
Freshness,
Reservation,
Policy,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OrderingPolicy {
StrictSequence,
Commutative,
LatestStateReconcile,
RejectUnordered,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TriggerEnvelope {
pub event_id: String,
pub event_type: String,
pub source: String,
pub schema_version: String,
pub tenant_id: TenantId,
pub subject_id: SubjectId,
pub aggregate_id: String,
pub source_sequence: Option<i64>,
pub observed_version: Option<String>,
pub occurred_at: DateTime<Utc>,
pub received_at: DateTime<Utc>,
pub watermark: Option<DateTime<Utc>>,
pub correlation_key: String,
pub dedup_key: String,
pub cursor: Option<String>,
pub payload: Value,
#[serde(default)]
pub trace_context: Value,
}
impl TriggerEnvelope {
pub fn validate(&self) -> Result<(), ContractError> {
for (name, value) in [
("event_id", self.event_id.as_str()),
("event_type", self.event_type.as_str()),
("source", self.source.as_str()),
("schema_version", self.schema_version.as_str()),
("tenant_id", self.tenant_id.as_str()),
("aggregate_id", self.aggregate_id.as_str()),
("correlation_key", self.correlation_key.as_str()),
("dedup_key", self.dedup_key.as_str()),
] {
required(name, value)?;
}
if self.occurred_at > self.received_at + chrono::Duration::minutes(5) {
return Err(ContractError::Invalid(
"occurred_at is implausibly ahead of received_at".into(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CompletionPolicy {
ExplicitStop,
FirstTrigger,
FirstMatch,
FirstActionTerminal,
FirstSuccess,
AfterMatchedEvaluations(u64),
AfterSuccessfulRuns(u64),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExpiryPolicy {
Drain,
Cancel,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ScheduleCadence {
Cron {
expression: String,
},
FixedRate {
milliseconds: u64,
},
FixedDelay {
milliseconds: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CatchUpPolicy {
Skip,
CatchUpOnce,
CatchUpAll {
limit: u32,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SchedulePolicy {
pub cadence: ScheduleCadence,
pub timezone: String,
pub catch_up: CatchUpPolicy,
}
impl SchedulePolicy {
pub fn validate(&self) -> Result<(), ContractError> {
self.timezone.parse::<chrono_tz::Tz>().map_err(|_| {
ContractError::Invalid(format!("unknown IANA timezone '{}'", self.timezone))
})?;
match &self.cadence {
ScheduleCadence::Cron { expression } => {
expression.parse::<cron::Schedule>().map_err(|error| {
ContractError::Invalid(format!("invalid cron '{expression}': {error}"))
})?;
}
ScheduleCadence::FixedRate { milliseconds }
| ScheduleCadence::FixedDelay { milliseconds }
if *milliseconds == 0 =>
{
return Err(ContractError::Invalid(
"schedule interval must be positive".into(),
));
}
_ => {}
}
if matches!(self.catch_up, CatchUpPolicy::CatchUpAll { limit: 0 }) {
return Err(ContractError::Invalid(
"catch_up_all limit must be positive".into(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct LifecyclePolicy {
pub starts_at: Option<DateTime<Utc>>,
pub expires_at: Option<DateTime<Utc>>,
pub completion: CompletionPolicy,
pub event_idle_timeout_ms: Option<u64>,
pub progress_timeout_ms: Option<u64>,
pub on_expiry: ExpiryPolicy,
pub drain_deadline: Option<DateTime<Utc>>,
}
impl LifecyclePolicy {
pub fn run_once() -> Self {
Self {
starts_at: None,
expires_at: None,
completion: CompletionPolicy::FirstSuccess,
event_idle_timeout_ms: None,
progress_timeout_ms: None,
on_expiry: ExpiryPolicy::Drain,
drain_deadline: None,
}
}
pub fn validate(&self) -> Result<(), ContractError> {
if self
.starts_at
.zip(self.expires_at)
.is_some_and(|(a, b)| a >= b)
{
return Err(ContractError::Invalid(
"lifecycle starts_at must be before expires_at".into(),
));
}
if self.drain_deadline.is_some() && self.on_expiry != ExpiryPolicy::Drain {
return Err(ContractError::Invalid(
"drain_deadline requires on_expiry=drain".into(),
));
}
if self
.expires_at
.zip(self.drain_deadline)
.is_some_and(|(expires, drain)| drain <= expires)
{
return Err(ContractError::Invalid(
"lifecycle drain_deadline must be after expires_at".into(),
));
}
if self.event_idle_timeout_ms == Some(0) || self.progress_timeout_ms == Some(0) {
return Err(ContractError::Invalid(
"lifecycle timeouts must be positive".into(),
));
}
if matches!(
self.completion,
CompletionPolicy::AfterMatchedEvaluations(0) | CompletionPolicy::AfterSuccessfulRuns(0)
) {
return Err(ContractError::Invalid(
"lifecycle completion count must be positive".into(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum StepOutcomeKind {
Succeeded,
Skipped,
Waiting,
Failed,
Cancelled,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StepOutcome {
pub kind: StepOutcomeKind,
pub reason_code: Option<String>,
#[serde(default)]
pub wake_condition: Value,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ActionState {
Prepared,
AwaitingConfirmation,
Authorized,
DispatchCommitted,
Executing,
Succeeded,
Rejected,
Retryable,
Unknown,
Reconciled,
}
impl ActionState {
pub fn can_transition_to(self, next: Self) -> bool {
use ActionState::*;
matches!(
(self, next),
(Prepared, AwaitingConfirmation | Authorized | Rejected)
| (AwaitingConfirmation, Authorized | Rejected)
| (Authorized, DispatchCommitted | Rejected)
| (
DispatchCommitted,
Executing | Succeeded | Rejected | Retryable | Unknown
)
| (Executing, Succeeded | Rejected | Retryable | Unknown)
| (Retryable, DispatchCommitted | Reconciled)
| (Unknown, Reconciled)
| (Succeeded | Rejected, Reconciled)
)
}
pub fn dispatch_committed(self) -> bool {
matches!(
self,
Self::DispatchCommitted
| Self::Executing
| Self::Succeeded
| Self::Rejected
| Self::Retryable
| Self::Unknown
| Self::Reconciled
)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ControlEpochs {
pub tenant: i64,
pub resource: i64,
pub instance: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ControlScope {
Tenant,
Resource,
Instance,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ControlMode {
Running,
Paused,
Stopped,
EmergencyStopped,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ControlCommand {
pub scope: ControlScope,
pub scope_id: String,
pub mode: ControlMode,
pub operator_subject_id: SubjectId,
pub reason: String,
pub override_expires_at: Option<DateTime<Utc>>,
}
impl ControlCommand {
pub fn validate(&self) -> Result<(), ContractError> {
required("control scope_id", &self.scope_id)?;
required("control operator_subject_id", &self.operator_subject_id)?;
required("control reason", &self.reason)?;
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ResourceReservationRef {
pub reservation_id: String,
pub fencing_token: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ActionIntent {
pub id: String,
pub tenant_id: TenantId,
pub instance_id: InstanceId,
pub run_id: RunId,
pub capability: CapabilityPin,
pub idempotency_key: String,
pub state: ActionState,
pub input: Value,
pub effect: Effect,
pub retry_class: IdempotencyMode,
pub control_epochs: ControlEpochs,
pub resource_scope_id: String,
pub lease_epoch: i64,
pub action_epoch: i64,
pub deadline: Option<DateTime<Utc>>,
pub reservation: Option<ResourceReservationRef>,
pub created_at: DateTime<Utc>,
}
impl ActionIntent {
pub fn validate(&self) -> Result<(), ContractError> {
required("action id", &self.id)?;
required("action idempotency_key", &self.idempotency_key)?;
if self.effect == Effect::Funds && self.reservation.is_none() {
return Err(ContractError::Invalid(
"funds action requires a resource reservation".into(),
));
}
if self.effect == Effect::Funds && self.resource_scope_id.trim().is_empty() {
return Err(ContractError::Invalid(
"funds action requires a resource scope".into(),
));
}
if self.effect == Effect::Funds && self.retry_class == IdempotencyMode::None {
return Err(ContractError::Invalid(
"funds action requires an explicit retry class".into(),
));
}
if self
.deadline
.is_some_and(|deadline| deadline <= self.created_at)
{
return Err(ContractError::Invalid(
"action deadline must be after creation".into(),
));
}
Ok(())
}
pub fn validate_prepared(&self) -> Result<(), ContractError> {
self.validate()?;
if self.state != ActionState::Prepared {
return Err(ContractError::Invalid(
"new action intent must start in prepared state".into(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ActionReceipt {
pub id: String,
pub action_intent_id: String,
pub provider_version: String,
pub received_at: DateTime<Utc>,
pub outcome: ActionState,
pub payload: Value,
pub raw_receipt_digest: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ActionObservation {
pub id: String,
pub action_intent_id: String,
pub provider_version: String,
pub observed_at: DateTime<Utc>,
pub state: String,
pub resource_ref: Option<Value>,
pub raw_receipt_digest: String,
pub terminal: bool,
#[serde(default)]
pub retry_authorized: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SourceGap {
pub tenant_id: TenantId,
pub source: String,
pub aggregate_id: String,
pub missing_from: i64,
pub missing_to: i64,
pub deadline: DateTime<Utc>,
pub status: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ReservationScope {
Tenant,
Resource,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ReservationState {
Reserved,
Consumed,
Released,
Expired,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ResourceReservation {
pub id: String,
pub tenant_id: TenantId,
pub scope: ReservationScope,
pub scope_id: String,
pub resource_kind: String,
pub amount: String,
pub policy_version: String,
pub expires_at: DateTime<Utc>,
pub fencing_token: i64,
pub provider_reservation_id: Option<String>,
pub state: ReservationState,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ObservedValue<T> {
pub value: T,
pub source: String,
pub observed_at: DateTime<Utc>,
pub received_at: DateTime<Utc>,
pub source_version: String,
pub quality: String,
pub digest: String,
}
impl<T> ObservedValue<T> {
pub fn is_accepted(
&self,
now: DateTime<Utc>,
max_age: chrono::Duration,
accepted_quality: &BTreeSet<String>,
) -> bool {
self.observed_at <= now
&& now - self.observed_at <= max_age
&& accepted_quality.contains(&self.quality)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DecisionSnapshot {
pub input_digests: Vec<String>,
pub instance_config_version: String,
pub policy_version: String,
pub capability_versions: Vec<CapabilityPin>,
pub execution_profile_revision_id: String,
#[serde(default)]
pub context_snapshot: Value,
#[serde(default)]
pub policy_snapshot: Value,
#[serde(default)]
pub agent_provenance: Value,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ParameterProvenance {
UserSupplied,
ProductDefault,
AgentInferred,
Derived,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ParameterValue {
pub value: Value,
pub provenance: ParameterProvenance,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ParameterSpec {
pub name: String,
pub required: bool,
pub required_explicit_for_live: bool,
}
pub fn validate_parameters(
mode: ExecutionMode,
specs: &[ParameterSpec],
values: &BTreeMap<String, ParameterValue>,
) -> Result<(), MissingRequirements> {
let missing = specs
.iter()
.filter(|spec| match values.get(&spec.name) {
None => {
spec.required || (mode == ExecutionMode::Live && spec.required_explicit_for_live)
}
Some(value) => {
mode == ExecutionMode::Live
&& spec.required_explicit_for_live
&& value.provenance != ParameterProvenance::UserSupplied
}
})
.map(|spec| spec.name.clone())
.collect::<Vec<_>>();
if missing.is_empty() {
Ok(())
} else {
Err(MissingRequirements {
parameters: missing,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, thiserror::Error)]
#[error("missing explicit workflow requirements: {parameters:?}")]
pub struct MissingRequirements {
pub parameters: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RootOperationBudget {
pub root_operation_id: String,
pub max_depth: u32,
pub descendant_limit: u32,
pub run_limit: u32,
pub token_budget: u64,
pub cost_budget_micros: u64,
pub action_budget: u32,
pub deadline: DateTime<Utc>,
#[serde(default)]
pub permission_ceiling: BTreeSet<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DecisionArtifact {
pub id: String,
pub root_operation_id: String,
pub content_digest: String,
pub model: String,
pub prompt_digest: String,
pub output: Value,
pub created_at: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DiagnosticRecord {
pub code: String,
pub message: String,
pub at: DateTime<Utc>,
#[serde(default)]
pub detail: Value,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ExecutionProjection {
pub schema_version: String,
pub source_sequence: i64,
pub current_steps: Vec<String>,
pub completed_steps: Vec<String>,
pub waiting_on: Option<Value>,
pub last_decision: Option<Value>,
pub next_possible_steps: Vec<String>,
pub next_trigger_at: Option<DateTime<Utc>>,
pub planned_actions: Vec<String>,
pub latest_diagnostic: Option<DiagnosticRecord>,
pub progress: Value,
pub expires_at: Option<DateTime<Utc>>,
}
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
pub enum ContractError {
#[error("invalid workflow contract: {0}")]
Invalid(String),
}
fn required(name: &str, value: &str) -> Result<(), ContractError> {
if value.trim().is_empty() {
Err(ContractError::Invalid(format!("{name} is required")))
} else {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn empty_spec() -> crate::Spec {
crate::Spec {
spec_id: "test".into(),
version: "1".into(),
description: String::new(),
aliases: vec![],
instance_config_schema: None,
display: BTreeMap::new(),
branches: vec![],
}
}
fn funds_manifest() -> CapabilityManifest {
CapabilityManifest::action(
"action.example",
"1",
"sha256:x",
Effect::Funds,
IdempotencyMode::ReconcileBeforeRetry,
true,
)
}
#[test]
fn funds_capability_requires_reconciliation_and_all_guards() {
assert!(funds_manifest().validate().is_ok());
let mut invalid = funds_manifest();
invalid.required_guards.remove(&GuardKind::Freshness);
assert!(invalid.validate().is_err());
invalid.required_guards.insert(GuardKind::Freshness);
invalid.supports_reconciliation = false;
assert!(invalid.validate().is_err());
}
#[test]
fn action_state_never_skips_dispatch_commit() {
assert!(ActionState::Authorized.can_transition_to(ActionState::DispatchCommitted));
assert!(!ActionState::Authorized.can_transition_to(ActionState::Succeeded));
assert!(ActionState::Unknown.can_transition_to(ActionState::Reconciled));
}
#[test]
fn live_parameters_cannot_be_silently_inferred() {
let specs = [ParameterSpec {
name: "account".into(),
required: true,
required_explicit_for_live: true,
}];
let inferred = BTreeMap::from([(
"account".into(),
ParameterValue {
value: Value::String("a".into()),
provenance: ParameterProvenance::AgentInferred,
},
)]);
assert_eq!(
validate_parameters(ExecutionMode::Live, &specs, &inferred)
.unwrap_err()
.parameters,
["account"]
);
assert!(validate_parameters(ExecutionMode::Paper, &specs, &inferred).is_ok());
}
#[test]
fn freshness_is_fail_closed() {
let now = Utc::now();
let observed = ObservedValue {
value: 1,
source: "source".into(),
observed_at: now - chrono::Duration::seconds(2),
received_at: now,
source_version: "1".into(),
quality: "good".into(),
digest: "d".into(),
};
assert!(observed.is_accepted(
now,
chrono::Duration::seconds(3),
&BTreeSet::from(["good".into()])
));
assert!(!observed.is_accepted(
now,
chrono::Duration::seconds(1),
&BTreeSet::from(["good".into()])
));
}
#[test]
fn revision_and_manifest_contracts_fail_closed() {
let mut revision = WorkflowRevision {
definition_id: "definition".into(),
revision: 1,
content_digest: "digest".into(),
kernel_abi_version: "1".into(),
dependency_set_digest: "dependencies".into(),
expression_versions: BTreeMap::new(),
capabilities: vec![CapabilityPin {
id: "action.example".into(),
contract_version: "1".into(),
content_digest: "capability-digest".into(),
}],
template_provenance: Value::Null,
spec: empty_spec(),
};
assert!(revision.validate().is_ok());
revision.capabilities.push(revision.capabilities[0].clone());
assert!(revision.validate().is_err());
revision.capabilities.pop();
revision.kernel_abi_version.clear();
assert!(revision.validate().is_err());
let external = CapabilityManifest::action(
"action.notify",
"1",
"digest",
Effect::ExternalWrite,
IdempotencyMode::NeverAutomaticRetry,
false,
);
assert_eq!(
external.required_guards,
BTreeSet::from([GuardKind::Authorization])
);
assert!(external.validate().is_ok());
let mut invalid = external;
invalid.retry.max_attempts = 0;
assert!(invalid.validate().is_err());
let mut wrong_kind = funds_manifest();
wrong_kind.kind = CapabilityKind::Expression;
assert!(wrong_kind.validate().is_err());
wrong_kind.kind = CapabilityKind::Action;
wrong_kind.idempotency_mode = IdempotencyMode::Native;
wrong_kind.supports_reconciliation = false;
assert!(wrong_kind.validate().is_ok());
wrong_kind.lifecycle = CapabilityLifecycle::Disabled;
assert!(!wrong_kind.can_start_new_work());
}
#[test]
fn trigger_schedule_lifecycle_and_control_validate_boundaries() {
let now = Utc::now();
let mut trigger = TriggerEnvelope {
event_id: "event".into(),
event_type: "example".into(),
source: "source".into(),
schema_version: "1".into(),
tenant_id: "tenant".parse().unwrap(),
subject_id: "subject".parse().unwrap(),
aggregate_id: "aggregate".into(),
source_sequence: Some(1),
observed_version: None,
occurred_at: now,
received_at: now,
watermark: None,
correlation_key: "key".into(),
dedup_key: "dedup".into(),
cursor: None,
payload: Value::Null,
trace_context: Value::Null,
};
assert!(trigger.validate().is_ok());
trigger.occurred_at = now + chrono::Duration::minutes(6);
assert!(trigger.validate().is_err());
let valid_schedule = SchedulePolicy {
cadence: ScheduleCadence::Cron {
expression: "0 0 * * * *".into(),
},
timezone: "Asia/Shanghai".into(),
catch_up: CatchUpPolicy::CatchUpOnce,
};
assert!(valid_schedule.validate().is_ok());
for invalid in [
SchedulePolicy {
timezone: "Nowhere/Invalid".into(),
..valid_schedule.clone()
},
SchedulePolicy {
cadence: ScheduleCadence::FixedRate { milliseconds: 0 },
..valid_schedule.clone()
},
SchedulePolicy {
catch_up: CatchUpPolicy::CatchUpAll { limit: 0 },
..valid_schedule
},
] {
assert!(invalid.validate().is_err());
}
let mut lifecycle = LifecyclePolicy::run_once();
assert!(lifecycle.validate().is_ok());
lifecycle.starts_at = Some(now);
lifecycle.expires_at = Some(now);
assert!(lifecycle.validate().is_err());
lifecycle.starts_at = None;
lifecycle.expires_at = None;
lifecycle.on_expiry = ExpiryPolicy::Cancel;
lifecycle.drain_deadline = Some(now);
assert!(lifecycle.validate().is_err());
let mut command = ControlCommand {
scope: ControlScope::Instance,
scope_id: "instance".into(),
mode: ControlMode::Paused,
operator_subject_id: "operator".parse().unwrap(),
reason: "maintenance".into(),
override_expires_at: None,
};
assert!(command.validate().is_ok());
command.reason.clear();
assert!(command.validate().is_err());
}
#[test]
fn funds_intent_requires_reservation_scope_and_retry_class() {
let now = Utc::now();
let mut intent = ActionIntent {
id: "intent".into(),
tenant_id: "tenant".parse().unwrap(),
instance_id: "instance".parse().unwrap(),
run_id: "run".parse().unwrap(),
capability: CapabilityPin {
id: "action.example".into(),
contract_version: "1".into(),
content_digest: "digest".into(),
},
idempotency_key: "idempotency".into(),
state: ActionState::Prepared,
input: Value::Null,
effect: Effect::Funds,
retry_class: IdempotencyMode::ReconcileBeforeRetry,
control_epochs: ControlEpochs::default(),
resource_scope_id: "resource".into(),
lease_epoch: 1,
action_epoch: 1,
deadline: None,
reservation: Some(ResourceReservationRef {
reservation_id: "reservation".into(),
fencing_token: 1,
}),
created_at: now,
};
assert!(intent.validate().is_ok());
assert!(intent.validate_prepared().is_ok());
intent.reservation = None;
assert!(intent.validate().is_err());
intent.reservation = Some(ResourceReservationRef {
reservation_id: "reservation".into(),
fencing_token: 1,
});
intent.resource_scope_id.clear();
assert!(intent.validate().is_err());
intent.resource_scope_id = "resource".into();
intent.retry_class = IdempotencyMode::None;
assert!(intent.validate().is_err());
intent.retry_class = IdempotencyMode::ReconcileBeforeRetry;
intent.state = ActionState::Authorized;
assert!(intent.validate_prepared().is_err());
assert!(!ActionState::Prepared.dispatch_committed());
assert!(ActionState::Unknown.dispatch_committed());
}
fn trigger_binding() -> TriggerBinding {
TriggerBinding {
id: "binding".into(),
revision: 1,
source: "source".into(),
event_type: "event".into(),
instance_id: "instance".parse().unwrap(),
predicate: serde_json::json!({}),
ordering: OrderingPolicy::Commutative,
starts_at: None,
expires_at: None,
gap_wait_ms: default_gap_wait_ms(),
gap_limit: default_gap_limit(),
}
}
#[test]
fn trigger_binding_rejects_blank_ids_bad_gaps_and_inverted_windows() {
assert!(trigger_binding().validate().is_ok());
let blank = TriggerBinding {
source: " ".into(),
..trigger_binding()
};
assert!(blank.validate().is_err());
for gap_wait_ms in [0, 3_600_001] {
let binding = TriggerBinding {
gap_wait_ms,
..trigger_binding()
};
assert!(binding.validate().is_err(), "gap_wait_ms {gap_wait_ms}");
}
for gap_limit in [0, 10_001] {
let binding = TriggerBinding {
gap_limit,
..trigger_binding()
};
assert!(binding.validate().is_err(), "gap_limit {gap_limit}");
}
let array_predicate = TriggerBinding {
predicate: serde_json::json!([1]),
..trigger_binding()
};
assert!(array_predicate.validate().is_err());
let now = Utc::now();
let inverted = TriggerBinding {
starts_at: Some(now),
expires_at: Some(now),
..trigger_binding()
};
assert!(inverted.validate().is_err());
let ordered = TriggerBinding {
starts_at: Some(now),
expires_at: Some(now + chrono::Duration::seconds(1)),
..trigger_binding()
};
assert!(ordered.validate().is_ok());
let defaults: TriggerBinding = serde_json::from_value(serde_json::json!({
"id": "binding", "revision": 1, "source": "s", "event_type": "e",
"instance_id": "i", "predicate": {}, "ordering": "commutative",
"starts_at": null, "expires_at": null
}))
.unwrap();
assert_eq!(defaults.gap_wait_ms, 30_000);
assert_eq!(defaults.gap_limit, 100);
}
#[test]
fn execution_profile_revision_requires_every_provider_name() {
let profile = ExecutionProfileRevision {
id: "profile".into(),
revision: 1,
content_digest: "sha256:profile".into(),
mode: ExecutionMode::Paper,
durability_grade: DurabilityGrade::Standard,
trigger_provider: "triggers".into(),
data_provider: "data".into(),
clock_model: "database".into(),
action_provider: "actions".into(),
models: BTreeMap::new(),
environment: Value::Null,
policy_bundle: Value::Null,
connection_bindings: BTreeMap::new(),
};
assert!(profile.validate().is_ok());
let blank_provider = ExecutionProfileRevision {
action_provider: String::new(),
..profile.clone()
};
assert!(blank_provider.validate().is_err());
let blank_digest = ExecutionProfileRevision {
content_digest: " ".into(),
..profile
};
assert!(blank_digest.validate().is_err());
}
#[test]
fn retry_backoff_grows_exponentially_with_bounded_jitter_and_cap() {
let policy = RetryPolicy {
max_attempts: 5,
timeout_ms: 1_000,
initial_backoff_ms: 100,
max_backoff_ms: 1_000,
};
assert_eq!(policy.backoff_ms(1, 0), 100);
assert_eq!(policy.backoff_ms(2, 0), 200);
assert_eq!(policy.backoff_ms(3, 0), 400);
assert_eq!(policy.backoff_ms(1, 24), 124);
assert_eq!(policy.backoff_ms(1, 25), 100);
assert_eq!(policy.backoff_ms(4, 249), 849);
assert_eq!(policy.backoff_ms(5, 0), 1_000);
assert_eq!(policy.backoff_ms(5, 249), 1_000);
assert_eq!(policy.backoff_ms(40, u64::MAX), 1_000);
assert_eq!(policy.backoff_ms(0, 0), 100);
}
#[test]
fn lifecycle_policy_rejects_inconsistent_windows_and_zero_counts() {
let now = Utc::now();
let later = now + chrono::Duration::hours(1);
assert!(LifecyclePolicy::run_once().validate().is_ok());
let inverted = LifecyclePolicy {
starts_at: Some(later),
expires_at: Some(now),
..LifecyclePolicy::run_once()
};
assert!(inverted.validate().is_err());
let drain_without_policy = LifecyclePolicy {
on_expiry: ExpiryPolicy::Cancel,
drain_deadline: Some(later),
..LifecyclePolicy::run_once()
};
assert!(drain_without_policy.validate().is_err());
let drain_before_expiry = LifecyclePolicy {
expires_at: Some(later),
drain_deadline: Some(now),
..LifecyclePolicy::run_once()
};
assert!(drain_before_expiry.validate().is_err());
let zero_timeout = LifecyclePolicy {
event_idle_timeout_ms: Some(0),
..LifecyclePolicy::run_once()
};
assert!(zero_timeout.validate().is_err());
for completion in [
CompletionPolicy::AfterMatchedEvaluations(0),
CompletionPolicy::AfterSuccessfulRuns(0),
] {
let zero_count = LifecyclePolicy {
completion,
..LifecyclePolicy::run_once()
};
assert!(zero_count.validate().is_err());
}
let bounded = LifecyclePolicy {
starts_at: Some(now),
expires_at: Some(later),
completion: CompletionPolicy::AfterSuccessfulRuns(2),
event_idle_timeout_ms: Some(1_000),
progress_timeout_ms: Some(1_000),
on_expiry: ExpiryPolicy::Drain,
drain_deadline: Some(later + chrono::Duration::minutes(1)),
};
assert!(bounded.validate().is_ok());
}
}