Skip to main content

af_workflow/
contract.rs

1//! Stable contracts for durable workflow execution.
2//!
3//! These types describe immutable definitions, accepted source facts and
4//! external effects. Products provide capabilities; the workflow kernel owns
5//! lifecycle, fencing and persistence semantics.
6
7use 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/// Named workflow whose revisions are immutable.
15#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
16pub struct WorkflowDefinition {
17    /// Stable identifier of this record.
18    pub id: String,
19    /// Display name.
20    pub name: String,
21}
22
23/// One durable instance pinned to a revision and execution profile for its lifetime.
24#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
25pub struct WorkflowInstance {
26    /// Stable identifier of this record.
27    pub id: String,
28    /// Tenant that owns this record.
29    pub tenant_id: TenantId,
30    /// Subject (user or service principal) acting on or owning this record.
31    pub subject_id: SubjectId,
32    /// Workflow definition this record belongs to.
33    pub definition_id: String,
34    /// Monotonic revision number.
35    pub revision: u64,
36    /// Execution profile the instance pins.
37    pub execution_profile_id: String,
38    /// Pinned execution profile revision.
39    pub execution_profile_revision: u64,
40    /// Lifecycle policy applied on database time.
41    pub lifecycle: LifecyclePolicy,
42    /// Current lifecycle status.
43    pub status: String,
44    /// CAS version of the authoritative state.
45    pub state_version: i64,
46    /// Sequence of the last authoritative event.
47    pub event_sequence: i64,
48    /// Instance control epoch; pause/resume bump it and stale workers fail.
49    pub control_epoch: i64,
50}
51
52/// Subscription of an instance to a source event type, with a JSON-containment predicate.
53#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
54pub struct TriggerBinding {
55    /// Stable identifier of this record.
56    pub id: String,
57    /// Monotonic revision number.
58    pub revision: u64,
59    /// Source the binding subscribes to.
60    pub source: String,
61    /// Stable machine-readable event type.
62    pub event_type: String,
63    /// Workflow instance this record refers to.
64    pub instance_id: InstanceId,
65    /// Target branch; omitted legacy requests resolve only for a single-branch spec.
66    #[serde(default)]
67    pub branch_id: Option<af_context::BranchId>,
68    /// JSON object the event payload must contain (`@>`) to create a delivery.
69    pub predicate: Value,
70    /// How out-of-order source sequences are handled.
71    pub ordering: OrderingPolicy,
72    /// Earliest time the binding or instance is active.
73    pub starts_at: Option<DateTime<Utc>>,
74    /// When the record stops being valid.
75    pub expires_at: Option<DateTime<Utc>>,
76    /// How long a missing source sequence may block before it is reported as a gap.
77    #[serde(default = "default_gap_wait_ms")]
78    pub gap_wait_ms: u64,
79    /// Maximum buffered out-of-order deliveries before the source is blocked.
80    #[serde(default = "default_gap_limit")]
81    pub gap_limit: u32,
82}
83
84impl TriggerBinding {
85    /// Reject blank identifiers, out-of-range gap settings and inverted windows.
86    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/// Authoritative, append-only fact recorded by a committed transition.
128#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
129pub struct WorkflowEvent {
130    /// Workflow instance this record refers to.
131    pub instance_id: InstanceId,
132    /// Position in the instance's event log.
133    pub sequence: i64,
134    /// Stable machine-readable event type.
135    pub event_type: String,
136    /// Structured payload.
137    pub payload: Value,
138    /// Content hash that makes the referenced artifact immutable.
139    pub content_digest: String,
140    /// When the event happened.
141    pub occurred_at: DateTime<Utc>,
142}
143
144/// Where an execution profile is allowed to act.
145#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
146#[serde(rename_all = "snake_case")]
147pub enum ExecutionMode {
148    /// In-memory dry run; no external effects.
149    Simulation,
150    /// Historical replay.
151    Backtest,
152    /// Live data, simulated effects.
153    Paper,
154    /// Real external effects.
155    Live,
156}
157
158/// Durability the deployment must provide for an execution profile.
159#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
160#[serde(rename_all = "snake_case")]
161pub enum DurabilityGrade {
162    /// Up to five minutes RPO for ordinary work.
163    Standard,
164    /// Committed intents synchronously preserved before dispatch; required for funds actions.
165    FundsGrade,
166}
167
168/// Immutable execution profile: mode, durability grade and provider bindings an instance pins.
169#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
170pub struct ExecutionProfileRevision {
171    /// Stable identifier of this record.
172    pub id: String,
173    /// Monotonic revision number.
174    pub revision: u64,
175    /// Content hash that makes the referenced artifact immutable.
176    pub content_digest: String,
177    /// Execution mode this record was produced under.
178    pub mode: ExecutionMode,
179    /// Durability the deployment guarantees for this profile.
180    pub durability_grade: DurabilityGrade,
181    /// Provider that feeds triggers.
182    pub trigger_provider: String,
183    /// Provider that feeds market or reference data.
184    pub data_provider: String,
185    /// Clock source; `database` is the only kernel-supported value.
186    pub clock_model: String,
187    /// Provider that dispatches actions.
188    pub action_provider: String,
189    /// Model bindings by role.
190    #[serde(default)]
191    pub models: BTreeMap<String, Value>,
192    /// Environment settings the product interprets.
193    #[serde(default)]
194    pub environment: Value,
195    /// Product policy bundle applied to permissions and guards.
196    #[serde(default)]
197    pub policy_bundle: Value,
198    /// Named external connections by role.
199    #[serde(default)]
200    pub connection_bindings: BTreeMap<String, String>,
201}
202
203impl ExecutionProfileRevision {
204    /// Reject blank identifiers and provider names.
205    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/// Exact capability identity: id, contract version and content digest.
224#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
225pub struct CapabilityPin {
226    /// Stable identifier of this record.
227    pub id: String,
228    /// Contract version of the capability.
229    pub contract_version: String,
230    /// Content hash that makes the referenced artifact immutable.
231    pub content_digest: String,
232}
233
234/// Immutable, digest-locked revision of a workflow spec with its pinned capabilities.
235#[derive(Debug, Clone, Serialize, Deserialize)]
236pub struct WorkflowRevision {
237    /// Workflow definition this record belongs to.
238    pub definition_id: String,
239    /// Monotonic revision number.
240    pub revision: u64,
241    /// Content hash that makes the referenced artifact immutable.
242    pub content_digest: String,
243    /// Kernel ABI the revision was validated against.
244    pub kernel_abi_version: String,
245    /// Digest over the pinned capabilities.
246    pub dependency_set_digest: String,
247    /// Versions of generic expressions used by the spec.
248    #[serde(default)]
249    pub expression_versions: BTreeMap<String, String>,
250    /// Capabilities the spec may dispatch.
251    #[serde(default)]
252    pub capabilities: Vec<CapabilityPin>,
253    /// Where the spec came from (template, draft command, author).
254    #[serde(default)]
255    pub template_provenance: Value,
256    /// The validated spec.
257    pub spec: crate::Spec,
258}
259
260impl WorkflowRevision {
261    /// Check this contract's invariants; returns the first violation.
262    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/// Role a capability plays in a graph.
290#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
291#[serde(rename_all = "snake_case")]
292pub enum CapabilityKind {
293    /// Produces events.
294    Trigger,
295    /// Pure or read-only transform.
296    Expression,
297    /// Authorization, freshness, reservation or policy check dominating an action.
298    Guard,
299    /// External effect dispatched through an intent.
300    Action,
301    /// Terminal consumer.
302    Sink,
303}
304
305/// Side-effect class of a capability.
306#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
307#[serde(rename_all = "snake_case")]
308pub enum Effect {
309    /// No I/O.
310    Pure,
311    /// Reads external state.
312    Read,
313    /// Writes kernel-owned state only.
314    InternalWrite,
315    /// Writes external state; requires authorization dominance.
316    ExternalWrite,
317    /// Moves value; requires authorization, freshness, reservation and funds-grade durability.
318    Funds,
319}
320
321/// Whether repeating a dispatch is safe.
322#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
323#[serde(rename_all = "snake_case")]
324pub enum IdempotencyMode {
325    /// Not idempotent; never retried automatically.
326    None,
327    /// The provider deduplicates by idempotency key; safe to retry.
328    Native,
329    /// Must reconcile the previous attempt before dispatching again.
330    ReconcileBeforeRetry,
331    /// Retry only through an explicit operator action.
332    NeverAutomaticRetry,
333}
334
335/// Deployment state of a capability provider.
336#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
337#[serde(rename_all = "snake_case")]
338pub enum CapabilityLifecycle {
339    /// Registered but not yet serving.
340    Installed,
341    /// Serving new and existing work.
342    Active,
343    /// Serving pinned work only; new pins rejected.
344    Deprecated,
345    /// Not serving; pinned work waits.
346    Disabled,
347    /// Temporarily unreachable.
348    Unavailable,
349    /// Blocks every action immediately.
350    EmergencyRevoked,
351}
352
353/// Automatic retry policy of a capability.
354#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
355pub struct RetryPolicy {
356    /// Total dispatch attempts allowed.
357    pub max_attempts: u32,
358    /// Per-attempt deadline.
359    pub timeout_ms: u64,
360    /// Backoff after the first failure.
361    pub initial_backoff_ms: u64,
362    /// Cap on exponential backoff.
363    pub max_backoff_ms: u64,
364}
365
366impl RetryPolicy {
367    /// Exponential backoff for the 1-based `attempt`, capped at `max_backoff_ms`
368    /// with up to 25% deterministic jitter derived from `jitter_seed`.
369    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/// Immutable contract of a trigger, expression, guard, action or sink.
384#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
385pub struct CapabilityManifest {
386    /// Stable identifier of this record.
387    pub id: String,
388    /// Contract version of the capability.
389    pub contract_version: String,
390    /// Content hash that makes the referenced artifact immutable.
391    pub content_digest: String,
392    /// Discriminator naming the variant of this record.
393    pub kind: CapabilityKind,
394    /// Role enforced by a guard capability.
395    #[serde(default, skip_serializing_if = "Option::is_none")]
396    pub guard_kind: Option<GuardKind>,
397    /// JSON Schema for node configuration exposed to authoring clients.
398    #[serde(default = "object_schema")]
399    pub config_schema: Value,
400    /// JSON Schema for provider input.
401    pub input_schema: Value,
402    /// JSON Schema for provider output.
403    pub output_schema: Value,
404    /// Side-effect class of the action.
405    pub effect: Effect,
406    /// Whether equal inputs always produce equal outputs.
407    pub deterministic: bool,
408    /// Retry safety class.
409    pub idempotency_mode: IdempotencyMode,
410    /// Retry policy.
411    pub retry: RetryPolicy,
412    /// Permissions granted or required.
413    #[serde(default)]
414    pub permissions: BTreeSet<String>,
415    /// Guards that must dominate the action.
416    #[serde(default)]
417    pub required_guards: BTreeSet<GuardKind>,
418    /// Exact guard capabilities that must each dominate, including guards of the same role.
419    #[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
420    pub required_guard_pins: BTreeSet<CapabilityPin>,
421    /// Which inputs may be LLM-tainted.
422    #[serde(default)]
423    pub taint_rules: Value,
424    /// Declared cost for budgeting.
425    #[serde(default)]
426    pub resource_cost: Value,
427    /// Usable in simulation.
428    pub supports_simulation: bool,
429    /// Usable in backtest replay.
430    pub supports_replay: bool,
431    /// Usable in paper mode.
432    pub supports_paper: bool,
433    /// Usable live.
434    pub supports_live: bool,
435    /// Clock guarantees the provider needs.
436    #[serde(default)]
437    pub clock_requirements: Value,
438    /// Data freshness the provider needs.
439    #[serde(default)]
440    pub data_requirements: Value,
441    /// Whether the provider can observe a dispatched action's outcome.
442    pub supports_reconciliation: bool,
443    /// Deployment state.
444    pub lifecycle: CapabilityLifecycle,
445}
446
447fn object_schema() -> Value {
448    serde_json::json!({"type":"object"})
449}
450
451impl CapabilityManifest {
452    /// Manifest for an action capability with conservative defaults (one attempt, 30 s timeout).
453    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    /// Check this contract's invariants; returns the first violation.
505    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    /// Whether new intents may pin this provider (`installed` or `active`).
570    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/// Role a guard plays on an action's dominating path.
581#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
582#[serde(rename_all = "snake_case")]
583pub enum GuardKind {
584    /// Caller may perform the effect.
585    Authorization,
586    /// Inputs are recent enough.
587    Freshness,
588    /// Resources are reserved and fenced.
589    Reservation,
590    /// Product policy allows it.
591    Policy,
592}
593
594/// How a binding treats source sequence order.
595#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
596#[serde(rename_all = "snake_case")]
597pub enum OrderingPolicy {
598    /// Reject a gap; the next accepted sequence must be previous + 1.
599    StrictSequence,
600    /// Order does not matter.
601    Commutative,
602    /// Only the latest event matters.
603    LatestStateReconcile,
604    /// Reject anything older than the last accepted sequence.
605    RejectUnordered,
606}
607
608/// Immutable source event as delivered by a trigger adapter.
609#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
610pub struct TriggerEnvelope {
611    /// Source-unique event id.
612    pub event_id: String,
613    /// Stable machine-readable event type.
614    pub event_type: String,
615    /// Trigger adapter that delivered the event.
616    pub source: String,
617    /// Projection schema version.
618    pub schema_version: String,
619    /// Tenant that owns this record.
620    pub tenant_id: TenantId,
621    /// Subject (user or service principal) acting on or owning this record.
622    pub subject_id: SubjectId,
623    /// Aggregate the event belongs to (for ordering).
624    pub aggregate_id: String,
625    /// Per-aggregate sequence from the source.
626    pub source_sequence: Option<i64>,
627    /// Version of the observed object, if any.
628    pub observed_version: Option<String>,
629    /// When the event happened.
630    pub occurred_at: DateTime<Utc>,
631    /// When the receipt arrived.
632    pub received_at: DateTime<Utc>,
633    /// Source watermark up to which events are complete.
634    pub watermark: Option<DateTime<Utc>>,
635    /// Key that correlates related events.
636    pub correlation_key: String,
637    /// Key that deduplicates redeliveries.
638    pub dedup_key: String,
639    /// Source cursor after this event.
640    pub cursor: Option<String>,
641    /// Structured payload.
642    pub payload: Value,
643    /// Distributed tracing context.
644    #[serde(default)]
645    pub trace_context: Value,
646}
647
648impl TriggerEnvelope {
649    /// Check this contract's invariants; returns the first violation.
650    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/// Cursor to resume an admitted source adapter after restart.
682#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
683pub struct SourceAdapterState {
684    /// Exact source.
685    pub source: af_context::WorkflowSourceId,
686    /// Opaque provider cursor.
687    pub cursor: String,
688    /// Event whose transaction last advanced the cursor; absent for an initial seed.
689    pub source_event_id: Option<String>,
690}
691/// Receipt returned only after source fact, binding snapshot, deliveries and cursor commit.
692#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
693pub struct SourceIngestionReceipt {
694    /// Durable event identity safe for adapter ACK.
695    pub event_id: String,
696    /// True for an identical committed retry.
697    pub duplicate: bool,
698    /// Newly inserted deliveries; duplicate retries report zero.
699    pub deliveries: u64,
700    /// Current committed resume cursor, including one advanced after this event.
701    pub committed_cursor: Option<String>,
702}
703
704/// When an instance completes on its own.
705#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
706#[serde(rename_all = "snake_case")]
707pub enum CompletionPolicy {
708    /// Only an operator stops it.
709    ExplicitStop,
710    /// After the first consumed trigger or timer.
711    FirstTrigger,
712    /// After the first evaluation that matched.
713    FirstMatch,
714    /// After the first action reaches a terminal state.
715    FirstActionTerminal,
716    /// After the first successful evaluation.
717    FirstSuccess,
718    /// After this many matched evaluations.
719    AfterMatchedEvaluations(u64),
720    /// After this many successful evaluations.
721    AfterSuccessfulRuns(u64),
722}
723
724/// What happens when a lifecycle deadline passes.
725#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
726#[serde(rename_all = "snake_case")]
727pub enum ExpiryPolicy {
728    /// Stop new work; let dispatched actions finish.
729    Drain,
730    /// Stop and request cancellation.
731    Cancel,
732}
733
734/// How a schedule fires.
735#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
736#[serde(rename_all = "snake_case")]
737pub enum ScheduleCadence {
738    /// Cron expression evaluated in the schedule's timezone.
739    Cron {
740        /// Six- or seven-field cron expression.
741        expression: String,
742    },
743    /// Fixed period measured from the scheduled time.
744    FixedRate {
745        /// Period in milliseconds.
746        milliseconds: u64,
747    },
748    /// Fixed delay measured from the previous completion.
749    FixedDelay {
750        /// Delay in milliseconds.
751        milliseconds: u64,
752    },
753}
754
755/// How missed schedule windows are handled.
756#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
757#[serde(rename_all = "snake_case")]
758pub enum CatchUpPolicy {
759    /// Run only the most recent due occurrence.
760    Skip,
761    /// Run once, listing the missed windows it covers.
762    CatchUpOnce,
763    /// Run every missed occurrence oldest first, at most `limit` per transition.
764    CatchUpAll {
765        /// Maximum occurrences per transition.
766        limit: u32,
767    },
768}
769
770/// Cadence, timezone and catch-up behavior of a schedule.
771#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
772pub struct SchedulePolicy {
773    /// Cron, fixed-rate or fixed-delay.
774    pub cadence: ScheduleCadence,
775    /// IANA timezone for cron evaluation and DST.
776    pub timezone: String,
777    /// Missed-window behavior.
778    pub catch_up: CatchUpPolicy,
779}
780
781impl SchedulePolicy {
782    /// Check this contract's invariants; returns the first violation.
783    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/// Start, completion, timeout and expiry rules evaluated on database time.
813#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
814pub struct LifecyclePolicy {
815    /// Earliest time the binding or instance is active.
816    pub starts_at: Option<DateTime<Utc>>,
817    /// When the record stops being valid.
818    pub expires_at: Option<DateTime<Utc>>,
819    /// When the instance completes.
820    pub completion: CompletionPolicy,
821    /// Terminalize after this long without a consumed event.
822    pub event_idle_timeout_ms: Option<u64>,
823    /// Terminalize after this long without a committed transition.
824    pub progress_timeout_ms: Option<u64>,
825    /// Drain or cancel on expiry.
826    pub on_expiry: ExpiryPolicy,
827    /// Hard stop for draining.
828    pub drain_deadline: Option<DateTime<Utc>>,
829}
830
831impl LifecyclePolicy {
832    /// One-shot lifecycle: completes on the first success, drains on expiry.
833    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    /// Check this contract's invariants; returns the first violation.
846    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/// How a step ended.
888#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
889#[serde(rename_all = "snake_case")]
890pub enum StepOutcomeKind {
891    /// Terminal success.
892    Succeeded,
893    /// Filtered out.
894    Skipped,
895    /// Parked on a wake condition.
896    Waiting,
897    /// Errored.
898    Failed,
899    /// Cancelled.
900    Cancelled,
901}
902
903/// Recorded outcome of one step.
904#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
905pub struct StepOutcome {
906    /// Discriminator naming the variant of this record.
907    pub kind: StepOutcomeKind,
908    /// Stable code for material outcomes.
909    pub reason_code: Option<String>,
910    /// What resumes a waiting step.
911    #[serde(default)]
912    pub wake_condition: Value,
913}
914
915/// Lifecycle of an action intent; `dispatch_committed` is the point of no return.
916#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
917#[serde(rename_all = "snake_case")]
918pub enum ActionState {
919    /// Created by a transition; not yet authorized.
920    Prepared,
921    /// Waiting for a human confirmation.
922    AwaitingConfirmation,
923    /// Guards passed; may be dispatched.
924    Authorized,
925    /// Effect committed; cancellation can no longer claim it was not performed.
926    DispatchCommitted,
927    /// Provider reports progress.
928    Executing,
929    /// Terminal success.
930    Succeeded,
931    /// Terminal rejection.
932    Rejected,
933    /// Failed; waits for its backoff, then dispatches again.
934    Retryable,
935    /// Outcome unknown; waits for reconciliation.
936    Unknown,
937    /// Outcome established by reconciliation.
938    Reconciled,
939}
940
941impl ActionState {
942    /// Whether the state machine allows `self -> next`.
943    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    /// Whether the effect may already have happened.
962    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/// Control epochs an intent was created under; any bump makes it stale.
977#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
978pub struct ControlEpochs {
979    /// Tenant control epoch.
980    pub tenant: i64,
981    /// Resource control epoch.
982    pub resource: i64,
983    /// Instance control epoch.
984    pub instance: i64,
985}
986
987/// Scope of a control command.
988#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
989#[serde(rename_all = "snake_case")]
990pub enum ControlScope {
991    /// Whole tenant.
992    Tenant,
993    /// One external resource.
994    Resource,
995    /// One instance.
996    Instance,
997}
998
999/// Operating mode set by a control command.
1000#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1001#[serde(rename_all = "snake_case")]
1002pub enum ControlMode {
1003    /// Normal operation.
1004    Running,
1005    /// No new evaluations or dispatches.
1006    Paused,
1007    /// Terminal stop.
1008    Stopped,
1009    /// Immediate stop that also revokes providers.
1010    EmergencyStopped,
1011}
1012
1013/// Operator command that bumps a control epoch.
1014#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1015pub struct ControlCommand {
1016    /// Scope kind.
1017    pub scope: ControlScope,
1018    /// Scope identity.
1019    pub scope_id: String,
1020    /// Execution mode this record was produced under.
1021    pub mode: ControlMode,
1022    /// Operator issuing the command.
1023    pub operator_subject_id: SubjectId,
1024    /// Human-readable reason.
1025    pub reason: String,
1026    /// When an override lapses.
1027    pub override_expires_at: Option<DateTime<Utc>>,
1028}
1029
1030impl ControlCommand {
1031    /// Check this contract's invariants; returns the first violation.
1032    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/// Reference to a product-owned fenced reservation.
1041#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1042pub struct ResourceReservationRef {
1043    /// Reservation identity.
1044    pub reservation_id: String,
1045    /// Fence the reservation was taken under.
1046    pub fencing_token: i64,
1047}
1048
1049/// Durable intent to perform one external effect, created inside a fenced transition.
1050#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1051pub struct ActionIntent {
1052    /// Stable identifier of this record.
1053    pub id: String,
1054    /// Tenant that owns this record.
1055    pub tenant_id: TenantId,
1056    /// Workflow instance this record refers to.
1057    pub instance_id: InstanceId,
1058    /// Run this record belongs to.
1059    pub run_id: RunId,
1060    /// Pinned capability (id, contract version, content digest).
1061    pub capability: CapabilityPin,
1062    /// Caller-supplied key that makes repeated submissions return the first result.
1063    pub idempotency_key: String,
1064    /// Current lifecycle state.
1065    pub state: ActionState,
1066    /// Structured input handed to the provider.
1067    pub input: Value,
1068    /// Side-effect class of the action.
1069    pub effect: Effect,
1070    /// Idempotency class that decides whether automatic retry is safe.
1071    pub retry_class: IdempotencyMode,
1072    /// Control epochs at creation.
1073    pub control_epochs: ControlEpochs,
1074    /// Resource the effect targets; empty when unscoped.
1075    pub resource_scope_id: String,
1076    /// Instance lease version at creation.
1077    pub lease_epoch: i64,
1078    /// Bumped on every claim; fences claim owners.
1079    pub action_epoch: i64,
1080    /// Latest time by which the work must finish.
1081    pub deadline: Option<DateTime<Utc>>,
1082    /// Fenced reservation for funds actions.
1083    pub reservation: Option<ResourceReservationRef>,
1084    /// When the record was created (database time).
1085    pub created_at: DateTime<Utc>,
1086}
1087
1088impl ActionIntent {
1089    /// Check this contract's invariants; returns the first violation.
1090    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    /// Validate an intent at the only creation boundary. Later states are
1120    /// reached exclusively through the fenced store transition.
1121    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/// Provider's answer to a dispatch.
1133#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1134pub struct ActionReceipt {
1135    /// Stable identifier of this record.
1136    pub id: String,
1137    /// Intent this refers to.
1138    pub action_intent_id: String,
1139    /// Provider version that produced it.
1140    pub provider_version: String,
1141    /// When the receipt arrived.
1142    pub received_at: DateTime<Utc>,
1143    /// State the provider reports.
1144    pub outcome: ActionState,
1145    /// Structured payload.
1146    pub payload: Value,
1147    /// Digest of the raw provider response; deduplicates receipts.
1148    pub raw_receipt_digest: String,
1149}
1150
1151/// Observed state of a dispatched action.
1152#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1153pub struct ActionObservation {
1154    /// Stable identifier of this record.
1155    pub id: String,
1156    /// Intent this refers to.
1157    pub action_intent_id: String,
1158    /// Provider version that produced it.
1159    pub provider_version: String,
1160    /// When the state was observed.
1161    pub observed_at: DateTime<Utc>,
1162    /// Current lifecycle state.
1163    pub state: String,
1164    /// External resource created or affected.
1165    pub resource_ref: Option<Value>,
1166    /// Digest of the raw provider response; deduplicates receipts.
1167    pub raw_receipt_digest: String,
1168    /// Whether this observation ends the action.
1169    pub terminal: bool,
1170    /// A reconciliation provider sets this only after proving that repeating
1171    /// the same idempotent action is safe.
1172    #[serde(default)]
1173    pub retry_authorized: bool,
1174}
1175
1176/// A missing source-sequence range that blocks ordered delivery for an aggregate.
1177#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1178pub struct SourceGap {
1179    /// Tenant that owns this record.
1180    pub tenant_id: TenantId,
1181    /// Source the gap belongs to.
1182    pub source: String,
1183    /// Aggregate whose sequence has the hole.
1184    pub aggregate_id: String,
1185    /// First missing sequence (inclusive).
1186    pub missing_from: i64,
1187    /// Last missing sequence (inclusive).
1188    pub missing_to: i64,
1189    /// When waiting for the gap expires and delivery blocks.
1190    pub deadline: DateTime<Utc>,
1191    /// Gap status (`waiting`, `blocked`, `unblocked`).
1192    pub status: String,
1193}
1194
1195/// Scope of a resource reservation.
1196#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1197#[serde(rename_all = "snake_case")]
1198pub enum ReservationScope {
1199    /// Whole tenant.
1200    Tenant,
1201    /// One external resource.
1202    Resource,
1203}
1204
1205/// Lifecycle of a reservation.
1206#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1207#[serde(rename_all = "snake_case")]
1208pub enum ReservationState {
1209    /// Held.
1210    Reserved,
1211    /// Used by a committed dispatch.
1212    Consumed,
1213    /// Given back.
1214    Released,
1215    /// Lapsed unused.
1216    Expired,
1217}
1218
1219/// Mirror of a product-owned fenced reservation.
1220#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1221pub struct ResourceReservation {
1222    /// Stable identifier of this record.
1223    pub id: String,
1224    /// Tenant that owns this record.
1225    pub tenant_id: TenantId,
1226    /// Scope kind.
1227    pub scope: ReservationScope,
1228    /// Scope identity.
1229    pub scope_id: String,
1230    /// Kind of resource reserved.
1231    pub resource_kind: String,
1232    /// Reserved amount as a decimal string.
1233    pub amount: String,
1234    /// Product policy version that granted it.
1235    pub policy_version: String,
1236    /// When the record stops being valid.
1237    pub expires_at: DateTime<Utc>,
1238    /// Fence the reservation was taken under.
1239    pub fencing_token: i64,
1240    /// Provider-side reservation id.
1241    pub provider_reservation_id: Option<String>,
1242    /// Current lifecycle state.
1243    pub state: ReservationState,
1244}
1245
1246/// A value with provenance and freshness metadata.
1247#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1248pub struct ObservedValue<T> {
1249    /// The value carried by this record.
1250    pub value: T,
1251    /// Where the value was observed.
1252    pub source: String,
1253    /// When the state was observed.
1254    pub observed_at: DateTime<Utc>,
1255    /// When the receipt arrived.
1256    pub received_at: DateTime<Utc>,
1257    /// Version of the source.
1258    pub source_version: String,
1259    /// Quality label from the source.
1260    pub quality: String,
1261    /// Digest of the value.
1262    pub digest: String,
1263}
1264
1265impl<T> ObservedValue<T> {
1266    /// Whether the value is fresh enough for `max_age` at `now`.
1267    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/// Everything a decision depended on, for replay and audit.
1280#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1281pub struct DecisionSnapshot {
1282    /// Digests of the inputs.
1283    pub input_digests: Vec<String>,
1284    /// Config version.
1285    pub instance_config_version: String,
1286    /// Product policy version that granted it.
1287    pub policy_version: String,
1288    /// Capabilities in force.
1289    pub capability_versions: Vec<CapabilityPin>,
1290    /// Execution profile revision in force.
1291    pub execution_profile_revision_id: String,
1292    /// Context values used.
1293    #[serde(default)]
1294    pub context_snapshot: Value,
1295    /// Policy values used.
1296    #[serde(default)]
1297    pub policy_snapshot: Value,
1298    /// Which Agent artifacts influenced the decision.
1299    #[serde(default)]
1300    pub agent_provenance: Value,
1301}
1302
1303/// Who supplied a parameter value.
1304#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1305#[serde(rename_all = "snake_case")]
1306pub enum ParameterProvenance {
1307    /// Explicitly supplied by the user.
1308    UserSupplied,
1309    /// Product default.
1310    ProductDefault,
1311    /// Inferred by an Agent.
1312    AgentInferred,
1313    /// Computed from other parameters.
1314    Derived,
1315}
1316
1317/// A parameter with its provenance.
1318#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1319pub struct ParameterValue {
1320    /// The value carried by this record.
1321    pub value: Value,
1322    /// Who supplied it.
1323    pub provenance: ParameterProvenance,
1324}
1325
1326/// Declared parameter and its live-mode requirements.
1327#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1328pub struct ParameterSpec {
1329    /// Display name.
1330    pub name: String,
1331    /// Must be present.
1332    pub required: bool,
1333    /// Must be user-supplied in live mode.
1334    pub required_explicit_for_live: bool,
1335}
1336
1337/// Check supplied parameters against their specs for `mode`.
1338pub 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/// Parameters that block a live run.
1367#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, thiserror::Error)]
1368#[error("missing explicit workflow requirements: {parameters:?}")]
1369pub struct MissingRequirements {
1370    /// Missing or non-explicit parameter names.
1371    pub parameters: Vec<String>,
1372}
1373
1374/// Budget reserved atomically before a root operation starts children.
1375#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1376pub struct RootOperationBudget {
1377    /// Root operation identity.
1378    pub root_operation_id: String,
1379    /// Maximum nesting depth.
1380    pub max_depth: u32,
1381    /// Maximum descendants.
1382    pub descendant_limit: u32,
1383    /// Maximum runs.
1384    pub run_limit: u32,
1385    /// Token budget.
1386    pub token_budget: u64,
1387    /// Cost budget in micro-units.
1388    pub cost_budget_micros: u64,
1389    /// Maximum actions.
1390    pub action_budget: u32,
1391    /// Latest time by which the work must finish.
1392    pub deadline: DateTime<Utc>,
1393    /// Permissions children may not exceed.
1394    #[serde(default)]
1395    pub permission_ceiling: BTreeSet<String>,
1396}
1397
1398/// Persisted Agent decision reused after recovery instead of re-asking the model.
1399#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1400pub struct DecisionArtifact {
1401    /// Stable identifier of this record.
1402    pub id: String,
1403    /// Root operation identity.
1404    pub root_operation_id: String,
1405    /// Content hash that makes the referenced artifact immutable.
1406    pub content_digest: String,
1407    /// Model identifier as registered in the model registry.
1408    pub model: String,
1409    /// Digest of the prompt.
1410    pub prompt_digest: String,
1411    /// Model output.
1412    pub output: Value,
1413    /// When the record was created (database time).
1414    pub created_at: DateTime<Utc>,
1415}
1416
1417/// Operator-facing diagnostic attached to an instance.
1418#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1419pub struct DiagnosticRecord {
1420    /// Stable diagnostic code.
1421    pub code: String,
1422    /// Human-readable message.
1423    pub message: String,
1424    /// When it was recorded.
1425    pub at: DateTime<Utc>,
1426    /// Structured detail.
1427    #[serde(default)]
1428    pub detail: Value,
1429}
1430
1431/// Rebuildable read model of an instance; never authorizes or schedules work.
1432#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1433pub struct ExecutionProjection {
1434    /// Projection schema version.
1435    pub schema_version: String,
1436    /// Authoritative event sequence the projection was built from.
1437    pub source_sequence: i64,
1438    /// Steps in progress.
1439    pub current_steps: Vec<String>,
1440    /// Steps completed.
1441    pub completed_steps: Vec<String>,
1442    /// What the instance waits for.
1443    pub waiting_on: Option<Value>,
1444    /// Most recent decision.
1445    pub last_decision: Option<Value>,
1446    /// Steps that may run next.
1447    pub next_possible_steps: Vec<String>,
1448    /// Next scheduled evaluation.
1449    pub next_trigger_at: Option<DateTime<Utc>>,
1450    /// Actions prepared but not dispatched.
1451    pub planned_actions: Vec<String>,
1452    /// Most recent diagnostic.
1453    pub latest_diagnostic: Option<DiagnosticRecord>,
1454    /// Product-defined progress.
1455    pub progress: Value,
1456    /// When the record stops being valid.
1457    pub expires_at: Option<DateTime<Utc>>,
1458}
1459
1460/// A contract value violates its invariants.
1461#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1462pub enum ContractError {
1463    /// Invalid workflow contract.
1464    #[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        // Jitter is deterministic and never exceeds 25% of the base.
1881        assert_eq!(policy.backoff_ms(1, 24), 124);
1882        assert_eq!(policy.backoff_ms(1, 25), 100);
1883        // Jitter is taken modulo the ceiling (800 / 4 = 200 here).
1884        assert_eq!(policy.backoff_ms(4, 249), 849);
1885        // The cap holds for large attempts and for base + jitter.
1886        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        // Attempt 0 behaves like the first attempt instead of underflowing.
1890        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}