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}