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 pub predicate: Value,
67 pub ordering: OrderingPolicy,
69 pub starts_at: Option<DateTime<Utc>>,
71 pub expires_at: Option<DateTime<Utc>>,
73 #[serde(default = "default_gap_wait_ms")]
75 pub gap_wait_ms: u64,
76 #[serde(default = "default_gap_limit")]
78 pub gap_limit: u32,
79}
80
81impl TriggerBinding {
82 pub fn validate(&self) -> Result<(), ContractError> {
84 for (name, value) in [
85 ("trigger binding id", self.id.as_str()),
86 ("trigger source", self.source.as_str()),
87 ("trigger event_type", self.event_type.as_str()),
88 ("trigger instance_id", self.instance_id.as_str()),
89 ] {
90 required(name, value)?;
91 }
92 if self.gap_wait_ms == 0 || self.gap_wait_ms > 3_600_000 {
93 return Err(ContractError::Invalid(
94 "trigger gap_wait_ms must be between 1 and 3600000".into(),
95 ));
96 }
97 if self.gap_limit == 0 || self.gap_limit > 10_000 {
98 return Err(ContractError::Invalid(
99 "trigger gap_limit must be between 1 and 10000".into(),
100 ));
101 }
102 if !self.predicate.is_object() {
103 return Err(ContractError::Invalid(
104 "trigger predicate must be a JSON object".into(),
105 ));
106 }
107 if matches!((self.starts_at, self.expires_at), (Some(start), Some(end)) if end <= start) {
108 return Err(ContractError::Invalid(
109 "trigger expires_at must be after starts_at".into(),
110 ));
111 }
112 Ok(())
113 }
114}
115
116const fn default_gap_wait_ms() -> u64 {
117 30_000
118}
119
120const fn default_gap_limit() -> u32 {
121 100
122}
123
124#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
126pub struct WorkflowEvent {
127 pub instance_id: InstanceId,
129 pub sequence: i64,
131 pub event_type: String,
133 pub payload: Value,
135 pub content_digest: String,
137 pub occurred_at: DateTime<Utc>,
139}
140
141#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
143#[serde(rename_all = "snake_case")]
144pub enum ExecutionMode {
145 Simulation,
147 Backtest,
149 Paper,
151 Live,
153}
154
155#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
157#[serde(rename_all = "snake_case")]
158pub enum DurabilityGrade {
159 Standard,
161 FundsGrade,
163}
164
165#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
167pub struct ExecutionProfileRevision {
168 pub id: String,
170 pub revision: u64,
172 pub content_digest: String,
174 pub mode: ExecutionMode,
176 pub durability_grade: DurabilityGrade,
178 pub trigger_provider: String,
180 pub data_provider: String,
182 pub clock_model: String,
184 pub action_provider: String,
186 #[serde(default)]
188 pub models: BTreeMap<String, Value>,
189 #[serde(default)]
191 pub environment: Value,
192 #[serde(default)]
194 pub policy_bundle: Value,
195 #[serde(default)]
197 pub connection_bindings: BTreeMap<String, String>,
198}
199
200impl ExecutionProfileRevision {
201 pub fn validate(&self) -> Result<(), ContractError> {
203 for (name, value) in [
204 ("execution profile id", self.id.as_str()),
205 (
206 "execution profile content_digest",
207 self.content_digest.as_str(),
208 ),
209 ("trigger_provider", self.trigger_provider.as_str()),
210 ("data_provider", self.data_provider.as_str()),
211 ("clock_model", self.clock_model.as_str()),
212 ("action_provider", self.action_provider.as_str()),
213 ] {
214 required(name, value)?;
215 }
216 Ok(())
217 }
218}
219
220#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
222pub struct CapabilityPin {
223 pub id: String,
225 pub contract_version: String,
227 pub content_digest: String,
229}
230
231#[derive(Debug, Clone, Serialize, Deserialize)]
233pub struct WorkflowRevision {
234 pub definition_id: String,
236 pub revision: u64,
238 pub content_digest: String,
240 pub kernel_abi_version: String,
242 pub dependency_set_digest: String,
244 #[serde(default)]
246 pub expression_versions: BTreeMap<String, String>,
247 #[serde(default)]
249 pub capabilities: Vec<CapabilityPin>,
250 #[serde(default)]
252 pub template_provenance: Value,
253 pub spec: crate::Spec,
255}
256
257impl WorkflowRevision {
258 pub fn validate(&self) -> Result<(), ContractError> {
260 required("definition_id", &self.definition_id)?;
261 required("content_digest", &self.content_digest)?;
262 required("kernel_abi_version", &self.kernel_abi_version)?;
263 required("dependency_set_digest", &self.dependency_set_digest)?;
264 self.spec
265 .validate_structure()
266 .map_err(|error| ContractError::Invalid(error.to_string()))?;
267 let unique = self
268 .capabilities
269 .iter()
270 .map(|pin| (&pin.id, &pin.contract_version))
271 .collect::<BTreeSet<_>>();
272 if unique.len() != self.capabilities.len() {
273 return Err(ContractError::Invalid(
274 "capability pins must be unique by id and contract version".into(),
275 ));
276 }
277 for pin in &self.capabilities {
278 required("capability id", &pin.id)?;
279 required("capability contract_version", &pin.contract_version)?;
280 required("capability content_digest", &pin.content_digest)?;
281 }
282 Ok(())
283 }
284}
285
286#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
288#[serde(rename_all = "snake_case")]
289pub enum CapabilityKind {
290 Trigger,
292 Expression,
294 Guard,
296 Action,
298 Sink,
300}
301
302#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
304#[serde(rename_all = "snake_case")]
305pub enum Effect {
306 Pure,
308 Read,
310 InternalWrite,
312 ExternalWrite,
314 Funds,
316}
317
318#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
320#[serde(rename_all = "snake_case")]
321pub enum IdempotencyMode {
322 None,
324 Native,
326 ReconcileBeforeRetry,
328 NeverAutomaticRetry,
330}
331
332#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
334#[serde(rename_all = "snake_case")]
335pub enum CapabilityLifecycle {
336 Installed,
338 Active,
340 Deprecated,
342 Disabled,
344 Unavailable,
346 EmergencyRevoked,
348}
349
350#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
352pub struct RetryPolicy {
353 pub max_attempts: u32,
355 pub timeout_ms: u64,
357 pub initial_backoff_ms: u64,
359 pub max_backoff_ms: u64,
361}
362
363impl RetryPolicy {
364 pub fn backoff_ms(&self, attempt: u32, jitter_seed: u64) -> u64 {
367 let factor = 1_u64
368 .checked_shl(attempt.saturating_sub(1).min(20))
369 .unwrap_or(u64::MAX);
370 let base = self
371 .initial_backoff_ms
372 .saturating_mul(factor)
373 .min(self.max_backoff_ms);
374 let jitter_ceiling = (base / 4).max(1);
375 base.saturating_add(jitter_seed % jitter_ceiling)
376 .min(self.max_backoff_ms)
377 }
378}
379
380#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
382pub struct CapabilityManifest {
383 pub id: String,
385 pub contract_version: String,
387 pub content_digest: String,
389 pub kind: CapabilityKind,
391 pub input_schema: Value,
393 pub output_schema: Value,
395 pub effect: Effect,
397 pub deterministic: bool,
399 pub idempotency_mode: IdempotencyMode,
401 pub retry: RetryPolicy,
403 #[serde(default)]
405 pub permissions: BTreeSet<String>,
406 #[serde(default)]
408 pub required_guards: BTreeSet<GuardKind>,
409 #[serde(default)]
411 pub taint_rules: Value,
412 #[serde(default)]
414 pub resource_cost: Value,
415 pub supports_simulation: bool,
417 pub supports_replay: bool,
419 pub supports_paper: bool,
421 pub supports_live: bool,
423 #[serde(default)]
425 pub clock_requirements: Value,
426 #[serde(default)]
428 pub data_requirements: Value,
429 pub supports_reconciliation: bool,
431 pub lifecycle: CapabilityLifecycle,
433}
434
435impl CapabilityManifest {
436 pub fn action(
438 id: impl Into<String>,
439 contract_version: impl Into<String>,
440 content_digest: impl Into<String>,
441 effect: Effect,
442 idempotency_mode: IdempotencyMode,
443 supports_reconciliation: bool,
444 ) -> Self {
445 let required_guards = match effect {
446 Effect::Funds => BTreeSet::from([
447 GuardKind::Authorization,
448 GuardKind::Freshness,
449 GuardKind::Reservation,
450 ]),
451 Effect::ExternalWrite => BTreeSet::from([GuardKind::Authorization]),
452 _ => BTreeSet::new(),
453 };
454 Self {
455 id: id.into(),
456 contract_version: contract_version.into(),
457 content_digest: content_digest.into(),
458 kind: CapabilityKind::Action,
459 input_schema: serde_json::json!({"type": "object"}),
460 output_schema: serde_json::json!({"type": "object"}),
461 effect,
462 deterministic: false,
463 idempotency_mode,
464 retry: RetryPolicy {
465 max_attempts: 1,
466 timeout_ms: 30_000,
467 initial_backoff_ms: 100,
468 max_backoff_ms: 5_000,
469 },
470 permissions: BTreeSet::new(),
471 required_guards,
472 taint_rules: Value::Null,
473 resource_cost: Value::Null,
474 supports_simulation: true,
475 supports_replay: false,
476 supports_paper: true,
477 supports_live: true,
478 clock_requirements: Value::Null,
479 data_requirements: Value::Null,
480 supports_reconciliation,
481 lifecycle: CapabilityLifecycle::Active,
482 }
483 }
484
485 pub fn validate(&self) -> Result<(), ContractError> {
487 required("capability id", &self.id)?;
488 required("contract_version", &self.contract_version)?;
489 required("content_digest", &self.content_digest)?;
490 if self.retry.max_attempts == 0
491 || self.retry.timeout_ms == 0
492 || self.retry.initial_backoff_ms == 0
493 || self.retry.max_backoff_ms < self.retry.initial_backoff_ms
494 {
495 return Err(ContractError::Invalid(
496 "retry attempts, timeout and backoff bounds are invalid".into(),
497 ));
498 }
499 if self.effect == Effect::Funds {
500 if self.kind != CapabilityKind::Action {
501 return Err(ContractError::Invalid(
502 "funds effect is only valid for action capabilities".into(),
503 ));
504 }
505 if !self.supports_reconciliation && self.idempotency_mode != IdempotencyMode::Native {
506 return Err(ContractError::Invalid(
507 "funds actions require native idempotency or reconciliation".into(),
508 ));
509 }
510 for guard in [
511 GuardKind::Authorization,
512 GuardKind::Freshness,
513 GuardKind::Reservation,
514 ] {
515 if !self.required_guards.contains(&guard) {
516 return Err(ContractError::Invalid(format!(
517 "funds action must require {guard:?} guard"
518 )));
519 }
520 }
521 }
522 Ok(())
523 }
524
525 pub fn can_start_new_work(&self) -> bool {
527 matches!(
528 self.lifecycle,
529 CapabilityLifecycle::Installed
530 | CapabilityLifecycle::Active
531 | CapabilityLifecycle::Deprecated
532 )
533 }
534}
535
536#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
538#[serde(rename_all = "snake_case")]
539pub enum GuardKind {
540 Authorization,
542 Freshness,
544 Reservation,
546 Policy,
548}
549
550#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
552#[serde(rename_all = "snake_case")]
553pub enum OrderingPolicy {
554 StrictSequence,
556 Commutative,
558 LatestStateReconcile,
560 RejectUnordered,
562}
563
564#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
566pub struct TriggerEnvelope {
567 pub event_id: String,
569 pub event_type: String,
571 pub source: String,
573 pub schema_version: String,
575 pub tenant_id: TenantId,
577 pub subject_id: SubjectId,
579 pub aggregate_id: String,
581 pub source_sequence: Option<i64>,
583 pub observed_version: Option<String>,
585 pub occurred_at: DateTime<Utc>,
587 pub received_at: DateTime<Utc>,
589 pub watermark: Option<DateTime<Utc>>,
591 pub correlation_key: String,
593 pub dedup_key: String,
595 pub cursor: Option<String>,
597 pub payload: Value,
599 #[serde(default)]
601 pub trace_context: Value,
602}
603
604impl TriggerEnvelope {
605 pub fn validate(&self) -> Result<(), ContractError> {
607 for (name, value) in [
608 ("event_id", self.event_id.as_str()),
609 ("event_type", self.event_type.as_str()),
610 ("source", self.source.as_str()),
611 ("schema_version", self.schema_version.as_str()),
612 ("tenant_id", self.tenant_id.as_str()),
613 ("aggregate_id", self.aggregate_id.as_str()),
614 ("correlation_key", self.correlation_key.as_str()),
615 ("dedup_key", self.dedup_key.as_str()),
616 ] {
617 required(name, value)?;
618 }
619 if self.occurred_at > self.received_at + chrono::Duration::minutes(5) {
620 return Err(ContractError::Invalid(
621 "occurred_at is implausibly ahead of received_at".into(),
622 ));
623 }
624 Ok(())
625 }
626}
627
628#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
630#[serde(rename_all = "snake_case")]
631pub enum CompletionPolicy {
632 ExplicitStop,
634 FirstTrigger,
636 FirstMatch,
638 FirstActionTerminal,
640 FirstSuccess,
642 AfterMatchedEvaluations(u64),
644 AfterSuccessfulRuns(u64),
646}
647
648#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
650#[serde(rename_all = "snake_case")]
651pub enum ExpiryPolicy {
652 Drain,
654 Cancel,
656}
657
658#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
660#[serde(rename_all = "snake_case")]
661pub enum ScheduleCadence {
662 Cron {
664 expression: String,
666 },
667 FixedRate {
669 milliseconds: u64,
671 },
672 FixedDelay {
674 milliseconds: u64,
676 },
677}
678
679#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
681#[serde(rename_all = "snake_case")]
682pub enum CatchUpPolicy {
683 Skip,
685 CatchUpOnce,
687 CatchUpAll {
689 limit: u32,
691 },
692}
693
694#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
696pub struct SchedulePolicy {
697 pub cadence: ScheduleCadence,
699 pub timezone: String,
701 pub catch_up: CatchUpPolicy,
703}
704
705impl SchedulePolicy {
706 pub fn validate(&self) -> Result<(), ContractError> {
708 self.timezone.parse::<chrono_tz::Tz>().map_err(|_| {
709 ContractError::Invalid(format!("unknown IANA timezone '{}'", self.timezone))
710 })?;
711 match &self.cadence {
712 ScheduleCadence::Cron { expression } => {
713 expression.parse::<cron::Schedule>().map_err(|error| {
714 ContractError::Invalid(format!("invalid cron '{expression}': {error}"))
715 })?;
716 }
717 ScheduleCadence::FixedRate { milliseconds }
718 | ScheduleCadence::FixedDelay { milliseconds }
719 if *milliseconds == 0 =>
720 {
721 return Err(ContractError::Invalid(
722 "schedule interval must be positive".into(),
723 ));
724 }
725 _ => {}
726 }
727 if matches!(self.catch_up, CatchUpPolicy::CatchUpAll { limit: 0 }) {
728 return Err(ContractError::Invalid(
729 "catch_up_all limit must be positive".into(),
730 ));
731 }
732 Ok(())
733 }
734}
735
736#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
738pub struct LifecyclePolicy {
739 pub starts_at: Option<DateTime<Utc>>,
741 pub expires_at: Option<DateTime<Utc>>,
743 pub completion: CompletionPolicy,
745 pub event_idle_timeout_ms: Option<u64>,
747 pub progress_timeout_ms: Option<u64>,
749 pub on_expiry: ExpiryPolicy,
751 pub drain_deadline: Option<DateTime<Utc>>,
753}
754
755impl LifecyclePolicy {
756 pub fn run_once() -> Self {
758 Self {
759 starts_at: None,
760 expires_at: None,
761 completion: CompletionPolicy::FirstSuccess,
762 event_idle_timeout_ms: None,
763 progress_timeout_ms: None,
764 on_expiry: ExpiryPolicy::Drain,
765 drain_deadline: None,
766 }
767 }
768
769 pub fn validate(&self) -> Result<(), ContractError> {
771 if self
772 .starts_at
773 .zip(self.expires_at)
774 .is_some_and(|(a, b)| a >= b)
775 {
776 return Err(ContractError::Invalid(
777 "lifecycle starts_at must be before expires_at".into(),
778 ));
779 }
780 if self.drain_deadline.is_some() && self.on_expiry != ExpiryPolicy::Drain {
781 return Err(ContractError::Invalid(
782 "drain_deadline requires on_expiry=drain".into(),
783 ));
784 }
785 if self
786 .expires_at
787 .zip(self.drain_deadline)
788 .is_some_and(|(expires, drain)| drain <= expires)
789 {
790 return Err(ContractError::Invalid(
791 "lifecycle drain_deadline must be after expires_at".into(),
792 ));
793 }
794 if self.event_idle_timeout_ms == Some(0) || self.progress_timeout_ms == Some(0) {
795 return Err(ContractError::Invalid(
796 "lifecycle timeouts must be positive".into(),
797 ));
798 }
799 if matches!(
800 self.completion,
801 CompletionPolicy::AfterMatchedEvaluations(0) | CompletionPolicy::AfterSuccessfulRuns(0)
802 ) {
803 return Err(ContractError::Invalid(
804 "lifecycle completion count must be positive".into(),
805 ));
806 }
807 Ok(())
808 }
809}
810
811#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
813#[serde(rename_all = "snake_case")]
814pub enum StepOutcomeKind {
815 Succeeded,
817 Skipped,
819 Waiting,
821 Failed,
823 Cancelled,
825}
826
827#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
829pub struct StepOutcome {
830 pub kind: StepOutcomeKind,
832 pub reason_code: Option<String>,
834 #[serde(default)]
836 pub wake_condition: Value,
837}
838
839#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
841#[serde(rename_all = "snake_case")]
842pub enum ActionState {
843 Prepared,
845 AwaitingConfirmation,
847 Authorized,
849 DispatchCommitted,
851 Executing,
853 Succeeded,
855 Rejected,
857 Retryable,
859 Unknown,
861 Reconciled,
863}
864
865impl ActionState {
866 pub fn can_transition_to(self, next: Self) -> bool {
868 use ActionState::*;
869 matches!(
870 (self, next),
871 (Prepared, AwaitingConfirmation | Authorized | Rejected)
872 | (AwaitingConfirmation, Authorized | Rejected)
873 | (Authorized, DispatchCommitted | Rejected)
874 | (
875 DispatchCommitted,
876 Executing | Succeeded | Rejected | Retryable | Unknown
877 )
878 | (Executing, Succeeded | Rejected | Retryable | Unknown)
879 | (Retryable, DispatchCommitted | Reconciled)
880 | (Unknown, Reconciled)
881 | (Succeeded | Rejected, Reconciled)
882 )
883 }
884
885 pub fn dispatch_committed(self) -> bool {
887 matches!(
888 self,
889 Self::DispatchCommitted
890 | Self::Executing
891 | Self::Succeeded
892 | Self::Rejected
893 | Self::Retryable
894 | Self::Unknown
895 | Self::Reconciled
896 )
897 }
898}
899
900#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
902pub struct ControlEpochs {
903 pub tenant: i64,
905 pub resource: i64,
907 pub instance: i64,
909}
910
911#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
913#[serde(rename_all = "snake_case")]
914pub enum ControlScope {
915 Tenant,
917 Resource,
919 Instance,
921}
922
923#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
925#[serde(rename_all = "snake_case")]
926pub enum ControlMode {
927 Running,
929 Paused,
931 Stopped,
933 EmergencyStopped,
935}
936
937#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
939pub struct ControlCommand {
940 pub scope: ControlScope,
942 pub scope_id: String,
944 pub mode: ControlMode,
946 pub operator_subject_id: SubjectId,
948 pub reason: String,
950 pub override_expires_at: Option<DateTime<Utc>>,
952}
953
954impl ControlCommand {
955 pub fn validate(&self) -> Result<(), ContractError> {
957 required("control scope_id", &self.scope_id)?;
958 required("control operator_subject_id", &self.operator_subject_id)?;
959 required("control reason", &self.reason)?;
960 Ok(())
961 }
962}
963
964#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
966pub struct ResourceReservationRef {
967 pub reservation_id: String,
969 pub fencing_token: i64,
971}
972
973#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
975pub struct ActionIntent {
976 pub id: String,
978 pub tenant_id: TenantId,
980 pub instance_id: InstanceId,
982 pub run_id: RunId,
984 pub capability: CapabilityPin,
986 pub idempotency_key: String,
988 pub state: ActionState,
990 pub input: Value,
992 pub effect: Effect,
994 pub retry_class: IdempotencyMode,
996 pub control_epochs: ControlEpochs,
998 pub resource_scope_id: String,
1000 pub lease_epoch: i64,
1002 pub action_epoch: i64,
1004 pub deadline: Option<DateTime<Utc>>,
1006 pub reservation: Option<ResourceReservationRef>,
1008 pub created_at: DateTime<Utc>,
1010}
1011
1012impl ActionIntent {
1013 pub fn validate(&self) -> Result<(), ContractError> {
1015 required("action id", &self.id)?;
1016 required("action idempotency_key", &self.idempotency_key)?;
1017 if self.effect == Effect::Funds && self.reservation.is_none() {
1018 return Err(ContractError::Invalid(
1019 "funds action requires a resource reservation".into(),
1020 ));
1021 }
1022 if self.effect == Effect::Funds && self.resource_scope_id.trim().is_empty() {
1023 return Err(ContractError::Invalid(
1024 "funds action requires a resource scope".into(),
1025 ));
1026 }
1027 if self.effect == Effect::Funds && self.retry_class == IdempotencyMode::None {
1028 return Err(ContractError::Invalid(
1029 "funds action requires an explicit retry class".into(),
1030 ));
1031 }
1032 if self
1033 .deadline
1034 .is_some_and(|deadline| deadline <= self.created_at)
1035 {
1036 return Err(ContractError::Invalid(
1037 "action deadline must be after creation".into(),
1038 ));
1039 }
1040 Ok(())
1041 }
1042
1043 pub fn validate_prepared(&self) -> Result<(), ContractError> {
1046 self.validate()?;
1047 if self.state != ActionState::Prepared {
1048 return Err(ContractError::Invalid(
1049 "new action intent must start in prepared state".into(),
1050 ));
1051 }
1052 Ok(())
1053 }
1054}
1055
1056#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1058pub struct ActionReceipt {
1059 pub id: String,
1061 pub action_intent_id: String,
1063 pub provider_version: String,
1065 pub received_at: DateTime<Utc>,
1067 pub outcome: ActionState,
1069 pub payload: Value,
1071 pub raw_receipt_digest: String,
1073}
1074
1075#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1077pub struct ActionObservation {
1078 pub id: String,
1080 pub action_intent_id: String,
1082 pub provider_version: String,
1084 pub observed_at: DateTime<Utc>,
1086 pub state: String,
1088 pub resource_ref: Option<Value>,
1090 pub raw_receipt_digest: String,
1092 pub terminal: bool,
1094 #[serde(default)]
1097 pub retry_authorized: bool,
1098}
1099
1100#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1102pub struct SourceGap {
1103 pub tenant_id: TenantId,
1105 pub source: String,
1107 pub aggregate_id: String,
1109 pub missing_from: i64,
1111 pub missing_to: i64,
1113 pub deadline: DateTime<Utc>,
1115 pub status: String,
1117}
1118
1119#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1121#[serde(rename_all = "snake_case")]
1122pub enum ReservationScope {
1123 Tenant,
1125 Resource,
1127}
1128
1129#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1131#[serde(rename_all = "snake_case")]
1132pub enum ReservationState {
1133 Reserved,
1135 Consumed,
1137 Released,
1139 Expired,
1141}
1142
1143#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1145pub struct ResourceReservation {
1146 pub id: String,
1148 pub tenant_id: TenantId,
1150 pub scope: ReservationScope,
1152 pub scope_id: String,
1154 pub resource_kind: String,
1156 pub amount: String,
1158 pub policy_version: String,
1160 pub expires_at: DateTime<Utc>,
1162 pub fencing_token: i64,
1164 pub provider_reservation_id: Option<String>,
1166 pub state: ReservationState,
1168}
1169
1170#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1172pub struct ObservedValue<T> {
1173 pub value: T,
1175 pub source: String,
1177 pub observed_at: DateTime<Utc>,
1179 pub received_at: DateTime<Utc>,
1181 pub source_version: String,
1183 pub quality: String,
1185 pub digest: String,
1187}
1188
1189impl<T> ObservedValue<T> {
1190 pub fn is_accepted(
1192 &self,
1193 now: DateTime<Utc>,
1194 max_age: chrono::Duration,
1195 accepted_quality: &BTreeSet<String>,
1196 ) -> bool {
1197 self.observed_at <= now
1198 && now - self.observed_at <= max_age
1199 && accepted_quality.contains(&self.quality)
1200 }
1201}
1202
1203#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1205pub struct DecisionSnapshot {
1206 pub input_digests: Vec<String>,
1208 pub instance_config_version: String,
1210 pub policy_version: String,
1212 pub capability_versions: Vec<CapabilityPin>,
1214 pub execution_profile_revision_id: String,
1216 #[serde(default)]
1218 pub context_snapshot: Value,
1219 #[serde(default)]
1221 pub policy_snapshot: Value,
1222 #[serde(default)]
1224 pub agent_provenance: Value,
1225}
1226
1227#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1229#[serde(rename_all = "snake_case")]
1230pub enum ParameterProvenance {
1231 UserSupplied,
1233 ProductDefault,
1235 AgentInferred,
1237 Derived,
1239}
1240
1241#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1243pub struct ParameterValue {
1244 pub value: Value,
1246 pub provenance: ParameterProvenance,
1248}
1249
1250#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1252pub struct ParameterSpec {
1253 pub name: String,
1255 pub required: bool,
1257 pub required_explicit_for_live: bool,
1259}
1260
1261pub fn validate_parameters(
1263 mode: ExecutionMode,
1264 specs: &[ParameterSpec],
1265 values: &BTreeMap<String, ParameterValue>,
1266) -> Result<(), MissingRequirements> {
1267 let missing = specs
1268 .iter()
1269 .filter(|spec| match values.get(&spec.name) {
1270 None => {
1271 spec.required || (mode == ExecutionMode::Live && spec.required_explicit_for_live)
1272 }
1273 Some(value) => {
1274 mode == ExecutionMode::Live
1275 && spec.required_explicit_for_live
1276 && value.provenance != ParameterProvenance::UserSupplied
1277 }
1278 })
1279 .map(|spec| spec.name.clone())
1280 .collect::<Vec<_>>();
1281 if missing.is_empty() {
1282 Ok(())
1283 } else {
1284 Err(MissingRequirements {
1285 parameters: missing,
1286 })
1287 }
1288}
1289
1290#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, thiserror::Error)]
1292#[error("missing explicit workflow requirements: {parameters:?}")]
1293pub struct MissingRequirements {
1294 pub parameters: Vec<String>,
1296}
1297
1298#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1300pub struct RootOperationBudget {
1301 pub root_operation_id: String,
1303 pub max_depth: u32,
1305 pub descendant_limit: u32,
1307 pub run_limit: u32,
1309 pub token_budget: u64,
1311 pub cost_budget_micros: u64,
1313 pub action_budget: u32,
1315 pub deadline: DateTime<Utc>,
1317 #[serde(default)]
1319 pub permission_ceiling: BTreeSet<String>,
1320}
1321
1322#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1324pub struct DecisionArtifact {
1325 pub id: String,
1327 pub root_operation_id: String,
1329 pub content_digest: String,
1331 pub model: String,
1333 pub prompt_digest: String,
1335 pub output: Value,
1337 pub created_at: DateTime<Utc>,
1339}
1340
1341#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1343pub struct DiagnosticRecord {
1344 pub code: String,
1346 pub message: String,
1348 pub at: DateTime<Utc>,
1350 #[serde(default)]
1352 pub detail: Value,
1353}
1354
1355#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1357pub struct ExecutionProjection {
1358 pub schema_version: String,
1360 pub source_sequence: i64,
1362 pub current_steps: Vec<String>,
1364 pub completed_steps: Vec<String>,
1366 pub waiting_on: Option<Value>,
1368 pub last_decision: Option<Value>,
1370 pub next_possible_steps: Vec<String>,
1372 pub next_trigger_at: Option<DateTime<Utc>>,
1374 pub planned_actions: Vec<String>,
1376 pub latest_diagnostic: Option<DiagnosticRecord>,
1378 pub progress: Value,
1380 pub expires_at: Option<DateTime<Utc>>,
1382}
1383
1384#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1386pub enum ContractError {
1387 #[error("invalid workflow contract: {0}")]
1389 Invalid(String),
1390}
1391
1392fn required(name: &str, value: &str) -> Result<(), ContractError> {
1393 if value.trim().is_empty() {
1394 Err(ContractError::Invalid(format!("{name} is required")))
1395 } else {
1396 Ok(())
1397 }
1398}
1399
1400#[cfg(test)]
1401mod tests {
1402 use super::*;
1403
1404 fn empty_spec() -> crate::Spec {
1405 crate::Spec {
1406 spec_id: "test".into(),
1407 version: "1".into(),
1408 description: String::new(),
1409 aliases: vec![],
1410 instance_config_schema: None,
1411 display: BTreeMap::new(),
1412 branches: vec![],
1413 }
1414 }
1415
1416 fn funds_manifest() -> CapabilityManifest {
1417 CapabilityManifest::action(
1418 "action.example",
1419 "1",
1420 "sha256:x",
1421 Effect::Funds,
1422 IdempotencyMode::ReconcileBeforeRetry,
1423 true,
1424 )
1425 }
1426
1427 #[test]
1428 fn funds_capability_requires_reconciliation_and_all_guards() {
1429 assert!(funds_manifest().validate().is_ok());
1430 let mut invalid = funds_manifest();
1431 invalid.required_guards.remove(&GuardKind::Freshness);
1432 assert!(invalid.validate().is_err());
1433 invalid.required_guards.insert(GuardKind::Freshness);
1434 invalid.supports_reconciliation = false;
1435 assert!(invalid.validate().is_err());
1436 }
1437
1438 #[test]
1439 fn action_state_never_skips_dispatch_commit() {
1440 assert!(ActionState::Authorized.can_transition_to(ActionState::DispatchCommitted));
1441 assert!(!ActionState::Authorized.can_transition_to(ActionState::Succeeded));
1442 assert!(ActionState::Unknown.can_transition_to(ActionState::Reconciled));
1443 }
1444
1445 #[test]
1446 fn live_parameters_cannot_be_silently_inferred() {
1447 let specs = [ParameterSpec {
1448 name: "account".into(),
1449 required: true,
1450 required_explicit_for_live: true,
1451 }];
1452 let inferred = BTreeMap::from([(
1453 "account".into(),
1454 ParameterValue {
1455 value: Value::String("a".into()),
1456 provenance: ParameterProvenance::AgentInferred,
1457 },
1458 )]);
1459 assert_eq!(
1460 validate_parameters(ExecutionMode::Live, &specs, &inferred)
1461 .unwrap_err()
1462 .parameters,
1463 ["account"]
1464 );
1465 assert!(validate_parameters(ExecutionMode::Paper, &specs, &inferred).is_ok());
1466 }
1467
1468 #[test]
1469 fn freshness_is_fail_closed() {
1470 let now = Utc::now();
1471 let observed = ObservedValue {
1472 value: 1,
1473 source: "source".into(),
1474 observed_at: now - chrono::Duration::seconds(2),
1475 received_at: now,
1476 source_version: "1".into(),
1477 quality: "good".into(),
1478 digest: "d".into(),
1479 };
1480 assert!(observed.is_accepted(
1481 now,
1482 chrono::Duration::seconds(3),
1483 &BTreeSet::from(["good".into()])
1484 ));
1485 assert!(!observed.is_accepted(
1486 now,
1487 chrono::Duration::seconds(1),
1488 &BTreeSet::from(["good".into()])
1489 ));
1490 }
1491
1492 #[test]
1493 fn revision_and_manifest_contracts_fail_closed() {
1494 let mut revision = WorkflowRevision {
1495 definition_id: "definition".into(),
1496 revision: 1,
1497 content_digest: "digest".into(),
1498 kernel_abi_version: "1".into(),
1499 dependency_set_digest: "dependencies".into(),
1500 expression_versions: BTreeMap::new(),
1501 capabilities: vec![CapabilityPin {
1502 id: "action.example".into(),
1503 contract_version: "1".into(),
1504 content_digest: "capability-digest".into(),
1505 }],
1506 template_provenance: Value::Null,
1507 spec: empty_spec(),
1508 };
1509 assert!(revision.validate().is_ok());
1510 revision.capabilities.push(revision.capabilities[0].clone());
1511 assert!(revision.validate().is_err());
1512 revision.capabilities.pop();
1513 revision.kernel_abi_version.clear();
1514 assert!(revision.validate().is_err());
1515
1516 let external = CapabilityManifest::action(
1517 "action.notify",
1518 "1",
1519 "digest",
1520 Effect::ExternalWrite,
1521 IdempotencyMode::NeverAutomaticRetry,
1522 false,
1523 );
1524 assert_eq!(
1525 external.required_guards,
1526 BTreeSet::from([GuardKind::Authorization])
1527 );
1528 assert!(external.validate().is_ok());
1529 let mut invalid = external;
1530 invalid.retry.max_attempts = 0;
1531 assert!(invalid.validate().is_err());
1532
1533 let mut wrong_kind = funds_manifest();
1534 wrong_kind.kind = CapabilityKind::Expression;
1535 assert!(wrong_kind.validate().is_err());
1536 wrong_kind.kind = CapabilityKind::Action;
1537 wrong_kind.idempotency_mode = IdempotencyMode::Native;
1538 wrong_kind.supports_reconciliation = false;
1539 assert!(wrong_kind.validate().is_ok());
1540 wrong_kind.lifecycle = CapabilityLifecycle::Disabled;
1541 assert!(!wrong_kind.can_start_new_work());
1542 }
1543
1544 #[test]
1545 fn trigger_schedule_lifecycle_and_control_validate_boundaries() {
1546 let now = Utc::now();
1547 let mut trigger = TriggerEnvelope {
1548 event_id: "event".into(),
1549 event_type: "example".into(),
1550 source: "source".into(),
1551 schema_version: "1".into(),
1552 tenant_id: "tenant".parse().unwrap(),
1553 subject_id: "subject".parse().unwrap(),
1554 aggregate_id: "aggregate".into(),
1555 source_sequence: Some(1),
1556 observed_version: None,
1557 occurred_at: now,
1558 received_at: now,
1559 watermark: None,
1560 correlation_key: "key".into(),
1561 dedup_key: "dedup".into(),
1562 cursor: None,
1563 payload: Value::Null,
1564 trace_context: Value::Null,
1565 };
1566 assert!(trigger.validate().is_ok());
1567 trigger.occurred_at = now + chrono::Duration::minutes(6);
1568 assert!(trigger.validate().is_err());
1569
1570 let valid_schedule = SchedulePolicy {
1571 cadence: ScheduleCadence::Cron {
1572 expression: "0 0 * * * *".into(),
1573 },
1574 timezone: "Asia/Shanghai".into(),
1575 catch_up: CatchUpPolicy::CatchUpOnce,
1576 };
1577 assert!(valid_schedule.validate().is_ok());
1578 for invalid in [
1579 SchedulePolicy {
1580 timezone: "Nowhere/Invalid".into(),
1581 ..valid_schedule.clone()
1582 },
1583 SchedulePolicy {
1584 cadence: ScheduleCadence::FixedRate { milliseconds: 0 },
1585 ..valid_schedule.clone()
1586 },
1587 SchedulePolicy {
1588 catch_up: CatchUpPolicy::CatchUpAll { limit: 0 },
1589 ..valid_schedule
1590 },
1591 ] {
1592 assert!(invalid.validate().is_err());
1593 }
1594
1595 let mut lifecycle = LifecyclePolicy::run_once();
1596 assert!(lifecycle.validate().is_ok());
1597 lifecycle.starts_at = Some(now);
1598 lifecycle.expires_at = Some(now);
1599 assert!(lifecycle.validate().is_err());
1600 lifecycle.starts_at = None;
1601 lifecycle.expires_at = None;
1602 lifecycle.on_expiry = ExpiryPolicy::Cancel;
1603 lifecycle.drain_deadline = Some(now);
1604 assert!(lifecycle.validate().is_err());
1605
1606 let mut command = ControlCommand {
1607 scope: ControlScope::Instance,
1608 scope_id: "instance".into(),
1609 mode: ControlMode::Paused,
1610 operator_subject_id: "operator".parse().unwrap(),
1611 reason: "maintenance".into(),
1612 override_expires_at: None,
1613 };
1614 assert!(command.validate().is_ok());
1615 command.reason.clear();
1616 assert!(command.validate().is_err());
1617 }
1618
1619 #[test]
1620 fn funds_intent_requires_reservation_scope_and_retry_class() {
1621 let now = Utc::now();
1622 let mut intent = ActionIntent {
1623 id: "intent".into(),
1624 tenant_id: "tenant".parse().unwrap(),
1625 instance_id: "instance".parse().unwrap(),
1626 run_id: "run".parse().unwrap(),
1627 capability: CapabilityPin {
1628 id: "action.example".into(),
1629 contract_version: "1".into(),
1630 content_digest: "digest".into(),
1631 },
1632 idempotency_key: "idempotency".into(),
1633 state: ActionState::Prepared,
1634 input: Value::Null,
1635 effect: Effect::Funds,
1636 retry_class: IdempotencyMode::ReconcileBeforeRetry,
1637 control_epochs: ControlEpochs::default(),
1638 resource_scope_id: "resource".into(),
1639 lease_epoch: 1,
1640 action_epoch: 1,
1641 deadline: None,
1642 reservation: Some(ResourceReservationRef {
1643 reservation_id: "reservation".into(),
1644 fencing_token: 1,
1645 }),
1646 created_at: now,
1647 };
1648 assert!(intent.validate().is_ok());
1649 assert!(intent.validate_prepared().is_ok());
1650 intent.reservation = None;
1651 assert!(intent.validate().is_err());
1652 intent.reservation = Some(ResourceReservationRef {
1653 reservation_id: "reservation".into(),
1654 fencing_token: 1,
1655 });
1656 intent.resource_scope_id.clear();
1657 assert!(intent.validate().is_err());
1658 intent.resource_scope_id = "resource".into();
1659 intent.retry_class = IdempotencyMode::None;
1660 assert!(intent.validate().is_err());
1661 intent.retry_class = IdempotencyMode::ReconcileBeforeRetry;
1662 intent.state = ActionState::Authorized;
1663 assert!(intent.validate_prepared().is_err());
1664 assert!(!ActionState::Prepared.dispatch_committed());
1665 assert!(ActionState::Unknown.dispatch_committed());
1666 }
1667
1668 fn trigger_binding() -> TriggerBinding {
1669 TriggerBinding {
1670 id: "binding".into(),
1671 revision: 1,
1672 source: "source".into(),
1673 event_type: "event".into(),
1674 instance_id: "instance".parse().unwrap(),
1675 predicate: serde_json::json!({}),
1676 ordering: OrderingPolicy::Commutative,
1677 starts_at: None,
1678 expires_at: None,
1679 gap_wait_ms: default_gap_wait_ms(),
1680 gap_limit: default_gap_limit(),
1681 }
1682 }
1683
1684 #[test]
1685 fn trigger_binding_rejects_blank_ids_bad_gaps_and_inverted_windows() {
1686 assert!(trigger_binding().validate().is_ok());
1687 let blank = TriggerBinding {
1688 source: " ".into(),
1689 ..trigger_binding()
1690 };
1691 assert!(blank.validate().is_err());
1692 for gap_wait_ms in [0, 3_600_001] {
1693 let binding = TriggerBinding {
1694 gap_wait_ms,
1695 ..trigger_binding()
1696 };
1697 assert!(binding.validate().is_err(), "gap_wait_ms {gap_wait_ms}");
1698 }
1699 for gap_limit in [0, 10_001] {
1700 let binding = TriggerBinding {
1701 gap_limit,
1702 ..trigger_binding()
1703 };
1704 assert!(binding.validate().is_err(), "gap_limit {gap_limit}");
1705 }
1706 let array_predicate = TriggerBinding {
1707 predicate: serde_json::json!([1]),
1708 ..trigger_binding()
1709 };
1710 assert!(array_predicate.validate().is_err());
1711 let now = Utc::now();
1712 let inverted = TriggerBinding {
1713 starts_at: Some(now),
1714 expires_at: Some(now),
1715 ..trigger_binding()
1716 };
1717 assert!(inverted.validate().is_err());
1718 let ordered = TriggerBinding {
1719 starts_at: Some(now),
1720 expires_at: Some(now + chrono::Duration::seconds(1)),
1721 ..trigger_binding()
1722 };
1723 assert!(ordered.validate().is_ok());
1724 let defaults: TriggerBinding = serde_json::from_value(serde_json::json!({
1725 "id": "binding", "revision": 1, "source": "s", "event_type": "e",
1726 "instance_id": "i", "predicate": {}, "ordering": "commutative",
1727 "starts_at": null, "expires_at": null
1728 }))
1729 .unwrap();
1730 assert_eq!(defaults.gap_wait_ms, 30_000);
1731 assert_eq!(defaults.gap_limit, 100);
1732 }
1733
1734 #[test]
1735 fn execution_profile_revision_requires_every_provider_name() {
1736 let profile = ExecutionProfileRevision {
1737 id: "profile".into(),
1738 revision: 1,
1739 content_digest: "sha256:profile".into(),
1740 mode: ExecutionMode::Paper,
1741 durability_grade: DurabilityGrade::Standard,
1742 trigger_provider: "triggers".into(),
1743 data_provider: "data".into(),
1744 clock_model: "database".into(),
1745 action_provider: "actions".into(),
1746 models: BTreeMap::new(),
1747 environment: Value::Null,
1748 policy_bundle: Value::Null,
1749 connection_bindings: BTreeMap::new(),
1750 };
1751 assert!(profile.validate().is_ok());
1752 let blank_provider = ExecutionProfileRevision {
1753 action_provider: String::new(),
1754 ..profile.clone()
1755 };
1756 assert!(blank_provider.validate().is_err());
1757 let blank_digest = ExecutionProfileRevision {
1758 content_digest: " ".into(),
1759 ..profile
1760 };
1761 assert!(blank_digest.validate().is_err());
1762 }
1763
1764 #[test]
1765 fn retry_backoff_grows_exponentially_with_bounded_jitter_and_cap() {
1766 let policy = RetryPolicy {
1767 max_attempts: 5,
1768 timeout_ms: 1_000,
1769 initial_backoff_ms: 100,
1770 max_backoff_ms: 1_000,
1771 };
1772 assert_eq!(policy.backoff_ms(1, 0), 100);
1773 assert_eq!(policy.backoff_ms(2, 0), 200);
1774 assert_eq!(policy.backoff_ms(3, 0), 400);
1775 assert_eq!(policy.backoff_ms(1, 24), 124);
1777 assert_eq!(policy.backoff_ms(1, 25), 100);
1778 assert_eq!(policy.backoff_ms(4, 249), 849);
1780 assert_eq!(policy.backoff_ms(5, 0), 1_000);
1782 assert_eq!(policy.backoff_ms(5, 249), 1_000);
1783 assert_eq!(policy.backoff_ms(40, u64::MAX), 1_000);
1784 assert_eq!(policy.backoff_ms(0, 0), 100);
1786 }
1787
1788 #[test]
1789 fn lifecycle_policy_rejects_inconsistent_windows_and_zero_counts() {
1790 let now = Utc::now();
1791 let later = now + chrono::Duration::hours(1);
1792 assert!(LifecyclePolicy::run_once().validate().is_ok());
1793 let inverted = LifecyclePolicy {
1794 starts_at: Some(later),
1795 expires_at: Some(now),
1796 ..LifecyclePolicy::run_once()
1797 };
1798 assert!(inverted.validate().is_err());
1799 let drain_without_policy = LifecyclePolicy {
1800 on_expiry: ExpiryPolicy::Cancel,
1801 drain_deadline: Some(later),
1802 ..LifecyclePolicy::run_once()
1803 };
1804 assert!(drain_without_policy.validate().is_err());
1805 let drain_before_expiry = LifecyclePolicy {
1806 expires_at: Some(later),
1807 drain_deadline: Some(now),
1808 ..LifecyclePolicy::run_once()
1809 };
1810 assert!(drain_before_expiry.validate().is_err());
1811 let zero_timeout = LifecyclePolicy {
1812 event_idle_timeout_ms: Some(0),
1813 ..LifecyclePolicy::run_once()
1814 };
1815 assert!(zero_timeout.validate().is_err());
1816 for completion in [
1817 CompletionPolicy::AfterMatchedEvaluations(0),
1818 CompletionPolicy::AfterSuccessfulRuns(0),
1819 ] {
1820 let zero_count = LifecyclePolicy {
1821 completion,
1822 ..LifecyclePolicy::run_once()
1823 };
1824 assert!(zero_count.validate().is_err());
1825 }
1826 let bounded = LifecyclePolicy {
1827 starts_at: Some(now),
1828 expires_at: Some(later),
1829 completion: CompletionPolicy::AfterSuccessfulRuns(2),
1830 event_idle_timeout_ms: Some(1_000),
1831 progress_timeout_ms: Some(1_000),
1832 on_expiry: ExpiryPolicy::Drain,
1833 drain_deadline: Some(later + chrono::Duration::minutes(1)),
1834 };
1835 assert!(bounded.validate().is_ok());
1836 }
1837}