1use af_context::{InstanceId, RunId, SubjectId, TenantId};
8use std::collections::{BTreeMap, BTreeSet};
9
10use chrono::{DateTime, Utc};
11use serde::{Deserialize, Serialize};
12use serde_json::Value;
13
14#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
16pub struct WorkflowDefinition {
17 pub id: String,
19 pub name: String,
21}
22
23#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
25pub struct WorkflowInstance {
26 pub id: String,
28 pub tenant_id: TenantId,
30 pub subject_id: SubjectId,
32 pub definition_id: String,
34 pub revision: u64,
36 pub execution_profile_id: String,
38 pub execution_profile_revision: u64,
40 pub lifecycle: LifecyclePolicy,
42 pub status: String,
44 pub state_version: i64,
46 pub event_sequence: i64,
48 pub control_epoch: i64,
50}
51
52#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
54pub struct TriggerBinding {
55 pub id: String,
57 pub revision: u64,
59 pub source: String,
61 pub event_type: String,
63 pub instance_id: InstanceId,
65 #[serde(default)]
67 pub branch_id: Option<af_context::BranchId>,
68 pub predicate: Value,
70 pub ordering: OrderingPolicy,
72 pub starts_at: Option<DateTime<Utc>>,
74 pub expires_at: Option<DateTime<Utc>>,
76 #[serde(default = "default_gap_wait_ms")]
78 pub gap_wait_ms: u64,
79 #[serde(default = "default_gap_limit")]
81 pub gap_limit: u32,
82}
83
84impl TriggerBinding {
85 pub fn validate(&self) -> Result<(), ContractError> {
87 for (name, value) in [
88 ("trigger binding id", self.id.as_str()),
89 ("trigger source", self.source.as_str()),
90 ("trigger event_type", self.event_type.as_str()),
91 ("trigger instance_id", self.instance_id.as_str()),
92 ] {
93 required(name, value)?;
94 }
95 if self.gap_wait_ms == 0 || self.gap_wait_ms > 3_600_000 {
96 return Err(ContractError::Invalid(
97 "trigger gap_wait_ms must be between 1 and 3600000".into(),
98 ));
99 }
100 if self.gap_limit == 0 || self.gap_limit > 10_000 {
101 return Err(ContractError::Invalid(
102 "trigger gap_limit must be between 1 and 10000".into(),
103 ));
104 }
105 if !self.predicate.is_object() {
106 return Err(ContractError::Invalid(
107 "trigger predicate must be a JSON object".into(),
108 ));
109 }
110 if matches!((self.starts_at, self.expires_at), (Some(start), Some(end)) if end <= start) {
111 return Err(ContractError::Invalid(
112 "trigger expires_at must be after starts_at".into(),
113 ));
114 }
115 Ok(())
116 }
117}
118
119const fn default_gap_wait_ms() -> u64 {
120 30_000
121}
122
123const fn default_gap_limit() -> u32 {
124 100
125}
126
127#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
129pub struct WorkflowEvent {
130 pub instance_id: InstanceId,
132 pub sequence: i64,
134 pub event_type: String,
136 pub payload: Value,
138 pub content_digest: String,
140 pub occurred_at: DateTime<Utc>,
142}
143
144#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
146#[serde(rename_all = "snake_case")]
147pub enum ExecutionMode {
148 Simulation,
150 Backtest,
152 Paper,
154 Live,
156}
157
158#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
160#[serde(rename_all = "snake_case")]
161pub enum DurabilityGrade {
162 Standard,
164 FundsGrade,
166}
167
168#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
170pub struct ExecutionProfileRevision {
171 pub id: String,
173 pub revision: u64,
175 pub content_digest: String,
177 pub mode: ExecutionMode,
179 pub durability_grade: DurabilityGrade,
181 pub trigger_provider: String,
183 pub data_provider: String,
185 pub clock_model: String,
187 pub action_provider: String,
189 #[serde(default)]
191 pub models: BTreeMap<String, Value>,
192 #[serde(default)]
194 pub environment: Value,
195 #[serde(default)]
197 pub policy_bundle: Value,
198 #[serde(default)]
200 pub connection_bindings: BTreeMap<String, String>,
201}
202
203impl ExecutionProfileRevision {
204 pub fn validate(&self) -> Result<(), ContractError> {
206 for (name, value) in [
207 ("execution profile id", self.id.as_str()),
208 (
209 "execution profile content_digest",
210 self.content_digest.as_str(),
211 ),
212 ("trigger_provider", self.trigger_provider.as_str()),
213 ("data_provider", self.data_provider.as_str()),
214 ("clock_model", self.clock_model.as_str()),
215 ("action_provider", self.action_provider.as_str()),
216 ] {
217 required(name, value)?;
218 }
219 Ok(())
220 }
221}
222
223#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
225pub struct CapabilityPin {
226 pub id: String,
228 pub contract_version: String,
230 pub content_digest: String,
232}
233
234#[derive(Debug, Clone, Serialize, Deserialize)]
236pub struct WorkflowRevision {
237 pub definition_id: String,
239 pub revision: u64,
241 pub content_digest: String,
243 pub kernel_abi_version: String,
245 pub dependency_set_digest: String,
247 #[serde(default)]
249 pub expression_versions: BTreeMap<String, String>,
250 #[serde(default)]
252 pub capabilities: Vec<CapabilityPin>,
253 #[serde(default)]
255 pub template_provenance: Value,
256 pub spec: crate::Spec,
258}
259
260impl WorkflowRevision {
261 pub fn validate(&self) -> Result<(), ContractError> {
263 required("definition_id", &self.definition_id)?;
264 required("content_digest", &self.content_digest)?;
265 required("kernel_abi_version", &self.kernel_abi_version)?;
266 required("dependency_set_digest", &self.dependency_set_digest)?;
267 self.spec
268 .validate_structure()
269 .map_err(|error| ContractError::Invalid(error.to_string()))?;
270 let unique = self
271 .capabilities
272 .iter()
273 .map(|pin| &pin.id)
274 .collect::<BTreeSet<_>>();
275 if unique.len() != self.capabilities.len() {
276 return Err(ContractError::Invalid(
277 "capability pins must be unique by id within one revision".into(),
278 ));
279 }
280 for pin in &self.capabilities {
281 required("capability id", &pin.id)?;
282 required("capability contract_version", &pin.contract_version)?;
283 required("capability content_digest", &pin.content_digest)?;
284 }
285 Ok(())
286 }
287}
288
289#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
291#[serde(rename_all = "snake_case")]
292pub enum CapabilityKind {
293 Trigger,
295 Expression,
297 Guard,
299 Action,
301 Sink,
303}
304
305#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
307#[serde(rename_all = "snake_case")]
308pub enum Effect {
309 Pure,
311 Read,
313 InternalWrite,
315 ExternalWrite,
317 Funds,
319}
320
321#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
323#[serde(rename_all = "snake_case")]
324pub enum IdempotencyMode {
325 None,
327 Native,
329 ReconcileBeforeRetry,
331 NeverAutomaticRetry,
333}
334
335#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
337#[serde(rename_all = "snake_case")]
338pub enum CapabilityLifecycle {
339 Installed,
341 Active,
343 Deprecated,
345 Disabled,
347 Unavailable,
349 EmergencyRevoked,
351}
352
353#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
355pub struct RetryPolicy {
356 pub max_attempts: u32,
358 pub timeout_ms: u64,
360 pub initial_backoff_ms: u64,
362 pub max_backoff_ms: u64,
364}
365
366impl RetryPolicy {
367 pub fn backoff_ms(&self, attempt: u32, jitter_seed: u64) -> u64 {
370 let factor = 1_u64
371 .checked_shl(attempt.saturating_sub(1).min(20))
372 .unwrap_or(u64::MAX);
373 let base = self
374 .initial_backoff_ms
375 .saturating_mul(factor)
376 .min(self.max_backoff_ms);
377 let jitter_ceiling = (base / 4).max(1);
378 base.saturating_add(jitter_seed % jitter_ceiling)
379 .min(self.max_backoff_ms)
380 }
381}
382
383#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
385pub struct CapabilityManifest {
386 pub id: String,
388 pub contract_version: String,
390 pub content_digest: String,
392 pub kind: CapabilityKind,
394 #[serde(default, skip_serializing_if = "Option::is_none")]
396 pub guard_kind: Option<GuardKind>,
397 #[serde(default = "object_schema")]
399 pub config_schema: Value,
400 pub input_schema: Value,
402 pub output_schema: Value,
404 pub effect: Effect,
406 pub deterministic: bool,
408 pub idempotency_mode: IdempotencyMode,
410 pub retry: RetryPolicy,
412 #[serde(default)]
414 pub permissions: BTreeSet<String>,
415 #[serde(default)]
417 pub required_guards: BTreeSet<GuardKind>,
418 #[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
420 pub required_guard_pins: BTreeSet<CapabilityPin>,
421 #[serde(default)]
423 pub taint_rules: Value,
424 #[serde(default)]
426 pub resource_cost: Value,
427 pub supports_simulation: bool,
429 pub supports_replay: bool,
431 pub supports_paper: bool,
433 pub supports_live: bool,
435 #[serde(default)]
437 pub clock_requirements: Value,
438 #[serde(default)]
440 pub data_requirements: Value,
441 pub supports_reconciliation: bool,
443 pub lifecycle: CapabilityLifecycle,
445}
446
447fn object_schema() -> Value {
448 serde_json::json!({"type":"object"})
449}
450
451impl CapabilityManifest {
452 pub fn action(
454 id: impl Into<String>,
455 contract_version: impl Into<String>,
456 content_digest: impl Into<String>,
457 effect: Effect,
458 idempotency_mode: IdempotencyMode,
459 supports_reconciliation: bool,
460 ) -> Self {
461 let required_guards = match effect {
462 Effect::Funds => BTreeSet::from([
463 GuardKind::Authorization,
464 GuardKind::Freshness,
465 GuardKind::Reservation,
466 ]),
467 Effect::ExternalWrite => BTreeSet::from([GuardKind::Authorization]),
468 _ => BTreeSet::new(),
469 };
470 Self {
471 id: id.into(),
472 contract_version: contract_version.into(),
473 content_digest: content_digest.into(),
474 kind: CapabilityKind::Action,
475 guard_kind: None,
476 config_schema: object_schema(),
477 input_schema: serde_json::json!({"type": "object"}),
478 output_schema: serde_json::json!({"type": "object"}),
479 effect,
480 deterministic: false,
481 idempotency_mode,
482 retry: RetryPolicy {
483 max_attempts: 1,
484 timeout_ms: 30_000,
485 initial_backoff_ms: 100,
486 max_backoff_ms: 5_000,
487 },
488 permissions: BTreeSet::new(),
489 required_guards,
490 required_guard_pins: BTreeSet::new(),
491 taint_rules: Value::Null,
492 resource_cost: Value::Null,
493 supports_simulation: true,
494 supports_replay: false,
495 supports_paper: true,
496 supports_live: true,
497 clock_requirements: Value::Null,
498 data_requirements: Value::Null,
499 supports_reconciliation,
500 lifecycle: CapabilityLifecycle::Active,
501 }
502 }
503
504 pub fn validate(&self) -> Result<(), ContractError> {
506 required("capability id", &self.id)?;
507 required("contract_version", &self.contract_version)?;
508 required("content_digest", &self.content_digest)?;
509 match (self.kind, self.guard_kind) {
510 (CapabilityKind::Guard, None) => {
511 return Err(ContractError::Invalid(
512 "guard capability must declare guard_kind".into(),
513 ));
514 }
515 (CapabilityKind::Guard, Some(_)) | (_, None) => {}
516 (_, Some(_)) => {
517 return Err(ContractError::Invalid(
518 "guard_kind is only valid for guard capabilities".into(),
519 ));
520 }
521 }
522 for pin in &self.required_guard_pins {
523 required("required guard id", &pin.id)?;
524 required("required guard contract_version", &pin.contract_version)?;
525 required("required guard content_digest", &pin.content_digest)?;
526 }
527 if !self.config_schema.is_object()
528 || jsonschema::validator_for(&self.config_schema).is_err()
529 {
530 return Err(ContractError::Invalid(
531 "capability config_schema must be a valid JSON Schema object".into(),
532 ));
533 }
534 if self.retry.max_attempts == 0
535 || self.retry.timeout_ms == 0
536 || self.retry.initial_backoff_ms == 0
537 || self.retry.max_backoff_ms < self.retry.initial_backoff_ms
538 {
539 return Err(ContractError::Invalid(
540 "retry attempts, timeout and backoff bounds are invalid".into(),
541 ));
542 }
543 if self.effect == Effect::Funds {
544 if self.kind != CapabilityKind::Action {
545 return Err(ContractError::Invalid(
546 "funds effect is only valid for action capabilities".into(),
547 ));
548 }
549 if !self.supports_reconciliation && self.idempotency_mode != IdempotencyMode::Native {
550 return Err(ContractError::Invalid(
551 "funds actions require native idempotency or reconciliation".into(),
552 ));
553 }
554 for guard in [
555 GuardKind::Authorization,
556 GuardKind::Freshness,
557 GuardKind::Reservation,
558 ] {
559 if !self.required_guards.contains(&guard) {
560 return Err(ContractError::Invalid(format!(
561 "funds action must require {guard:?} guard"
562 )));
563 }
564 }
565 }
566 Ok(())
567 }
568
569 pub fn can_start_new_work(&self) -> bool {
571 matches!(
572 self.lifecycle,
573 CapabilityLifecycle::Installed
574 | CapabilityLifecycle::Active
575 | CapabilityLifecycle::Deprecated
576 )
577 }
578}
579
580#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
582#[serde(rename_all = "snake_case")]
583pub enum GuardKind {
584 Authorization,
586 Freshness,
588 Reservation,
590 Policy,
592}
593
594#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
596#[serde(rename_all = "snake_case")]
597pub enum OrderingPolicy {
598 StrictSequence,
600 Commutative,
602 LatestStateReconcile,
604 RejectUnordered,
606}
607
608#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
610pub struct TriggerEnvelope {
611 pub event_id: String,
613 pub event_type: String,
615 pub source: String,
617 pub schema_version: String,
619 pub tenant_id: TenantId,
621 pub subject_id: SubjectId,
623 pub aggregate_id: String,
625 pub source_sequence: Option<i64>,
627 pub observed_version: Option<String>,
629 pub occurred_at: DateTime<Utc>,
631 pub received_at: DateTime<Utc>,
633 pub watermark: Option<DateTime<Utc>>,
635 pub correlation_key: String,
637 pub dedup_key: String,
639 pub cursor: Option<String>,
641 pub payload: Value,
643 #[serde(default)]
645 pub trace_context: Value,
646}
647
648impl TriggerEnvelope {
649 pub fn validate(&self) -> Result<(), ContractError> {
651 for (name, value) in [
652 ("event_id", self.event_id.as_str()),
653 ("event_type", self.event_type.as_str()),
654 ("source", self.source.as_str()),
655 ("schema_version", self.schema_version.as_str()),
656 ("tenant_id", self.tenant_id.as_str()),
657 ("aggregate_id", self.aggregate_id.as_str()),
658 ("correlation_key", self.correlation_key.as_str()),
659 ("dedup_key", self.dedup_key.as_str()),
660 ] {
661 required(name, value)?;
662 }
663 if self
664 .cursor
665 .as_ref()
666 .is_some_and(|cursor| cursor.is_empty() || cursor.len() > 4096)
667 {
668 return Err(ContractError::Invalid(
669 "source cursor requires 1..4096 bytes".into(),
670 ));
671 }
672 if self.occurred_at > self.received_at + chrono::Duration::minutes(5) {
673 return Err(ContractError::Invalid(
674 "occurred_at is implausibly ahead of received_at".into(),
675 ));
676 }
677 Ok(())
678 }
679}
680
681#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
683pub struct SourceAdapterState {
684 pub source: af_context::WorkflowSourceId,
686 pub cursor: String,
688 pub source_event_id: Option<String>,
690}
691#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
693pub struct SourceIngestionReceipt {
694 pub event_id: String,
696 pub duplicate: bool,
698 pub deliveries: u64,
700 pub committed_cursor: Option<String>,
702}
703
704#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
706#[serde(rename_all = "snake_case")]
707pub enum CompletionPolicy {
708 ExplicitStop,
710 FirstTrigger,
712 FirstMatch,
714 FirstActionTerminal,
716 FirstSuccess,
718 AfterMatchedEvaluations(u64),
720 AfterSuccessfulRuns(u64),
722}
723
724#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
726#[serde(rename_all = "snake_case")]
727pub enum ExpiryPolicy {
728 Drain,
730 Cancel,
732}
733
734#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
736#[serde(rename_all = "snake_case")]
737pub enum ScheduleCadence {
738 Cron {
740 expression: String,
742 },
743 FixedRate {
745 milliseconds: u64,
747 },
748 FixedDelay {
750 milliseconds: u64,
752 },
753}
754
755#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
757#[serde(rename_all = "snake_case")]
758pub enum CatchUpPolicy {
759 Skip,
761 CatchUpOnce,
763 CatchUpAll {
765 limit: u32,
767 },
768}
769
770#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
772pub struct SchedulePolicy {
773 pub cadence: ScheduleCadence,
775 pub timezone: String,
777 pub catch_up: CatchUpPolicy,
779}
780
781impl SchedulePolicy {
782 pub fn validate(&self) -> Result<(), ContractError> {
784 self.timezone.parse::<chrono_tz::Tz>().map_err(|_| {
785 ContractError::Invalid(format!("unknown IANA timezone '{}'", self.timezone))
786 })?;
787 match &self.cadence {
788 ScheduleCadence::Cron { expression } => {
789 expression.parse::<cron::Schedule>().map_err(|error| {
790 ContractError::Invalid(format!("invalid cron '{expression}': {error}"))
791 })?;
792 }
793 ScheduleCadence::FixedRate { milliseconds }
794 | ScheduleCadence::FixedDelay { milliseconds }
795 if *milliseconds == 0 =>
796 {
797 return Err(ContractError::Invalid(
798 "schedule interval must be positive".into(),
799 ));
800 }
801 _ => {}
802 }
803 if matches!(self.catch_up, CatchUpPolicy::CatchUpAll { limit: 0 }) {
804 return Err(ContractError::Invalid(
805 "catch_up_all limit must be positive".into(),
806 ));
807 }
808 Ok(())
809 }
810}
811
812#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
814pub struct LifecyclePolicy {
815 pub starts_at: Option<DateTime<Utc>>,
817 pub expires_at: Option<DateTime<Utc>>,
819 pub completion: CompletionPolicy,
821 pub event_idle_timeout_ms: Option<u64>,
823 pub progress_timeout_ms: Option<u64>,
825 pub on_expiry: ExpiryPolicy,
827 pub drain_deadline: Option<DateTime<Utc>>,
829}
830
831impl LifecyclePolicy {
832 pub fn run_once() -> Self {
834 Self {
835 starts_at: None,
836 expires_at: None,
837 completion: CompletionPolicy::FirstSuccess,
838 event_idle_timeout_ms: None,
839 progress_timeout_ms: None,
840 on_expiry: ExpiryPolicy::Drain,
841 drain_deadline: None,
842 }
843 }
844
845 pub fn validate(&self) -> Result<(), ContractError> {
847 if self
848 .starts_at
849 .zip(self.expires_at)
850 .is_some_and(|(a, b)| a >= b)
851 {
852 return Err(ContractError::Invalid(
853 "lifecycle starts_at must be before expires_at".into(),
854 ));
855 }
856 if self.drain_deadline.is_some() && self.on_expiry != ExpiryPolicy::Drain {
857 return Err(ContractError::Invalid(
858 "drain_deadline requires on_expiry=drain".into(),
859 ));
860 }
861 if self
862 .expires_at
863 .zip(self.drain_deadline)
864 .is_some_and(|(expires, drain)| drain <= expires)
865 {
866 return Err(ContractError::Invalid(
867 "lifecycle drain_deadline must be after expires_at".into(),
868 ));
869 }
870 if self.event_idle_timeout_ms == Some(0) || self.progress_timeout_ms == Some(0) {
871 return Err(ContractError::Invalid(
872 "lifecycle timeouts must be positive".into(),
873 ));
874 }
875 if matches!(
876 self.completion,
877 CompletionPolicy::AfterMatchedEvaluations(0) | CompletionPolicy::AfterSuccessfulRuns(0)
878 ) {
879 return Err(ContractError::Invalid(
880 "lifecycle completion count must be positive".into(),
881 ));
882 }
883 Ok(())
884 }
885}
886
887#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
889#[serde(rename_all = "snake_case")]
890pub enum StepOutcomeKind {
891 Succeeded,
893 Skipped,
895 Waiting,
897 Failed,
899 Cancelled,
901}
902
903#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
905pub struct StepOutcome {
906 pub kind: StepOutcomeKind,
908 pub reason_code: Option<String>,
910 #[serde(default)]
912 pub wake_condition: Value,
913}
914
915#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
917#[serde(rename_all = "snake_case")]
918pub enum ActionState {
919 Prepared,
921 AwaitingConfirmation,
923 Authorized,
925 DispatchCommitted,
927 Executing,
929 Succeeded,
931 Rejected,
933 Retryable,
935 Unknown,
937 Reconciled,
939}
940
941impl ActionState {
942 pub fn can_transition_to(self, next: Self) -> bool {
944 use ActionState::*;
945 matches!(
946 (self, next),
947 (Prepared, AwaitingConfirmation | Authorized | Rejected)
948 | (AwaitingConfirmation, Authorized | Rejected)
949 | (Authorized, DispatchCommitted | Rejected)
950 | (
951 DispatchCommitted,
952 Executing | Succeeded | Rejected | Retryable | Unknown
953 )
954 | (Executing, Succeeded | Rejected | Retryable | Unknown)
955 | (Retryable, DispatchCommitted | Reconciled)
956 | (Unknown, Reconciled)
957 | (Succeeded | Rejected, Reconciled)
958 )
959 }
960
961 pub fn dispatch_committed(self) -> bool {
963 matches!(
964 self,
965 Self::DispatchCommitted
966 | Self::Executing
967 | Self::Succeeded
968 | Self::Rejected
969 | Self::Retryable
970 | Self::Unknown
971 | Self::Reconciled
972 )
973 }
974}
975
976#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
978pub struct ControlEpochs {
979 pub tenant: i64,
981 pub resource: i64,
983 pub instance: i64,
985}
986
987#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
989#[serde(rename_all = "snake_case")]
990pub enum ControlScope {
991 Tenant,
993 Resource,
995 Instance,
997}
998
999#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1001#[serde(rename_all = "snake_case")]
1002pub enum ControlMode {
1003 Running,
1005 Paused,
1007 Stopped,
1009 EmergencyStopped,
1011}
1012
1013#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1015pub struct ControlCommand {
1016 pub scope: ControlScope,
1018 pub scope_id: String,
1020 pub mode: ControlMode,
1022 pub operator_subject_id: SubjectId,
1024 pub reason: String,
1026 pub override_expires_at: Option<DateTime<Utc>>,
1028}
1029
1030impl ControlCommand {
1031 pub fn validate(&self) -> Result<(), ContractError> {
1033 required("control scope_id", &self.scope_id)?;
1034 required("control operator_subject_id", &self.operator_subject_id)?;
1035 required("control reason", &self.reason)?;
1036 Ok(())
1037 }
1038}
1039
1040#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1042pub struct ResourceReservationRef {
1043 pub reservation_id: String,
1045 pub fencing_token: i64,
1047}
1048
1049#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1051pub struct ActionIntent {
1052 pub id: String,
1054 pub tenant_id: TenantId,
1056 pub instance_id: InstanceId,
1058 pub run_id: RunId,
1060 pub capability: CapabilityPin,
1062 pub idempotency_key: String,
1064 pub state: ActionState,
1066 pub input: Value,
1068 pub effect: Effect,
1070 pub retry_class: IdempotencyMode,
1072 pub control_epochs: ControlEpochs,
1074 pub resource_scope_id: String,
1076 pub lease_epoch: i64,
1078 pub action_epoch: i64,
1080 pub deadline: Option<DateTime<Utc>>,
1082 pub reservation: Option<ResourceReservationRef>,
1084 pub created_at: DateTime<Utc>,
1086}
1087
1088impl ActionIntent {
1089 pub fn validate(&self) -> Result<(), ContractError> {
1091 required("action id", &self.id)?;
1092 required("action idempotency_key", &self.idempotency_key)?;
1093 if self.effect == Effect::Funds && self.reservation.is_none() {
1094 return Err(ContractError::Invalid(
1095 "funds action requires a resource reservation".into(),
1096 ));
1097 }
1098 if self.effect == Effect::Funds && self.resource_scope_id.trim().is_empty() {
1099 return Err(ContractError::Invalid(
1100 "funds action requires a resource scope".into(),
1101 ));
1102 }
1103 if self.effect == Effect::Funds && self.retry_class == IdempotencyMode::None {
1104 return Err(ContractError::Invalid(
1105 "funds action requires an explicit retry class".into(),
1106 ));
1107 }
1108 if self
1109 .deadline
1110 .is_some_and(|deadline| deadline <= self.created_at)
1111 {
1112 return Err(ContractError::Invalid(
1113 "action deadline must be after creation".into(),
1114 ));
1115 }
1116 Ok(())
1117 }
1118
1119 pub fn validate_prepared(&self) -> Result<(), ContractError> {
1122 self.validate()?;
1123 if self.state != ActionState::Prepared {
1124 return Err(ContractError::Invalid(
1125 "new action intent must start in prepared state".into(),
1126 ));
1127 }
1128 Ok(())
1129 }
1130}
1131
1132#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1134pub struct ActionReceipt {
1135 pub id: String,
1137 pub action_intent_id: String,
1139 pub provider_version: String,
1141 pub received_at: DateTime<Utc>,
1143 pub outcome: ActionState,
1145 pub payload: Value,
1147 pub raw_receipt_digest: String,
1149}
1150
1151#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1153pub struct ActionObservation {
1154 pub id: String,
1156 pub action_intent_id: String,
1158 pub provider_version: String,
1160 pub observed_at: DateTime<Utc>,
1162 pub state: String,
1164 pub resource_ref: Option<Value>,
1166 pub raw_receipt_digest: String,
1168 pub terminal: bool,
1170 #[serde(default)]
1173 pub retry_authorized: bool,
1174}
1175
1176#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1178pub struct SourceGap {
1179 pub tenant_id: TenantId,
1181 pub source: String,
1183 pub aggregate_id: String,
1185 pub missing_from: i64,
1187 pub missing_to: i64,
1189 pub deadline: DateTime<Utc>,
1191 pub status: String,
1193}
1194
1195#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1197#[serde(rename_all = "snake_case")]
1198pub enum ReservationScope {
1199 Tenant,
1201 Resource,
1203}
1204
1205#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1207#[serde(rename_all = "snake_case")]
1208pub enum ReservationState {
1209 Reserved,
1211 Consumed,
1213 Released,
1215 Expired,
1217}
1218
1219#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1221pub struct ResourceReservation {
1222 pub id: String,
1224 pub tenant_id: TenantId,
1226 pub scope: ReservationScope,
1228 pub scope_id: String,
1230 pub resource_kind: String,
1232 pub amount: String,
1234 pub policy_version: String,
1236 pub expires_at: DateTime<Utc>,
1238 pub fencing_token: i64,
1240 pub provider_reservation_id: Option<String>,
1242 pub state: ReservationState,
1244}
1245
1246#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1248pub struct ObservedValue<T> {
1249 pub value: T,
1251 pub source: String,
1253 pub observed_at: DateTime<Utc>,
1255 pub received_at: DateTime<Utc>,
1257 pub source_version: String,
1259 pub quality: String,
1261 pub digest: String,
1263}
1264
1265impl<T> ObservedValue<T> {
1266 pub fn is_accepted(
1268 &self,
1269 now: DateTime<Utc>,
1270 max_age: chrono::Duration,
1271 accepted_quality: &BTreeSet<String>,
1272 ) -> bool {
1273 self.observed_at <= now
1274 && now - self.observed_at <= max_age
1275 && accepted_quality.contains(&self.quality)
1276 }
1277}
1278
1279#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1281pub struct DecisionSnapshot {
1282 pub input_digests: Vec<String>,
1284 pub instance_config_version: String,
1286 pub policy_version: String,
1288 pub capability_versions: Vec<CapabilityPin>,
1290 pub execution_profile_revision_id: String,
1292 #[serde(default)]
1294 pub context_snapshot: Value,
1295 #[serde(default)]
1297 pub policy_snapshot: Value,
1298 #[serde(default)]
1300 pub agent_provenance: Value,
1301}
1302
1303#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1305#[serde(rename_all = "snake_case")]
1306pub enum ParameterProvenance {
1307 UserSupplied,
1309 ProductDefault,
1311 AgentInferred,
1313 Derived,
1315}
1316
1317#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1319pub struct ParameterValue {
1320 pub value: Value,
1322 pub provenance: ParameterProvenance,
1324}
1325
1326#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1328pub struct ParameterSpec {
1329 pub name: String,
1331 pub required: bool,
1333 pub required_explicit_for_live: bool,
1335}
1336
1337pub fn validate_parameters(
1339 mode: ExecutionMode,
1340 specs: &[ParameterSpec],
1341 values: &BTreeMap<String, ParameterValue>,
1342) -> Result<(), MissingRequirements> {
1343 let missing = specs
1344 .iter()
1345 .filter(|spec| match values.get(&spec.name) {
1346 None => {
1347 spec.required || (mode == ExecutionMode::Live && spec.required_explicit_for_live)
1348 }
1349 Some(value) => {
1350 mode == ExecutionMode::Live
1351 && spec.required_explicit_for_live
1352 && value.provenance != ParameterProvenance::UserSupplied
1353 }
1354 })
1355 .map(|spec| spec.name.clone())
1356 .collect::<Vec<_>>();
1357 if missing.is_empty() {
1358 Ok(())
1359 } else {
1360 Err(MissingRequirements {
1361 parameters: missing,
1362 })
1363 }
1364}
1365
1366#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, thiserror::Error)]
1368#[error("missing explicit workflow requirements: {parameters:?}")]
1369pub struct MissingRequirements {
1370 pub parameters: Vec<String>,
1372}
1373
1374#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1376pub struct RootOperationBudget {
1377 pub root_operation_id: String,
1379 pub max_depth: u32,
1381 pub descendant_limit: u32,
1383 pub run_limit: u32,
1385 pub token_budget: u64,
1387 pub cost_budget_micros: u64,
1389 pub action_budget: u32,
1391 pub deadline: DateTime<Utc>,
1393 #[serde(default)]
1395 pub permission_ceiling: BTreeSet<String>,
1396}
1397
1398#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1400pub struct DecisionArtifact {
1401 pub id: String,
1403 pub root_operation_id: String,
1405 pub content_digest: String,
1407 pub model: String,
1409 pub prompt_digest: String,
1411 pub output: Value,
1413 pub created_at: DateTime<Utc>,
1415}
1416
1417#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1419pub struct DiagnosticRecord {
1420 pub code: String,
1422 pub message: String,
1424 pub at: DateTime<Utc>,
1426 #[serde(default)]
1428 pub detail: Value,
1429}
1430
1431#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1433pub struct ExecutionProjection {
1434 pub schema_version: String,
1436 pub source_sequence: i64,
1438 pub current_steps: Vec<String>,
1440 pub completed_steps: Vec<String>,
1442 pub waiting_on: Option<Value>,
1444 pub last_decision: Option<Value>,
1446 pub next_possible_steps: Vec<String>,
1448 pub next_trigger_at: Option<DateTime<Utc>>,
1450 pub planned_actions: Vec<String>,
1452 pub latest_diagnostic: Option<DiagnosticRecord>,
1454 pub progress: Value,
1456 pub expires_at: Option<DateTime<Utc>>,
1458}
1459
1460#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1462pub enum ContractError {
1463 #[error("invalid workflow contract: {0}")]
1465 Invalid(String),
1466}
1467
1468fn required(name: &str, value: &str) -> Result<(), ContractError> {
1469 if value.trim().is_empty() {
1470 Err(ContractError::Invalid(format!("{name} is required")))
1471 } else {
1472 Ok(())
1473 }
1474}
1475
1476#[cfg(test)]
1477mod tests {
1478 use super::*;
1479
1480 fn empty_spec() -> crate::Spec {
1481 crate::Spec {
1482 spec_id: "test".into(),
1483 version: "1".into(),
1484 description: String::new(),
1485 aliases: vec![],
1486 instance_config_schema: None,
1487 display: BTreeMap::new(),
1488 branches: vec![],
1489 }
1490 }
1491
1492 fn funds_manifest() -> CapabilityManifest {
1493 CapabilityManifest::action(
1494 "action.example",
1495 "1",
1496 "sha256:x",
1497 Effect::Funds,
1498 IdempotencyMode::ReconcileBeforeRetry,
1499 true,
1500 )
1501 }
1502
1503 #[test]
1504 fn funds_capability_requires_reconciliation_and_all_guards() {
1505 assert!(funds_manifest().validate().is_ok());
1506 let mut invalid = funds_manifest();
1507 invalid.required_guards.remove(&GuardKind::Freshness);
1508 assert!(invalid.validate().is_err());
1509 invalid.required_guards.insert(GuardKind::Freshness);
1510 invalid.supports_reconciliation = false;
1511 assert!(invalid.validate().is_err());
1512 }
1513
1514 #[test]
1515 fn action_state_never_skips_dispatch_commit() {
1516 assert!(ActionState::Authorized.can_transition_to(ActionState::DispatchCommitted));
1517 assert!(!ActionState::Authorized.can_transition_to(ActionState::Succeeded));
1518 assert!(ActionState::Unknown.can_transition_to(ActionState::Reconciled));
1519 }
1520
1521 #[test]
1522 fn live_parameters_cannot_be_silently_inferred() {
1523 let specs = [ParameterSpec {
1524 name: "account".into(),
1525 required: true,
1526 required_explicit_for_live: true,
1527 }];
1528 let inferred = BTreeMap::from([(
1529 "account".into(),
1530 ParameterValue {
1531 value: Value::String("a".into()),
1532 provenance: ParameterProvenance::AgentInferred,
1533 },
1534 )]);
1535 assert_eq!(
1536 validate_parameters(ExecutionMode::Live, &specs, &inferred)
1537 .unwrap_err()
1538 .parameters,
1539 ["account"]
1540 );
1541 assert!(validate_parameters(ExecutionMode::Paper, &specs, &inferred).is_ok());
1542 }
1543
1544 #[test]
1545 fn freshness_is_fail_closed() {
1546 let now = Utc::now();
1547 let observed = ObservedValue {
1548 value: 1,
1549 source: "source".into(),
1550 observed_at: now - chrono::Duration::seconds(2),
1551 received_at: now,
1552 source_version: "1".into(),
1553 quality: "good".into(),
1554 digest: "d".into(),
1555 };
1556 assert!(observed.is_accepted(
1557 now,
1558 chrono::Duration::seconds(3),
1559 &BTreeSet::from(["good".into()])
1560 ));
1561 assert!(!observed.is_accepted(
1562 now,
1563 chrono::Duration::seconds(1),
1564 &BTreeSet::from(["good".into()])
1565 ));
1566 }
1567
1568 #[test]
1569 fn revision_and_manifest_contracts_fail_closed() {
1570 let mut revision = WorkflowRevision {
1571 definition_id: "definition".into(),
1572 revision: 1,
1573 content_digest: "digest".into(),
1574 kernel_abi_version: "1".into(),
1575 dependency_set_digest: "dependencies".into(),
1576 expression_versions: BTreeMap::new(),
1577 capabilities: vec![CapabilityPin {
1578 id: "action.example".into(),
1579 contract_version: "1".into(),
1580 content_digest: "capability-digest".into(),
1581 }],
1582 template_provenance: Value::Null,
1583 spec: empty_spec(),
1584 };
1585 assert!(revision.validate().is_ok());
1586 revision.capabilities.push(CapabilityPin {
1587 id: "action.example".into(),
1588 contract_version: "2".into(),
1589 content_digest: "capability-digest-v2".into(),
1590 });
1591 assert!(revision.validate().is_err());
1592 revision.capabilities.pop();
1593 revision.kernel_abi_version.clear();
1594 assert!(revision.validate().is_err());
1595
1596 let external = CapabilityManifest::action(
1597 "action.notify",
1598 "1",
1599 "digest",
1600 Effect::ExternalWrite,
1601 IdempotencyMode::NeverAutomaticRetry,
1602 false,
1603 );
1604 assert_eq!(
1605 external.required_guards,
1606 BTreeSet::from([GuardKind::Authorization])
1607 );
1608 assert!(external.validate().is_ok());
1609 assert_eq!(external.config_schema, serde_json::json!({"type":"object"}));
1610 assert!(serde_json::to_value(&external)
1611 .unwrap()
1612 .get("guard_kind")
1613 .is_none());
1614 let mut guard = external.clone();
1615 guard.kind = CapabilityKind::Guard;
1616 assert!(guard.validate().is_err());
1617 let old_guard: CapabilityManifest =
1618 serde_json::from_value(serde_json::to_value(&guard).unwrap()).unwrap();
1619 assert_eq!(old_guard.guard_kind, None);
1620 assert!(old_guard.validate().is_err());
1621 guard.guard_kind = Some(GuardKind::Authorization);
1622 assert!(guard.validate().is_ok());
1623 assert_eq!(
1624 serde_json::to_value(&guard).unwrap()["guard_kind"],
1625 serde_json::json!("authorization")
1626 );
1627 let mut wrong_guard_owner = external.clone();
1628 wrong_guard_owner.guard_kind = Some(GuardKind::Authorization);
1629 assert!(wrong_guard_owner.validate().is_err());
1630 let mut invalid_schema = external.clone();
1631 invalid_schema.config_schema = serde_json::json!({"type":7});
1632 assert!(invalid_schema.validate().is_err());
1633 let mut invalid = external;
1634 invalid.retry.max_attempts = 0;
1635 assert!(invalid.validate().is_err());
1636
1637 let mut wrong_kind = funds_manifest();
1638 wrong_kind.kind = CapabilityKind::Expression;
1639 assert!(wrong_kind.validate().is_err());
1640 wrong_kind.kind = CapabilityKind::Action;
1641 wrong_kind.idempotency_mode = IdempotencyMode::Native;
1642 wrong_kind.supports_reconciliation = false;
1643 assert!(wrong_kind.validate().is_ok());
1644 wrong_kind.lifecycle = CapabilityLifecycle::Disabled;
1645 assert!(!wrong_kind.can_start_new_work());
1646 }
1647
1648 #[test]
1649 fn trigger_schedule_lifecycle_and_control_validate_boundaries() {
1650 let now = Utc::now();
1651 let mut trigger = TriggerEnvelope {
1652 event_id: "event".into(),
1653 event_type: "example".into(),
1654 source: "source".into(),
1655 schema_version: "1".into(),
1656 tenant_id: "tenant".parse().unwrap(),
1657 subject_id: "subject".parse().unwrap(),
1658 aggregate_id: "aggregate".into(),
1659 source_sequence: Some(1),
1660 observed_version: None,
1661 occurred_at: now,
1662 received_at: now,
1663 watermark: None,
1664 correlation_key: "key".into(),
1665 dedup_key: "dedup".into(),
1666 cursor: None,
1667 payload: Value::Null,
1668 trace_context: Value::Null,
1669 };
1670 assert!(trigger.validate().is_ok());
1671 trigger.occurred_at = now + chrono::Duration::minutes(6);
1672 assert!(trigger.validate().is_err());
1673
1674 let valid_schedule = SchedulePolicy {
1675 cadence: ScheduleCadence::Cron {
1676 expression: "0 0 * * * *".into(),
1677 },
1678 timezone: "Asia/Shanghai".into(),
1679 catch_up: CatchUpPolicy::CatchUpOnce,
1680 };
1681 assert!(valid_schedule.validate().is_ok());
1682 for invalid in [
1683 SchedulePolicy {
1684 timezone: "Nowhere/Invalid".into(),
1685 ..valid_schedule.clone()
1686 },
1687 SchedulePolicy {
1688 cadence: ScheduleCadence::FixedRate { milliseconds: 0 },
1689 ..valid_schedule.clone()
1690 },
1691 SchedulePolicy {
1692 catch_up: CatchUpPolicy::CatchUpAll { limit: 0 },
1693 ..valid_schedule
1694 },
1695 ] {
1696 assert!(invalid.validate().is_err());
1697 }
1698
1699 let mut lifecycle = LifecyclePolicy::run_once();
1700 assert!(lifecycle.validate().is_ok());
1701 lifecycle.starts_at = Some(now);
1702 lifecycle.expires_at = Some(now);
1703 assert!(lifecycle.validate().is_err());
1704 lifecycle.starts_at = None;
1705 lifecycle.expires_at = None;
1706 lifecycle.on_expiry = ExpiryPolicy::Cancel;
1707 lifecycle.drain_deadline = Some(now);
1708 assert!(lifecycle.validate().is_err());
1709
1710 let mut command = ControlCommand {
1711 scope: ControlScope::Instance,
1712 scope_id: "instance".into(),
1713 mode: ControlMode::Paused,
1714 operator_subject_id: "operator".parse().unwrap(),
1715 reason: "maintenance".into(),
1716 override_expires_at: None,
1717 };
1718 assert!(command.validate().is_ok());
1719 command.reason.clear();
1720 assert!(command.validate().is_err());
1721 }
1722
1723 #[test]
1724 fn funds_intent_requires_reservation_scope_and_retry_class() {
1725 let now = Utc::now();
1726 let mut intent = ActionIntent {
1727 id: "intent".into(),
1728 tenant_id: "tenant".parse().unwrap(),
1729 instance_id: "instance".parse().unwrap(),
1730 run_id: "run".parse().unwrap(),
1731 capability: CapabilityPin {
1732 id: "action.example".into(),
1733 contract_version: "1".into(),
1734 content_digest: "digest".into(),
1735 },
1736 idempotency_key: "idempotency".into(),
1737 state: ActionState::Prepared,
1738 input: Value::Null,
1739 effect: Effect::Funds,
1740 retry_class: IdempotencyMode::ReconcileBeforeRetry,
1741 control_epochs: ControlEpochs::default(),
1742 resource_scope_id: "resource".into(),
1743 lease_epoch: 1,
1744 action_epoch: 1,
1745 deadline: None,
1746 reservation: Some(ResourceReservationRef {
1747 reservation_id: "reservation".into(),
1748 fencing_token: 1,
1749 }),
1750 created_at: now,
1751 };
1752 assert!(intent.validate().is_ok());
1753 assert!(intent.validate_prepared().is_ok());
1754 intent.reservation = None;
1755 assert!(intent.validate().is_err());
1756 intent.reservation = Some(ResourceReservationRef {
1757 reservation_id: "reservation".into(),
1758 fencing_token: 1,
1759 });
1760 intent.resource_scope_id.clear();
1761 assert!(intent.validate().is_err());
1762 intent.resource_scope_id = "resource".into();
1763 intent.retry_class = IdempotencyMode::None;
1764 assert!(intent.validate().is_err());
1765 intent.retry_class = IdempotencyMode::ReconcileBeforeRetry;
1766 intent.state = ActionState::Authorized;
1767 assert!(intent.validate_prepared().is_err());
1768 assert!(!ActionState::Prepared.dispatch_committed());
1769 assert!(ActionState::Unknown.dispatch_committed());
1770 }
1771
1772 fn trigger_binding() -> TriggerBinding {
1773 TriggerBinding {
1774 id: "binding".into(),
1775 revision: 1,
1776 source: "source".into(),
1777 event_type: "event".into(),
1778 instance_id: "instance".parse().unwrap(),
1779 branch_id: None,
1780 predicate: serde_json::json!({}),
1781 ordering: OrderingPolicy::Commutative,
1782 starts_at: None,
1783 expires_at: None,
1784 gap_wait_ms: default_gap_wait_ms(),
1785 gap_limit: default_gap_limit(),
1786 }
1787 }
1788
1789 #[test]
1790 fn trigger_binding_rejects_blank_ids_bad_gaps_and_inverted_windows() {
1791 assert!(trigger_binding().validate().is_ok());
1792 let blank = TriggerBinding {
1793 source: " ".into(),
1794 ..trigger_binding()
1795 };
1796 assert!(blank.validate().is_err());
1797 for gap_wait_ms in [0, 3_600_001] {
1798 let binding = TriggerBinding {
1799 gap_wait_ms,
1800 ..trigger_binding()
1801 };
1802 assert!(binding.validate().is_err(), "gap_wait_ms {gap_wait_ms}");
1803 }
1804 for gap_limit in [0, 10_001] {
1805 let binding = TriggerBinding {
1806 gap_limit,
1807 ..trigger_binding()
1808 };
1809 assert!(binding.validate().is_err(), "gap_limit {gap_limit}");
1810 }
1811 let array_predicate = TriggerBinding {
1812 predicate: serde_json::json!([1]),
1813 ..trigger_binding()
1814 };
1815 assert!(array_predicate.validate().is_err());
1816 let now = Utc::now();
1817 let inverted = TriggerBinding {
1818 starts_at: Some(now),
1819 expires_at: Some(now),
1820 ..trigger_binding()
1821 };
1822 assert!(inverted.validate().is_err());
1823 let ordered = TriggerBinding {
1824 starts_at: Some(now),
1825 expires_at: Some(now + chrono::Duration::seconds(1)),
1826 ..trigger_binding()
1827 };
1828 assert!(ordered.validate().is_ok());
1829 let defaults: TriggerBinding = serde_json::from_value(serde_json::json!({
1830 "id": "binding", "revision": 1, "source": "s", "event_type": "e",
1831 "instance_id": "i", "predicate": {}, "ordering": "commutative",
1832 "starts_at": null, "expires_at": null
1833 }))
1834 .unwrap();
1835 assert_eq!(defaults.gap_wait_ms, 30_000);
1836 assert_eq!(defaults.gap_limit, 100);
1837 }
1838
1839 #[test]
1840 fn execution_profile_revision_requires_every_provider_name() {
1841 let profile = ExecutionProfileRevision {
1842 id: "profile".into(),
1843 revision: 1,
1844 content_digest: "sha256:profile".into(),
1845 mode: ExecutionMode::Paper,
1846 durability_grade: DurabilityGrade::Standard,
1847 trigger_provider: "triggers".into(),
1848 data_provider: "data".into(),
1849 clock_model: "database".into(),
1850 action_provider: "actions".into(),
1851 models: BTreeMap::new(),
1852 environment: Value::Null,
1853 policy_bundle: Value::Null,
1854 connection_bindings: BTreeMap::new(),
1855 };
1856 assert!(profile.validate().is_ok());
1857 let blank_provider = ExecutionProfileRevision {
1858 action_provider: String::new(),
1859 ..profile.clone()
1860 };
1861 assert!(blank_provider.validate().is_err());
1862 let blank_digest = ExecutionProfileRevision {
1863 content_digest: " ".into(),
1864 ..profile
1865 };
1866 assert!(blank_digest.validate().is_err());
1867 }
1868
1869 #[test]
1870 fn retry_backoff_grows_exponentially_with_bounded_jitter_and_cap() {
1871 let policy = RetryPolicy {
1872 max_attempts: 5,
1873 timeout_ms: 1_000,
1874 initial_backoff_ms: 100,
1875 max_backoff_ms: 1_000,
1876 };
1877 assert_eq!(policy.backoff_ms(1, 0), 100);
1878 assert_eq!(policy.backoff_ms(2, 0), 200);
1879 assert_eq!(policy.backoff_ms(3, 0), 400);
1880 assert_eq!(policy.backoff_ms(1, 24), 124);
1882 assert_eq!(policy.backoff_ms(1, 25), 100);
1883 assert_eq!(policy.backoff_ms(4, 249), 849);
1885 assert_eq!(policy.backoff_ms(5, 0), 1_000);
1887 assert_eq!(policy.backoff_ms(5, 249), 1_000);
1888 assert_eq!(policy.backoff_ms(40, u64::MAX), 1_000);
1889 assert_eq!(policy.backoff_ms(0, 0), 100);
1891 }
1892
1893 #[test]
1894 fn lifecycle_policy_rejects_inconsistent_windows_and_zero_counts() {
1895 let now = Utc::now();
1896 let later = now + chrono::Duration::hours(1);
1897 assert!(LifecyclePolicy::run_once().validate().is_ok());
1898 let inverted = LifecyclePolicy {
1899 starts_at: Some(later),
1900 expires_at: Some(now),
1901 ..LifecyclePolicy::run_once()
1902 };
1903 assert!(inverted.validate().is_err());
1904 let drain_without_policy = LifecyclePolicy {
1905 on_expiry: ExpiryPolicy::Cancel,
1906 drain_deadline: Some(later),
1907 ..LifecyclePolicy::run_once()
1908 };
1909 assert!(drain_without_policy.validate().is_err());
1910 let drain_before_expiry = LifecyclePolicy {
1911 expires_at: Some(later),
1912 drain_deadline: Some(now),
1913 ..LifecyclePolicy::run_once()
1914 };
1915 assert!(drain_before_expiry.validate().is_err());
1916 let zero_timeout = LifecyclePolicy {
1917 event_idle_timeout_ms: Some(0),
1918 ..LifecyclePolicy::run_once()
1919 };
1920 assert!(zero_timeout.validate().is_err());
1921 for completion in [
1922 CompletionPolicy::AfterMatchedEvaluations(0),
1923 CompletionPolicy::AfterSuccessfulRuns(0),
1924 ] {
1925 let zero_count = LifecyclePolicy {
1926 completion,
1927 ..LifecyclePolicy::run_once()
1928 };
1929 assert!(zero_count.validate().is_err());
1930 }
1931 let bounded = LifecyclePolicy {
1932 starts_at: Some(now),
1933 expires_at: Some(later),
1934 completion: CompletionPolicy::AfterSuccessfulRuns(2),
1935 event_idle_timeout_ms: Some(1_000),
1936 progress_timeout_ms: Some(1_000),
1937 on_expiry: ExpiryPolicy::Drain,
1938 drain_deadline: Some(later + chrono::Duration::minutes(1)),
1939 };
1940 assert!(bounded.validate().is_ok());
1941 }
1942}