Skip to main content

ironflow_store/entities/
run.rs

1//! [`Run`] entity and related request/update types.
2
3use std::collections::{HashMap, HashSet};
4use std::fmt;
5use std::time::Duration;
6
7use chrono::{DateTime, TimeDelta, Utc};
8use rust_decimal::Decimal;
9use serde::{Deserialize, Serialize};
10use serde_json::Value;
11use thiserror::Error;
12use uuid::Uuid;
13
14use super::{FsmState, ProviderKind, RunActor, RunStatus, TriggerKind, WorkerCapabilities};
15
16/// A workflow execution record.
17///
18/// Represents a single invocation of a workflow, tracking its status through
19/// the [`RunStatus`] FSM (SQL-side via [`lib_fsm`](crate::postgres::helpers::lib_fsm)),
20/// aggregated metrics, and timestamps.
21///
22/// # Examples
23///
24/// ```
25/// use ironflow_store::entities::Run;
26///
27/// // Runs are created by RunStore::create_run, not directly.
28/// ```
29#[derive(Debug, Clone, Serialize, Deserialize)]
30#[non_exhaustive]
31pub struct Run {
32    /// Unique identifier (UUIDv7, sortable by creation time).
33    pub id: Uuid,
34    /// Name of the workflow that was executed.
35    pub workflow_name: String,
36    /// Current FSM status — embeds state + state_machine_id for SQL-side transitions.
37    pub status: FsmState<RunStatus>,
38    /// How this run was triggered.
39    pub trigger: TriggerKind,
40    /// Trigger-specific payload (e.g. webhook body).
41    pub payload: Value,
42    /// Error message if the run failed.
43    pub error: Option<String>,
44    /// Number of times this run has been retried after a handler failure.
45    ///
46    /// Each retry starts a new attempt (`retry_count + 1`). A recovery after a
47    /// lost worker lease does not change it: see [`Run::lease_recoveries`].
48    pub retry_count: u32,
49    /// Maximum number of retries allowed.
50    pub max_retries: u32,
51    /// Aggregated cost across all agent steps, in USD.
52    pub cost_usd: Decimal,
53    /// Aggregated wall-clock duration across all steps, in milliseconds.
54    pub duration_ms: u64,
55    /// When the run was created (enqueued).
56    pub created_at: DateTime<Utc>,
57    /// When the run record was last updated.
58    pub updated_at: DateTime<Utc>,
59    /// When execution started (transitioned to Running).
60    pub started_at: Option<DateTime<Utc>>,
61    /// When execution finished (transitioned to a terminal state).
62    pub completed_at: Option<DateTime<Utc>>,
63    /// Version of the handler that created this run.
64    pub handler_version: Option<String>,
65    /// User-defined key-value labels for categorization and filtering.
66    #[serde(default)]
67    pub labels: HashMap<String, String>,
68    /// When the run should start executing. `None` means immediately.
69    #[serde(default)]
70    pub scheduled_at: Option<DateTime<Utc>>,
71    /// The authenticated principal that created this run.
72    ///
73    /// `None` for cron, webhook, and programmatic triggers.
74    #[serde(default)]
75    pub created_by: Option<RunActor>,
76    /// Human-readable label for [`Run::created_by`].
77    ///
78    /// Read-only projection resolved at read time from the referenced user and
79    /// API key — never written by [`crate::store::RunStore::create_run`]. `None`
80    /// when there is no actor, or when the referenced user or key no longer exists.
81    #[serde(default)]
82    pub created_by_label: Option<String>,
83    /// Client-supplied idempotency key that produced this run, if any.
84    ///
85    /// See [`IDEMPOTENCY_WINDOW`] for how long a key stays bound to its run.
86    #[serde(default)]
87    pub idempotency_key: Option<String>,
88    /// Concurrency key held by this run while it is not terminal, if any.
89    ///
90    /// See [`NewRun::concurrency_key`] for the exclusivity rule.
91    #[serde(default)]
92    pub concurrency_key: Option<String>,
93    /// Concurrency groups this run belongs to, each with its own limit.
94    ///
95    /// See [`NewRun::concurrency_limits`] for the gating rule. Empty means the
96    /// run is not limited by any group.
97    #[serde(default)]
98    pub concurrency_limits: Vec<ConcurrencyLimit>,
99    /// Maximum cumulative cost allowed for this run, in USD.
100    ///
101    /// Resolved once at run creation and frozen for the lifetime of the run.
102    /// `None` means no cap.
103    #[serde(default)]
104    pub max_cost_usd: Option<Decimal>,
105    /// Identifier of the worker currently holding the lease on this run.
106    ///
107    /// Set when a worker picks the run up, cleared as soon as the run leaves
108    /// `Running`. `None` means no worker owns this run (runs executed inline or
109    /// resumed in-process by the API server never hold a lease).
110    #[serde(default)]
111    pub worker_id: Option<String>,
112    /// When the worker lease expires.
113    ///
114    /// The worker refreshes this while it executes the run. Once it is in the
115    /// past, the reaper may requeue the run.
116    #[serde(default)]
117    pub lease_expires_at: Option<DateTime<Utc>>,
118    /// Output the workflow handler set with `WorkflowContext::set_output`.
119    ///
120    /// Written when an execution ends (completed, warning, failed or
121    /// cancelled). `None` when the handler never set an output.
122    #[serde(default)]
123    pub output: Option<Value>,
124    /// Number of times the reaper recovered this run after its worker lease
125    /// expired.
126    ///
127    /// Bounded by [`Run::max_retries`] and counted independently of
128    /// [`Run::retry_count`]: a recovered run stays in the same attempt, so the
129    /// steps it already finished are replayed instead of executed again.
130    #[serde(default)]
131    pub lease_recoveries: u32,
132    /// Provider kind (e.g. `"claude_subscription"`) the run is waiting for.
133    ///
134    /// Set only while the run is `Sleeping` because every targeted Provider
135    /// Account was rate limited; `None` in every other state. Adding,
136    /// re-enabling or renewing an account of that kind wakes the run early.
137    #[serde(default)]
138    pub capacity_wait_kind: Option<ProviderKind>,
139    /// Queue priority of the run, between [`MIN_PRIORITY`] and [`MAX_PRIORITY`].
140    ///
141    /// Workers pick the highest priority first, and runs of equal priority in
142    /// creation order. See [`NewRun::priority`].
143    #[serde(default)]
144    pub priority: i16,
145    /// Worker tags a worker must carry to pick this run up.
146    ///
147    /// Sorted and deduplicated. Empty means any worker may take the run. See
148    /// [`WorkerCapabilities::can_take`] for the routing rule.
149    #[serde(default)]
150    pub worker_tags: Vec<String>,
151    /// State the run returns to on resume.
152    ///
153    /// `Some` only while the run is `Paused`: it records the state the run was
154    /// paused from, or the state a decision taken during the pause (approval,
155    /// human input, signal) moved it to. `None` in every other state.
156    #[serde(default)]
157    pub resume_status: Option<RunStatus>,
158}
159
160/// How long a client-supplied idempotency key stays bound to its run.
161///
162/// Past this window a replayed key no longer resolves to the original run:
163/// the key is released and a fresh run is created.
164///
165/// # Examples
166///
167/// ```
168/// use ironflow_store::entities::IDEMPOTENCY_WINDOW;
169///
170/// assert_eq!(IDEMPOTENCY_WINDOW.num_hours(), 24);
171/// ```
172pub const IDEMPOTENCY_WINDOW: TimeDelta = TimeDelta::hours(24);
173
174/// Label set on every child run of a sub-workflow step, holding the id of the
175/// run that started it.
176///
177/// The root of the chain is recorded under `ironflow.io/root-run-id`. Both are
178/// set on the child run when it is created, so a suspended child can be
179/// listed by label and resumed through its root, and the descendants of a run
180/// can be found by [`RunStore::list_active_descendants`](crate::store::RunStore::list_active_descendants).
181///
182/// # Examples
183///
184/// ```
185/// use std::collections::HashMap;
186/// use ironflow_store::entities::{PARENT_RUN_ID_LABEL, RunFilter};
187/// use uuid::Uuid;
188///
189/// let parent = Uuid::now_v7();
190/// let children = RunFilter {
191///     labels: Some(HashMap::from([(PARENT_RUN_ID_LABEL.to_string(), parent.to_string())])),
192///     ..RunFilter::default()
193/// };
194/// assert!(children.labels.is_some());
195/// ```
196pub const PARENT_RUN_ID_LABEL: &str = "ironflow.io/parent-run-id";
197
198/// Maximum accepted length of an idempotency key, in bytes.
199///
200/// # Examples
201///
202/// ```
203/// use ironflow_store::entities::MAX_IDEMPOTENCY_KEY_LEN;
204///
205/// assert_eq!(MAX_IDEMPOTENCY_KEY_LEN, 255);
206/// ```
207pub const MAX_IDEMPOTENCY_KEY_LEN: usize = 255;
208
209/// Maximum accepted length of a concurrency key, in bytes.
210///
211/// # Examples
212///
213/// ```
214/// use ironflow_store::entities::MAX_CONCURRENCY_KEY_LEN;
215///
216/// assert_eq!(MAX_CONCURRENCY_KEY_LEN, 255);
217/// ```
218pub const MAX_CONCURRENCY_KEY_LEN: usize = 255;
219
220/// Maximum accepted length of a concurrency group name, in bytes.
221///
222/// # Examples
223///
224/// ```
225/// use ironflow_store::entities::MAX_CONCURRENCY_GROUP_LEN;
226///
227/// assert_eq!(MAX_CONCURRENCY_GROUP_LEN, 255);
228/// ```
229pub const MAX_CONCURRENCY_GROUP_LEN: usize = 255;
230
231/// Lowest accepted run priority.
232///
233/// # Examples
234///
235/// ```
236/// use ironflow_store::entities::MIN_PRIORITY;
237///
238/// assert_eq!(MIN_PRIORITY, -100);
239/// ```
240pub const MIN_PRIORITY: i16 = -100;
241
242/// Highest accepted run priority.
243///
244/// # Examples
245///
246/// ```
247/// use ironflow_store::entities::MAX_PRIORITY;
248///
249/// assert_eq!(MAX_PRIORITY, 100);
250/// ```
251pub const MAX_PRIORITY: i16 = 100;
252
253/// Validate a run priority before persisting it.
254///
255/// # Errors
256///
257/// Returns a message naming the accepted range when `priority` is lower than
258/// [`MIN_PRIORITY`] or higher than [`MAX_PRIORITY`].
259///
260/// # Examples
261///
262/// ```
263/// use ironflow_store::entities::validate_priority;
264///
265/// assert!(validate_priority(0).is_ok());
266/// assert!(validate_priority(100).is_ok());
267/// assert!(validate_priority(101).is_err());
268/// ```
269pub fn validate_priority(priority: i16) -> Result<(), String> {
270    if (MIN_PRIORITY..=MAX_PRIORITY).contains(&priority) {
271        Ok(())
272    } else {
273        Err(format!(
274            "priority must be between {MIN_PRIORITY} and {MAX_PRIORITY}"
275        ))
276    }
277}
278
279/// Membership of a run in a concurrency group, with the limit the run accepts.
280///
281/// A run is only moved to `Running` while fewer than `limit` root runs
282/// carrying `group` are running. Each run is compared against its own limit,
283/// so two runs of the same group may carry different limits.
284///
285/// # Examples
286///
287/// ```
288/// use ironflow_store::entities::ConcurrencyLimit;
289///
290/// let limit = ConcurrencyLimit::new("repo:acme/api", 2);
291/// assert_eq!(limit.group, "repo:acme/api");
292/// assert_eq!(limit.limit, 2);
293/// ```
294#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
295#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
296pub struct ConcurrencyLimit {
297    /// Name of the concurrency group (1 to 255 bytes).
298    pub group: String,
299    /// Maximum number of root runs of this group running at once (at least 1).
300    pub limit: u32,
301}
302
303impl ConcurrencyLimit {
304    /// Build a concurrency limit. Validation happens in
305    /// [`validate_concurrency_limits`].
306    ///
307    /// # Examples
308    ///
309    /// ```
310    /// use ironflow_store::entities::ConcurrencyLimit;
311    ///
312    /// let limit = ConcurrencyLimit::new("tenant:42", 1);
313    /// assert_eq!(limit.limit, 1);
314    /// ```
315    pub fn new(group: impl Into<String>, limit: u32) -> Self {
316        Self {
317            group: group.into(),
318            limit,
319        }
320    }
321}
322
323/// Why a list of concurrency limits was refused.
324///
325/// # Examples
326///
327/// ```
328/// use ironflow_store::entities::ConcurrencyLimitError;
329///
330/// let err = ConcurrencyLimitError::EmptyGroup;
331/// assert_eq!(err.to_string(), "concurrency group must not be empty");
332/// ```
333#[derive(Debug, Clone, PartialEq, Eq, Error)]
334#[non_exhaustive]
335pub enum ConcurrencyLimitError {
336    /// A group name was empty or whitespace only.
337    #[error("concurrency group must not be empty")]
338    EmptyGroup,
339    /// A group name exceeds [`MAX_CONCURRENCY_GROUP_LEN`] bytes.
340    #[error("concurrency group '{group}' exceeds {max} bytes")]
341    GroupTooLong {
342        /// The offending group name.
343        group: String,
344        /// The maximum accepted length, in bytes.
345        max: usize,
346    },
347    /// A limit was zero, which would hold the run back forever.
348    #[error("concurrency limit for group '{group}' must be at least 1")]
349    ZeroLimit {
350        /// The group carrying the zero limit.
351        group: String,
352    },
353    /// The same group appears twice in one run.
354    #[error("concurrency group '{group}' is listed more than once")]
355    DuplicateGroup {
356        /// The duplicated group name.
357        group: String,
358    },
359}
360
361/// Validate the concurrency limits of a run before persisting it.
362///
363/// # Errors
364///
365/// Returns [`ConcurrencyLimitError::EmptyGroup`] for an empty or
366/// whitespace-only group, [`ConcurrencyLimitError::GroupTooLong`] for a group
367/// longer than [`MAX_CONCURRENCY_GROUP_LEN`] bytes,
368/// [`ConcurrencyLimitError::ZeroLimit`] for a limit of zero and
369/// [`ConcurrencyLimitError::DuplicateGroup`] when a group is listed twice.
370///
371/// # Examples
372///
373/// ```
374/// use ironflow_store::entities::{ConcurrencyLimit, validate_concurrency_limits};
375///
376/// assert!(validate_concurrency_limits(&[]).is_ok());
377/// assert!(validate_concurrency_limits(&[ConcurrencyLimit::new("repo:acme", 2)]).is_ok());
378/// assert!(validate_concurrency_limits(&[ConcurrencyLimit::new("repo:acme", 0)]).is_err());
379/// ```
380pub fn validate_concurrency_limits(
381    limits: &[ConcurrencyLimit],
382) -> Result<(), ConcurrencyLimitError> {
383    let mut seen: HashSet<&str> = HashSet::with_capacity(limits.len());
384    for limit in limits {
385        if limit.group.trim().is_empty() {
386            return Err(ConcurrencyLimitError::EmptyGroup);
387        }
388        if limit.group.len() > MAX_CONCURRENCY_GROUP_LEN {
389            return Err(ConcurrencyLimitError::GroupTooLong {
390                group: limit.group.clone(),
391                max: MAX_CONCURRENCY_GROUP_LEN,
392            });
393        }
394        if limit.limit == 0 {
395            return Err(ConcurrencyLimitError::ZeroLimit {
396                group: limit.group.clone(),
397            });
398        }
399        if !seen.insert(limit.group.as_str()) {
400            return Err(ConcurrencyLimitError::DuplicateGroup {
401                group: limit.group.clone(),
402            });
403        }
404    }
405    Ok(())
406}
407
408/// Number of due runs held back because a concurrency group is saturated.
409///
410/// Produced by
411/// [`RunStore::count_blocked_runs_by_group`](crate::store::RunStore::count_blocked_runs_by_group).
412///
413/// # Examples
414///
415/// ```
416/// use ironflow_store::entities::ConcurrencyGroupBacklog;
417///
418/// let backlog = ConcurrencyGroupBacklog {
419///     group: "repo:acme/api".to_string(),
420///     blocked_runs: 3,
421/// };
422/// assert_eq!(backlog.blocked_runs, 3);
423/// ```
424#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
425#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
426pub struct ConcurrencyGroupBacklog {
427    /// Name of the saturated concurrency group.
428    pub group: String,
429    /// Number of due pending or retrying runs held back by this group.
430    pub blocked_runs: u64,
431}
432
433/// Outcome of [`RunStore::create_run`](crate::store::RunStore::create_run).
434///
435/// A request carrying an idempotency key already bound to a live run does not
436/// insert anything: the store returns the original run as [`RunCreation::Existing`].
437///
438/// # Examples
439///
440/// ```
441/// use std::collections::HashMap;
442/// use ironflow_store::entities::{NewRun, RunCreation, TriggerKind};
443/// use ironflow_store::memory::InMemoryStore;
444/// use ironflow_store::store::RunStore;
445/// use serde_json::json;
446///
447/// # async fn example() -> Result<(), ironflow_store::error::StoreError> {
448/// let store = InMemoryStore::new();
449/// let req = NewRun {
450///     workflow_name: "deploy".to_string(),
451///     trigger: TriggerKind::Manual,
452///     payload: json!({}),
453///     max_retries: 3,
454///     handler_version: None,
455///     labels: HashMap::new(),
456///     scheduled_at: None,
457///     created_by: None,
458///     idempotency_key: Some("deploy-2026-07-26".to_string()),
459///     concurrency_key: None,
460///     priority: 0,
461///     concurrency_limits: Vec::new(),
462///     max_cost_usd: None,
463///     worker_tags: Vec::new(),
464/// };
465///
466/// assert!(store.create_run(req.clone()).await?.is_created());
467/// assert!(!store.create_run(req).await?.is_created());
468/// # Ok(())
469/// # }
470/// ```
471#[derive(Debug, Clone, Serialize, Deserialize)]
472pub enum RunCreation {
473    /// A new run was inserted.
474    Created(Run),
475    /// The idempotency key already resolved to this run; nothing was inserted.
476    Existing(Run),
477}
478
479impl RunCreation {
480    /// Return the run, discarding whether it was created or replayed.
481    ///
482    /// # Examples
483    ///
484    /// ```no_run
485    /// # use ironflow_store::entities::RunCreation;
486    /// # fn example(creation: RunCreation) {
487    /// let run = creation.into_run();
488    /// # }
489    /// ```
490    pub fn into_run(self) -> Run {
491        match self {
492            RunCreation::Created(run) | RunCreation::Existing(run) => run,
493        }
494    }
495
496    /// Borrow the run, discarding whether it was created or replayed.
497    ///
498    /// # Examples
499    ///
500    /// ```no_run
501    /// # use ironflow_store::entities::RunCreation;
502    /// # fn example(creation: &RunCreation) {
503    /// let id = creation.run().id;
504    /// # }
505    /// ```
506    pub fn run(&self) -> &Run {
507        match self {
508            RunCreation::Created(run) | RunCreation::Existing(run) => run,
509        }
510    }
511
512    /// Whether a new run was actually inserted.
513    ///
514    /// # Examples
515    ///
516    /// ```no_run
517    /// # use ironflow_store::entities::RunCreation;
518    /// # fn example(creation: &RunCreation) {
519    /// if creation.is_created() {
520    ///     // publish a RunCreated event
521    /// }
522    /// # }
523    /// ```
524    pub fn is_created(&self) -> bool {
525        matches!(self, RunCreation::Created(_))
526    }
527}
528
529/// Request to acquire or renew a worker lease on a run.
530///
531/// # Examples
532///
533/// ```
534/// use std::time::Duration;
535/// use ironflow_store::entities::LeaseRequest;
536///
537/// let lease = LeaseRequest {
538///     worker_id: "worker-1".to_string(),
539///     ttl: Duration::from_secs(90),
540/// };
541/// assert_eq!(lease.ttl.as_secs(), 90);
542/// ```
543#[derive(Debug, Clone, PartialEq, Eq)]
544pub struct LeaseRequest {
545    /// Identifier of the worker acquiring the lease.
546    pub worker_id: String,
547    /// How long the lease stays valid without a refresh.
548    pub ttl: Duration,
549}
550
551impl LeaseRequest {
552    /// Compute the lease expiry from a reference instant.
553    ///
554    /// A TTL too large to be represented saturates to
555    /// [`DateTime::<Utc>::MAX_UTC`] instead of panicking.
556    ///
557    /// # Examples
558    ///
559    /// ```
560    /// use std::time::Duration;
561    /// use chrono::{TimeZone, Utc};
562    /// use ironflow_store::entities::LeaseRequest;
563    ///
564    /// let lease = LeaseRequest {
565    ///     worker_id: "worker-1".to_string(),
566    ///     ttl: Duration::from_secs(90),
567    /// };
568    /// let now = Utc.with_ymd_and_hms(2026, 1, 1, 0, 0, 0).unwrap();
569    /// assert_eq!(lease.expires_at(now).timestamp(), now.timestamp() + 90);
570    /// ```
571    pub fn expires_at(&self, from: DateTime<Utc>) -> DateTime<Utc> {
572        TimeDelta::from_std(self.ttl)
573            .ok()
574            .and_then(|ttl| from.checked_add_signed(ttl))
575            .unwrap_or(DateTime::<Utc>::MAX_UTC)
576    }
577}
578
579/// Change to the worker lease of a run, applied by [`RunUpdate::lease`].
580///
581/// Lets the engine hand a lease over from one run to another, atomically
582/// with a status transition: a root run resumed through its child takes the
583/// child's lease so the worker keeps renewing it and the reaper still sees it.
584///
585/// # Examples
586///
587/// ```
588/// use chrono::Utc;
589/// use ironflow_store::entities::{LeaseUpdate, RunStatus, RunUpdate};
590///
591/// let update = RunUpdate {
592///     status: Some(RunStatus::Running),
593///     lease: Some(LeaseUpdate::Set {
594///         worker_id: "worker-1".to_string(),
595///         expires_at: Utc::now(),
596///     }),
597///     ..RunUpdate::default()
598/// };
599/// assert!(matches!(update.lease, Some(LeaseUpdate::Set { .. })));
600///
601/// let release = LeaseUpdate::Release;
602/// assert_eq!(release, LeaseUpdate::Release);
603/// ```
604#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
605pub enum LeaseUpdate {
606    /// Hand the lease to `worker_id` until `expires_at`.
607    Set {
608        /// Identifier of the worker holding the lease.
609        worker_id: String,
610        /// When the lease expires without a renewal.
611        expires_at: DateTime<Utc>,
612    },
613    /// Drop the lease while the run stays in its state.
614    Release,
615}
616
617/// A run recovered by the reaper after its worker lease expired.
618///
619/// # Examples
620///
621/// ```
622/// use ironflow_store::entities::{ReapedRun, RunStatus};
623///
624/// // Reaped runs are produced by RunStore::reap_expired_leases.
625/// fn was_requeued(reaped: &ReapedRun) -> bool {
626///     reaped.to == RunStatus::Pending
627/// }
628/// ```
629#[derive(Debug, Clone)]
630#[non_exhaustive]
631pub struct ReapedRun {
632    /// The run after recovery.
633    pub run: Run,
634    /// Status the run held before recovery (always [`RunStatus::Running`]).
635    pub from: RunStatus,
636    /// Status the run was moved to: [`RunStatus::Pending`] when retries remain,
637    /// [`RunStatus::Failed`] once `max_retries` is exhausted.
638    pub to: RunStatus,
639}
640
641/// Request to create a new run.
642///
643/// # Examples
644///
645/// ```
646/// use std::collections::HashMap;
647/// use ironflow_store::entities::{NewRun, TriggerKind};
648/// use serde_json::json;
649///
650/// let req = NewRun {
651///     workflow_name: "deploy".to_string(),
652///     trigger: TriggerKind::Manual,
653///     payload: json!({}),
654///     max_retries: 3,
655///     handler_version: None,
656///     labels: HashMap::new(),
657///     scheduled_at: None,
658///     created_by: None,
659///     idempotency_key: None,
660///     concurrency_key: None,
661///     priority: 0,
662///     concurrency_limits: Vec::new(),
663///     max_cost_usd: None,
664///     worker_tags: Vec::new(),
665/// };
666/// ```
667#[derive(Debug, Clone, Serialize, Deserialize)]
668pub struct NewRun {
669    /// Workflow name.
670    pub workflow_name: String,
671    /// How the run was triggered.
672    pub trigger: TriggerKind,
673    /// Trigger-specific payload.
674    pub payload: Value,
675    /// Maximum retry attempts.
676    pub max_retries: u32,
677    /// Version of the handler at the time of run creation.
678    pub handler_version: Option<String>,
679    /// User-defined key-value labels for categorization and filtering.
680    #[serde(default)]
681    pub labels: HashMap<String, String>,
682    /// When the run should start executing. `None` means immediately.
683    #[serde(default)]
684    pub scheduled_at: Option<DateTime<Utc>>,
685    /// The authenticated principal creating this run.
686    ///
687    /// Defaults to `None` when absent from the payload, so an older worker that
688    /// does not send the field keeps working against a newer API.
689    #[serde(default)]
690    pub created_by: Option<RunActor>,
691    /// Optional idempotency key binding this request to a single run.
692    ///
693    /// When set and already bound to a run created within [`IDEMPOTENCY_WINDOW`],
694    /// the store returns that run instead of inserting a new one.
695    #[serde(default)]
696    pub idempotency_key: Option<String>,
697    /// Optional concurrency key making this run exclusive.
698    ///
699    /// At most one non-terminal run may hold a given key; creation fails with
700    /// [`StoreError::ConcurrencyConflict`](crate::error::StoreError::ConcurrencyConflict)
701    /// otherwise. The key is released when the run reaches a terminal state.
702    #[serde(default)]
703    pub concurrency_key: Option<String>,
704    /// Concurrency groups this run belongs to, each with its own limit.
705    ///
706    /// The run is only moved to `Running` while, for every listed group, fewer
707    /// than its `limit` root runs carrying that group are running. Sub-workflow
708    /// runs never count. Empty means no group limit.
709    #[serde(default)]
710    pub concurrency_limits: Vec<ConcurrencyLimit>,
711    /// Maximum cumulative cost allowed for this run, in USD. `None` means no cap.
712    #[serde(default)]
713    pub max_cost_usd: Option<Decimal>,
714    /// Queue priority, between [`MIN_PRIORITY`] and [`MAX_PRIORITY`]. `0` is the
715    /// default.
716    ///
717    /// [`RunStore::pick_next_pending`](crate::store::RunStore::pick_next_pending)
718    /// serves due runs by priority, highest first, then by creation order. A
719    /// running run is never preempted, and nothing ages a waiting run: a
720    /// continuous stream of higher priority runs keeps lower priority runs
721    /// waiting. Defaults to `0` when absent from the payload, so an older worker
722    /// keeps working against a newer API.
723    #[serde(default)]
724    pub priority: i16,
725    /// Worker tags a worker must carry to pick this run up.
726    ///
727    /// Validated with [`validate_worker_tags`](super::validate_worker_tags) and
728    /// stored normalized (sorted, deduplicated). Empty means any worker.
729    #[serde(default)]
730    pub worker_tags: Vec<String>,
731}
732
733/// Filters for listing runs.
734///
735/// All fields are optional; `None` means "no filter" for that field.
736///
737/// # Examples
738///
739/// ```
740/// use ironflow_store::entities::{RunFilter, RunStatus};
741///
742/// let filter = RunFilter {
743///     workflow_name: Some("deploy".to_string()),
744///     status: Some(RunStatus::Completed),
745///     ..RunFilter::default()
746/// };
747/// ```
748#[derive(Debug, Clone, Default)]
749pub struct RunFilter {
750    /// Filter by workflow name (exact match).
751    pub workflow_name: Option<String>,
752    /// Filter by run status.
753    pub status: Option<RunStatus>,
754    /// Only include runs created after this timestamp.
755    pub created_after: Option<DateTime<Utc>>,
756    /// Only include runs created before this timestamp.
757    pub created_before: Option<DateTime<Utc>>,
758    /// When `Some(true)`, only include runs that have at least one step.
759    /// When `Some(false)`, only include runs with no steps.
760    /// When `None`, no filtering on steps.
761    pub has_steps: Option<bool>,
762    /// Filter by label key-value pair. Only include runs that have ALL specified labels.
763    pub labels: Option<HashMap<String, String>>,
764    /// Filter by author. Matches runs created by this user directly, and runs
765    /// created by one of this user's API keys.
766    pub created_by_user_id: Option<Uuid>,
767    /// Filter by concurrency group. Only include runs whose concurrency limits
768    /// contain this group.
769    pub concurrency_group: Option<String>,
770    /// Filter by queue priority (exact match).
771    pub priority: Option<i16>,
772    /// Only include runs the worker with these capabilities could take (see
773    /// [`WorkerCapabilities::can_take`]). `None` means no filter.
774    pub eligible_for: Option<WorkerCapabilities>,
775}
776
777/// Partial update for a run.
778///
779/// # Examples
780///
781/// ```
782/// use ironflow_store::entities::{RunUpdate, RunStatus};
783///
784/// let update = RunUpdate {
785///     status: Some(RunStatus::Completed),
786///     ..RunUpdate::default()
787/// };
788/// ```
789#[derive(Debug, Clone, Default, Serialize, Deserialize)]
790pub struct RunUpdate {
791    /// New status.
792    pub status: Option<RunStatus>,
793    /// Error message.
794    pub error: Option<String>,
795    /// Increment retry count.
796    pub increment_retry: bool,
797    /// Aggregated cost.
798    pub cost_usd: Option<Decimal>,
799    /// Aggregated duration.
800    pub duration_ms: Option<u64>,
801    /// When execution started.
802    pub started_at: Option<DateTime<Utc>>,
803    /// When execution completed.
804    pub completed_at: Option<DateTime<Utc>>,
805    /// When the run should next be picked up. Used to arm the retry backoff.
806    #[serde(default)]
807    pub scheduled_at: Option<DateTime<Utc>>,
808    /// Output set by the workflow handler. `None` leaves the stored output unchanged.
809    #[serde(default)]
810    pub output: Option<Value>,
811    /// Lease change applied after the status transition. `None` leaves the
812    /// lease as the status transition left it (a run leaving `Running` always
813    /// loses its lease). See [`LeaseUpdate`].
814    ///
815    /// # Examples
816    ///
817    /// ```
818    /// use ironflow_store::entities::{LeaseUpdate, RunUpdate};
819    ///
820    /// let update = RunUpdate {
821    ///     lease: Some(LeaseUpdate::Release),
822    ///     ..RunUpdate::default()
823    /// };
824    /// assert_eq!(update.lease, Some(LeaseUpdate::Release));
825    /// ```
826    #[serde(default)]
827    pub lease: Option<LeaseUpdate>,
828    /// Provider kind the run waits for, see [`Run::capacity_wait_kind`].
829    ///
830    /// Applied only with `status: Some(Sleeping)`; any other status clears it.
831    #[serde(default)]
832    pub capacity_wait_kind: Option<ProviderKind>,
833    /// New target state of a `Paused` run, see [`Run::resume_status`].
834    ///
835    /// Applied only with `status: None`, and refused with
836    /// [`StoreError::InvalidTransition`](crate::error::StoreError::InvalidTransition)
837    /// unless the run is `Paused`: it lets a decision taken during the pause
838    /// change where the run goes on resume without leaving `Paused`.
839    #[serde(default)]
840    pub resume_status: Option<RunStatus>,
841}
842
843/// Retention policy for purging old runs.
844///
845/// Runs are eligible for purging when they are in a terminal state
846/// ([`RunStatus::is_terminal`]) **and** exceed either the age limit or the
847/// per-workflow count limit.
848///
849/// # Examples
850///
851/// ```
852/// use ironflow_store::entities::PurgePolicy;
853///
854/// let policy = PurgePolicy {
855///     max_age_days: 90,
856///     max_runs_per_workflow: 1000,
857///     dry_run: false,
858/// };
859/// assert_eq!(policy.max_age_days, 90);
860/// ```
861#[derive(Debug, Clone)]
862pub struct PurgePolicy {
863    /// Runs older than this many days are eligible for purging.
864    pub max_age_days: u32,
865    /// When a workflow has more runs than this, the oldest terminal runs are
866    /// eligible for purging.
867    pub max_runs_per_workflow: u32,
868    /// When `true`, the purger logs what would be deleted but does not delete.
869    pub dry_run: bool,
870}
871
872/// Why a run was selected for purging.
873///
874/// # Examples
875///
876/// ```
877/// use ironflow_store::entities::PurgeReason;
878///
879/// let reason = PurgeReason::TooOld;
880/// assert_eq!(format!("{reason}"), "too_old");
881/// ```
882#[derive(Debug, Clone, Copy, PartialEq, Eq)]
883pub enum PurgeReason {
884    /// The run exceeded [`PurgePolicy::max_age_days`].
885    TooOld,
886    /// The workflow exceeded [`PurgePolicy::max_runs_per_workflow`].
887    ExceedsWorkflowLimit,
888}
889
890impl fmt::Display for PurgeReason {
891    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
892        match self {
893            PurgeReason::TooOld => f.write_str("too_old"),
894            PurgeReason::ExceedsWorkflowLimit => f.write_str("exceeds_workflow_limit"),
895        }
896    }
897}
898
899/// A run selected for purging by [`RunStore::list_purgeable_runs`](crate::store::RunStore::list_purgeable_runs).
900///
901/// # Examples
902///
903/// ```
904/// use ironflow_store::entities::{PurgeReason, PurgeableRun};
905/// use uuid::Uuid;
906///
907/// let purgeable = PurgeableRun {
908///     run_id: Uuid::now_v7(),
909///     workflow_name: "deploy".to_string(),
910///     reason: PurgeReason::TooOld,
911/// };
912/// assert_eq!(purgeable.reason, PurgeReason::TooOld);
913/// ```
914#[derive(Debug, Clone)]
915pub struct PurgeableRun {
916    /// The run to purge.
917    pub run_id: Uuid,
918    /// Workflow the run belongs to.
919    pub workflow_name: String,
920    /// Why this run was selected.
921    pub reason: PurgeReason,
922}
923
924#[cfg(test)]
925mod tests {
926    use std::collections::HashMap;
927
928    use super::*;
929    use serde_json::json;
930
931    #[test]
932    fn newrun_serde_roundtrip() {
933        let new_run = NewRun {
934            created_by: None,
935            workflow_name: "deploy".to_string(),
936            trigger: TriggerKind::Manual,
937            payload: json!({"key": "value"}),
938            max_retries: 3,
939            handler_version: Some("1.2.0".to_string()),
940            labels: HashMap::from([("env".to_string(), "prod".to_string())]),
941            scheduled_at: None,
942            idempotency_key: None,
943            concurrency_key: None,
944            priority: 0,
945            concurrency_limits: Vec::new(),
946            max_cost_usd: Some(Decimal::new(250, 2)),
947            worker_tags: Vec::new(),
948        };
949
950        let json = serde_json::to_string(&new_run).expect("serialize");
951        let back: NewRun = serde_json::from_str(&json).expect("deserialize");
952        assert_eq!(back.max_cost_usd, new_run.max_cost_usd);
953        assert_eq!(back.workflow_name, new_run.workflow_name);
954        assert_eq!(back.trigger, new_run.trigger);
955        assert_eq!(back.payload, new_run.payload);
956        assert_eq!(back.max_retries, new_run.max_retries);
957        assert_eq!(back.handler_version, new_run.handler_version);
958        assert_eq!(back.labels, new_run.labels);
959        assert_eq!(back.scheduled_at, new_run.scheduled_at);
960        assert_eq!(back.created_by, new_run.created_by);
961        assert_eq!(back.idempotency_key, new_run.idempotency_key);
962    }
963
964    #[test]
965    fn newrun_serde_roundtrip_with_actor() {
966        let actor = RunActor::ApiKey {
967            api_key_id: Uuid::now_v7(),
968            user_id: Uuid::now_v7(),
969        };
970        let new_run = NewRun {
971            workflow_name: "deploy".to_string(),
972            trigger: TriggerKind::Api,
973            payload: json!({}),
974            max_retries: 0,
975            handler_version: None,
976            labels: HashMap::new(),
977            scheduled_at: None,
978            created_by: Some(actor.clone()),
979            idempotency_key: None,
980            concurrency_key: None,
981            priority: 0,
982            concurrency_limits: Vec::new(),
983            max_cost_usd: None,
984            worker_tags: Vec::new(),
985        };
986
987        let json = serde_json::to_string(&new_run).expect("serialize");
988        let back: NewRun = serde_json::from_str(&json).expect("deserialize");
989        assert_eq!(back.created_by, Some(actor));
990    }
991
992    #[test]
993    fn newrun_deserializes_without_created_by() {
994        // An older worker POSTs a payload with no `created_by` field.
995        let raw = json!({
996            "workflow_name": "deploy",
997            "trigger": {"kind": "workflow"},
998            "payload": {},
999            "max_retries": 0,
1000            "handler_version": null,
1001        });
1002
1003        let new_run: NewRun = serde_json::from_value(raw).expect("deserialize");
1004        assert!(new_run.created_by.is_none());
1005    }
1006
1007    #[test]
1008    fn run_serde_preserves_all_fields() {
1009        use crate::entities::FsmState;
1010        use chrono::Utc;
1011        use uuid::Uuid;
1012
1013        let now = Utc::now();
1014        let run = Run {
1015            id: Uuid::now_v7(),
1016            workflow_name: "test-wf".to_string(),
1017            status: FsmState::new(RunStatus::Running, Uuid::now_v7()),
1018            trigger: TriggerKind::Webhook {
1019                path: "/hooks/test".to_string(),
1020            },
1021            payload: json!({"data": 123}),
1022            error: Some("test error".to_string()),
1023            retry_count: 2,
1024            max_retries: 5,
1025            cost_usd: Decimal::new(1234, 2),
1026            duration_ms: 5000,
1027            created_at: now,
1028            updated_at: now,
1029            started_at: Some(now),
1030            completed_at: Some(now),
1031            handler_version: Some("2.0.0".to_string()),
1032            labels: HashMap::from([
1033                ("env".to_string(), "staging".to_string()),
1034                ("team".to_string(), "platform".to_string()),
1035            ]),
1036            scheduled_at: Some(now),
1037            created_by: Some(RunActor::User {
1038                user_id: Uuid::now_v7(),
1039            }),
1040            created_by_label: Some("alice".to_string()),
1041            idempotency_key: Some("gh:abc-123".to_string()),
1042            concurrency_key: Some("issue:12".to_string()),
1043            priority: 7,
1044            concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 2)],
1045            max_cost_usd: Some(Decimal::new(500, 2)),
1046            worker_id: Some("worker-1".to_string()),
1047            lease_expires_at: Some(now),
1048            output: Some(json!({"verdict": "approved", "score": 9})),
1049            lease_recoveries: 1,
1050            capacity_wait_kind: Some(ProviderKind::new("claude_subscription")),
1051            resume_status: None,
1052            worker_tags: vec!["gpu".to_string(), "region:eu".to_string()],
1053        };
1054
1055        let json = serde_json::to_string(&run).expect("serialize");
1056        let back: Run = serde_json::from_str(&json).expect("deserialize");
1057
1058        assert_eq!(back.id, run.id);
1059        assert_eq!(back.workflow_name, run.workflow_name);
1060        assert_eq!(back.status.state, run.status.state);
1061        assert_eq!(back.trigger, run.trigger);
1062        assert_eq!(back.payload, run.payload);
1063        assert_eq!(back.error, run.error);
1064        assert_eq!(back.retry_count, run.retry_count);
1065        assert_eq!(back.max_retries, run.max_retries);
1066        assert_eq!(back.cost_usd, run.cost_usd);
1067        assert_eq!(back.duration_ms, run.duration_ms);
1068        assert_eq!(back.started_at, run.started_at);
1069        assert_eq!(back.completed_at, run.completed_at);
1070        assert_eq!(back.handler_version, run.handler_version);
1071        assert_eq!(back.labels, run.labels);
1072        assert_eq!(back.scheduled_at, run.scheduled_at);
1073        assert_eq!(back.created_by, run.created_by);
1074        assert_eq!(back.created_by_label, run.created_by_label);
1075        assert_eq!(back.idempotency_key, run.idempotency_key);
1076        assert_eq!(back.concurrency_key, run.concurrency_key);
1077        assert_eq!(back.concurrency_limits, run.concurrency_limits);
1078        assert_eq!(back.max_cost_usd, run.max_cost_usd);
1079        assert_eq!(back.worker_id, run.worker_id);
1080        assert_eq!(back.lease_expires_at, run.lease_expires_at);
1081        assert_eq!(back.output, run.output);
1082        assert_eq!(back.lease_recoveries, run.lease_recoveries);
1083        assert_eq!(back.capacity_wait_kind, run.capacity_wait_kind);
1084        assert_eq!(back.priority, 7);
1085        assert_eq!(back.worker_tags, run.worker_tags);
1086    }
1087
1088    #[test]
1089    fn run_without_output_field_deserializes_to_none_output() {
1090        // A run serialized before the `output` column existed has no such key.
1091        let now = Utc::now();
1092        let run = Run {
1093            id: Uuid::now_v7(),
1094            workflow_name: "legacy".to_string(),
1095            status: FsmState::new(RunStatus::Completed, Uuid::now_v7()),
1096            trigger: TriggerKind::Manual,
1097            payload: json!({}),
1098            error: None,
1099            retry_count: 0,
1100            max_retries: 0,
1101            cost_usd: Decimal::ZERO,
1102            duration_ms: 0,
1103            created_at: now,
1104            updated_at: now,
1105            started_at: None,
1106            completed_at: None,
1107            handler_version: None,
1108            labels: HashMap::new(),
1109            scheduled_at: None,
1110            created_by: None,
1111            created_by_label: None,
1112            idempotency_key: None,
1113            concurrency_key: None,
1114            priority: 0,
1115            concurrency_limits: Vec::new(),
1116            max_cost_usd: None,
1117            worker_id: None,
1118            lease_expires_at: None,
1119            output: Some(json!("set")),
1120            lease_recoveries: 2,
1121            capacity_wait_kind: None,
1122            resume_status: None,
1123            worker_tags: vec!["gpu".to_string()],
1124        };
1125        let mut raw = serde_json::to_value(&run).expect("serialize");
1126        raw.as_object_mut().expect("object").remove("output");
1127        raw.as_object_mut()
1128            .expect("object")
1129            .remove("lease_recoveries");
1130        raw.as_object_mut().expect("object").remove("worker_tags");
1131
1132        let back: Run = serde_json::from_value(raw).expect("deserialize");
1133        assert!(back.output.is_none());
1134        assert_eq!(back.lease_recoveries, 0);
1135        assert!(back.worker_tags.is_empty());
1136    }
1137
1138    #[test]
1139    fn runupdate_output_round_trips_and_defaults_to_none() {
1140        let update = RunUpdate {
1141            output: Some(json!({"verdict": "rejected"})),
1142            ..RunUpdate::default()
1143        };
1144        let json = serde_json::to_string(&update).expect("serialize");
1145        let back: RunUpdate = serde_json::from_str(&json).expect("deserialize");
1146        assert_eq!(back.output, update.output);
1147
1148        let mut legacy = serde_json::to_value(RunUpdate::default()).expect("serialize");
1149        legacy.as_object_mut().expect("object").remove("output");
1150        let parsed: RunUpdate = serde_json::from_value(legacy).expect("deserialize");
1151        assert!(parsed.output.is_none());
1152    }
1153
1154    #[test]
1155    fn newrun_max_cost_usd_defaults_to_none_when_absent() {
1156        let without_cap = NewRun {
1157            workflow_name: "deploy".to_string(),
1158            trigger: TriggerKind::Manual,
1159            payload: json!({}),
1160            max_retries: 0,
1161            handler_version: None,
1162            labels: HashMap::new(),
1163            scheduled_at: None,
1164            created_by: None,
1165            idempotency_key: None,
1166            concurrency_key: None,
1167            priority: 0,
1168            concurrency_limits: Vec::new(),
1169            max_cost_usd: None,
1170            worker_tags: Vec::new(),
1171        };
1172        let mut value = serde_json::to_value(&without_cap).expect("serialize");
1173        value
1174            .as_object_mut()
1175            .expect("object")
1176            .remove("max_cost_usd");
1177
1178        let parsed: NewRun = serde_json::from_value(value).expect("deserialize");
1179        assert!(parsed.max_cost_usd.is_none());
1180    }
1181
1182    #[test]
1183    fn newrun_priority_defaults_to_zero_when_absent() {
1184        let raw = json!({
1185            "workflow_name": "deploy",
1186            "trigger": {"kind": "manual"},
1187            "payload": {},
1188            "max_retries": 0,
1189            "handler_version": null,
1190        });
1191
1192        let new_run: NewRun = serde_json::from_value(raw).expect("deserialize");
1193        assert_eq!(new_run.priority, 0);
1194    }
1195
1196    #[test]
1197    fn validate_priority_accepts_bounds() {
1198        assert!(validate_priority(MIN_PRIORITY).is_ok());
1199        assert!(validate_priority(-100).is_ok());
1200        assert!(validate_priority(0).is_ok());
1201        assert!(validate_priority(100).is_ok());
1202        assert!(validate_priority(MAX_PRIORITY).is_ok());
1203    }
1204
1205    #[test]
1206    fn validate_priority_rejects_out_of_range() {
1207        assert_eq!(
1208            validate_priority(101),
1209            Err("priority must be between -100 and 100".to_string())
1210        );
1211        assert!(validate_priority(-101).is_err());
1212        assert!(validate_priority(i16::MAX).is_err());
1213        assert!(validate_priority(i16::MIN).is_err());
1214    }
1215
1216    #[test]
1217    fn newrun_concurrency_limits_default_to_empty_when_absent() {
1218        let raw = json!({
1219            "workflow_name": "deploy",
1220            "trigger": {"kind": "manual"},
1221            "payload": {},
1222            "max_retries": 0,
1223            "handler_version": null,
1224        });
1225
1226        let new_run: NewRun = serde_json::from_value(raw).expect("deserialize");
1227        assert!(new_run.concurrency_limits.is_empty());
1228    }
1229
1230    #[test]
1231    fn newrun_serde_roundtrip_keeps_concurrency_limits() {
1232        let new_run = NewRun {
1233            workflow_name: "deploy".to_string(),
1234            trigger: TriggerKind::Manual,
1235            payload: json!({}),
1236            max_retries: 0,
1237            handler_version: None,
1238            labels: HashMap::new(),
1239            scheduled_at: None,
1240            created_by: None,
1241            idempotency_key: None,
1242            concurrency_key: None,
1243            priority: 0,
1244            concurrency_limits: vec![
1245                ConcurrencyLimit::new("repo:acme", 2),
1246                ConcurrencyLimit::new("tenant:42", 5),
1247            ],
1248            max_cost_usd: None,
1249            worker_tags: Vec::new(),
1250        };
1251
1252        let json = serde_json::to_string(&new_run).expect("serialize");
1253        let back: NewRun = serde_json::from_str(&json).expect("deserialize");
1254        assert_eq!(back.concurrency_limits, new_run.concurrency_limits);
1255    }
1256
1257    #[test]
1258    fn validate_concurrency_limits_accepts_empty_and_valid() {
1259        assert_eq!(validate_concurrency_limits(&[]), Ok(()));
1260        assert_eq!(
1261            validate_concurrency_limits(&[
1262                ConcurrencyLimit::new("repo:acme", 1),
1263                ConcurrencyLimit::new("tenant:42", 10),
1264            ]),
1265            Ok(())
1266        );
1267    }
1268
1269    #[test]
1270    fn validate_concurrency_limits_rejects_empty_group() {
1271        assert_eq!(
1272            validate_concurrency_limits(&[ConcurrencyLimit::new("", 1)]),
1273            Err(ConcurrencyLimitError::EmptyGroup)
1274        );
1275        assert_eq!(
1276            validate_concurrency_limits(&[ConcurrencyLimit::new("   ", 1)]),
1277            Err(ConcurrencyLimitError::EmptyGroup)
1278        );
1279    }
1280
1281    #[test]
1282    fn validate_concurrency_limits_rejects_zero_limit() {
1283        assert_eq!(
1284            validate_concurrency_limits(&[ConcurrencyLimit::new("repo:acme", 0)]),
1285            Err(ConcurrencyLimitError::ZeroLimit {
1286                group: "repo:acme".to_string(),
1287            })
1288        );
1289    }
1290
1291    #[test]
1292    fn validate_concurrency_limits_rejects_duplicate_group() {
1293        assert_eq!(
1294            validate_concurrency_limits(&[
1295                ConcurrencyLimit::new("repo:acme", 1),
1296                ConcurrencyLimit::new("repo:acme", 3),
1297            ]),
1298            Err(ConcurrencyLimitError::DuplicateGroup {
1299                group: "repo:acme".to_string(),
1300            })
1301        );
1302    }
1303
1304    #[test]
1305    fn validate_concurrency_limits_rejects_too_long_group() {
1306        let at_max = "g".repeat(MAX_CONCURRENCY_GROUP_LEN);
1307        assert_eq!(
1308            validate_concurrency_limits(&[ConcurrencyLimit::new(at_max, 1)]),
1309            Ok(())
1310        );
1311
1312        let too_long = "g".repeat(MAX_CONCURRENCY_GROUP_LEN + 1);
1313        assert_eq!(
1314            validate_concurrency_limits(&[ConcurrencyLimit::new(too_long.clone(), 1)]),
1315            Err(ConcurrencyLimitError::GroupTooLong {
1316                group: too_long,
1317                max: MAX_CONCURRENCY_GROUP_LEN,
1318            })
1319        );
1320    }
1321
1322    #[test]
1323    fn validate_concurrency_limits_accepts_unicode_group() {
1324        let group = "d\u{e9}p\u{f4}t:caf\u{e9}-\u{2615}";
1325        assert_eq!(
1326            validate_concurrency_limits(&[ConcurrencyLimit::new(group, 2)]),
1327            Ok(())
1328        );
1329        // Length is counted in bytes: 128 two-byte characters exceed 255 bytes.
1330        let too_long = "\u{e9}".repeat(128);
1331        assert!(matches!(
1332            validate_concurrency_limits(&[ConcurrencyLimit::new(too_long, 1)]),
1333            Err(ConcurrencyLimitError::GroupTooLong { .. })
1334        ));
1335    }
1336
1337    #[test]
1338    fn runupdate_serde_roundtrip() {
1339        let update = RunUpdate {
1340            status: Some(RunStatus::Completed),
1341            error: Some("test error".to_string()),
1342            increment_retry: true,
1343            cost_usd: Some(Decimal::new(5000, 2)),
1344            duration_ms: Some(3000),
1345            started_at: None,
1346            completed_at: None,
1347            scheduled_at: Some(Utc::now()),
1348            output: None,
1349            lease: None,
1350            capacity_wait_kind: None,
1351            resume_status: None,
1352        };
1353
1354        let json = serde_json::to_string(&update).expect("serialize");
1355        let back: RunUpdate = serde_json::from_str(&json).expect("deserialize");
1356
1357        assert_eq!(back.status, update.status);
1358        assert_eq!(back.error, update.error);
1359        assert_eq!(back.increment_retry, update.increment_retry);
1360        assert_eq!(back.cost_usd, update.cost_usd);
1361        assert_eq!(back.duration_ms, update.duration_ms);
1362        assert_eq!(back.scheduled_at, update.scheduled_at);
1363    }
1364
1365    #[test]
1366    fn runupdate_lease_serde_roundtrip() {
1367        let set = RunUpdate {
1368            status: Some(RunStatus::Running),
1369            lease: Some(LeaseUpdate::Set {
1370                worker_id: "worker-1".to_string(),
1371                expires_at: Utc::now(),
1372            }),
1373            ..RunUpdate::default()
1374        };
1375        let json = serde_json::to_string(&set).expect("serialize");
1376        let back: RunUpdate = serde_json::from_str(&json).expect("deserialize");
1377        assert_eq!(back.lease, set.lease);
1378
1379        let release = RunUpdate {
1380            lease: Some(LeaseUpdate::Release),
1381            ..RunUpdate::default()
1382        };
1383        let json = serde_json::to_string(&release).expect("serialize");
1384        let back: RunUpdate = serde_json::from_str(&json).expect("deserialize");
1385        assert_eq!(back.lease, Some(LeaseUpdate::Release));
1386    }
1387
1388    #[test]
1389    fn runupdate_without_lease_field_deserializes_to_none() {
1390        let json = json!({
1391            "status": "completed",
1392            "error": null,
1393            "increment_retry": false,
1394            "cost_usd": null,
1395            "duration_ms": null,
1396            "started_at": null,
1397            "completed_at": null,
1398        });
1399        let back: RunUpdate = serde_json::from_value(json).expect("deserialize");
1400        assert_eq!(back.status, Some(RunStatus::Completed));
1401        assert!(back.lease.is_none());
1402    }
1403
1404    #[test]
1405    fn runfilter_default_is_no_filters() {
1406        let filter = RunFilter::default();
1407        assert!(filter.workflow_name.is_none());
1408        assert!(filter.status.is_none());
1409        assert!(filter.created_after.is_none());
1410        assert!(filter.created_before.is_none());
1411        assert!(filter.created_by_user_id.is_none());
1412    }
1413
1414    #[test]
1415    fn runfilter_with_multiple_criteria() {
1416        let filter = RunFilter {
1417            workflow_name: Some("deploy".to_string()),
1418            status: Some(RunStatus::Running),
1419            ..RunFilter::default()
1420        };
1421
1422        assert_eq!(filter.workflow_name, Some("deploy".to_string()));
1423        assert_eq!(filter.status, Some(RunStatus::Running));
1424        assert!(filter.created_after.is_none());
1425        assert!(filter.created_before.is_none());
1426    }
1427}