1use chrono::{DateTime, Utc};
2use std::sync::Arc;
3use uuid::Uuid;
4
5use crate::error::{FlowError, Result};
6use crate::model::{
7 project_run, ActiveHookSnapshot, FlowEvent, FlowEventEnvelope, HookStatus, RuntimeCommand,
8 ScheduledWakeup, ScheduledWakeupKind, StepStatus, WaitStatus, WorkflowRunSnapshot,
9 WorkflowRunStatus, WorkflowRunSummary, WorkflowRunSuspension, WorkflowSpec,
10};
11use crate::observe::{FlowEventObserver, NoopFlowEventObserver};
12use crate::runtime::{FlowRuntime, WorkflowInvocation};
13use crate::store::{scheduled_wakeups_for_snapshot, FlowEventStore, InMemoryEventStore};
14
15mod operations;
16mod steps;
17mod validation;
18use steps::StepExecutionContext;
19use validation::{
20 ensure_child_operation_matches, ensure_hook_command_matches, ensure_progress_matches,
21 ensure_same_start, ensure_step_batch_valid, ensure_step_command_matches,
22 ensure_wait_command_matches, is_event_conflict, validate_run_id,
23};
24
25pub struct FlowEngineBuilder {
27 store: Arc<dyn FlowEventStore>,
28 runtime: Arc<dyn FlowRuntime>,
29 observer: Arc<dyn FlowEventObserver>,
30 max_replay_iterations: usize,
31}
32
33impl FlowEngineBuilder {
34 pub fn new(runtime: Arc<dyn FlowRuntime>) -> Self {
35 Self {
36 store: Arc::new(InMemoryEventStore::new()),
37 runtime,
38 observer: Arc::new(NoopFlowEventObserver),
39 max_replay_iterations: 1024,
40 }
41 }
42
43 pub fn with_store(mut self, store: Arc<dyn FlowEventStore>) -> Self {
44 self.store = store;
45 self
46 }
47
48 pub fn with_observer(mut self, observer: Arc<dyn FlowEventObserver>) -> Self {
49 self.observer = observer;
50 self
51 }
52
53 pub fn with_max_replay_iterations(mut self, max_replay_iterations: usize) -> Self {
54 self.max_replay_iterations = max_replay_iterations.max(1);
55 self
56 }
57
58 pub fn build(self) -> FlowEngine {
59 FlowEngine {
60 store: self.store,
61 runtime: self.runtime,
62 observer: self.observer,
63 max_replay_iterations: self.max_replay_iterations,
64 }
65 }
66}
67
68#[derive(Clone)]
70pub struct FlowEngine {
71 store: Arc<dyn FlowEventStore>,
72 runtime: Arc<dyn FlowRuntime>,
73 observer: Arc<dyn FlowEventObserver>,
74 max_replay_iterations: usize,
75}
76
77impl FlowEngine {
78 pub fn builder(runtime: Arc<dyn FlowRuntime>) -> FlowEngineBuilder {
79 FlowEngineBuilder::new(runtime)
80 }
81
82 pub fn new(store: Arc<dyn FlowEventStore>, runtime: Arc<dyn FlowRuntime>) -> Self {
83 Self {
84 store,
85 runtime,
86 observer: Arc::new(NoopFlowEventObserver),
87 max_replay_iterations: 1024,
88 }
89 }
90
91 pub fn in_memory(runtime: Arc<dyn FlowRuntime>) -> Self {
92 Self::new(Arc::new(InMemoryEventStore::new()), runtime)
93 }
94
95 pub fn store(&self) -> Arc<dyn FlowEventStore> {
96 Arc::clone(&self.store)
97 }
98
99 pub fn observer(&self) -> Arc<dyn FlowEventObserver> {
100 Arc::clone(&self.observer)
101 }
102
103 pub async fn start(&self, spec: WorkflowSpec, input: serde_json::Value) -> Result<String> {
105 let run_id = Uuid::new_v4().to_string();
106 self.start_with_id(run_id, spec, input).await
107 }
108
109 pub async fn start_with_id(
114 &self,
115 run_id: impl Into<String>,
116 spec: WorkflowSpec,
117 input: serde_json::Value,
118 ) -> Result<String> {
119 spec.validate()?;
120 let run_id = run_id.into();
121 validate_run_id(&run_id)?;
122
123 for _ in 0..self.max_replay_iterations {
124 match self.store.list(&run_id).await {
125 Ok(history) => {
126 let snapshot = project_run(&run_id, &history)?;
127 ensure_same_start(&run_id, &snapshot, &spec, &input)?;
128 if !history
129 .iter()
130 .any(|event| matches!(event.event, FlowEvent::RunStarted))
131 {
132 match self
133 .record_event_at(&run_id, snapshot.last_sequence, FlowEvent::RunStarted)
134 .await
135 {
136 Ok(_) => {}
137 Err(err) if is_event_conflict(&err) => continue,
138 Err(err) => return Err(err),
139 }
140 }
141 match self.drive(&run_id).await {
142 Ok(_) => return Ok(run_id),
143 Err(err) if is_event_conflict(&err) => continue,
144 Err(err) => return Err(err),
145 }
146 }
147 Err(FlowError::RunNotFound(_)) => {
148 let created = match self
149 .record_event_at(
150 &run_id,
151 0,
152 FlowEvent::RunCreated {
153 spec: spec.clone(),
154 input: input.clone(),
155 },
156 )
157 .await
158 {
159 Ok(created) => created,
160 Err(err) if is_event_conflict(&err) => continue,
161 Err(err) => return Err(err),
162 };
163 match self
164 .record_event_at(&run_id, created.sequence, FlowEvent::RunStarted)
165 .await
166 {
167 Ok(_) => {}
168 Err(err) if is_event_conflict(&err) => continue,
169 Err(err) => return Err(err),
170 }
171 match self.drive(&run_id).await {
172 Ok(_) => return Ok(run_id),
173 Err(err) if is_event_conflict(&err) => continue,
174 Err(err) => return Err(err),
175 }
176 }
177 Err(err) => return Err(err),
178 }
179 }
180
181 Err(FlowError::ReplayLimitExceeded(self.max_replay_iterations))
182 }
183
184 pub async fn resume_wait(&self, run_id: &str, wait_id: &str) -> Result<()> {
186 for _ in 0..self.max_replay_iterations {
187 let snapshot = self.snapshot(run_id).await?;
188 if snapshot.status.is_terminal() {
189 return Err(FlowError::RunTerminal(run_id.to_string()));
190 }
191 match snapshot.waits.get(wait_id) {
192 Some(wait) if wait.status == WaitStatus::Waiting => {
193 match self
194 .record_event_at(
195 run_id,
196 snapshot.last_sequence,
197 FlowEvent::WaitCompleted {
198 wait_id: wait_id.to_string(),
199 },
200 )
201 .await
202 {
203 Ok(_) => {}
204 Err(err) if is_event_conflict(&err) => continue,
205 Err(err) => return Err(err),
206 }
207 match self.drive(run_id).await {
208 Ok(_) => return Ok(()),
209 Err(err) if is_event_conflict(&err) => continue,
210 Err(err) => return Err(err),
211 }
212 }
213 Some(_) => match self.drive(run_id).await {
214 Ok(_) => return Ok(()),
215 Err(err) if is_event_conflict(&err) => continue,
216 Err(err) => return Err(err),
217 },
218 None => {
219 return Err(FlowError::InvalidTransition(format!(
220 "wait {wait_id} does not exist for run {run_id}"
221 )))
222 }
223 }
224 }
225
226 Err(FlowError::ReplayLimitExceeded(self.max_replay_iterations))
227 }
228
229 pub async fn resume_hook(
231 &self,
232 run_id: &str,
233 hook_id: &str,
234 payload: serde_json::Value,
235 ) -> Result<()> {
236 for _ in 0..self.max_replay_iterations {
237 let snapshot = self.snapshot(run_id).await?;
238 if snapshot.status.is_terminal() {
239 return Err(FlowError::RunTerminal(run_id.to_string()));
240 }
241 match snapshot.hooks.get(hook_id) {
242 Some(hook) if hook.status == HookStatus::Active => {
243 match self
244 .record_event_at(
245 run_id,
246 snapshot.last_sequence,
247 FlowEvent::HookReceived {
248 hook_id: hook_id.to_string(),
249 payload: payload.clone(),
250 },
251 )
252 .await
253 {
254 Ok(_) => {}
255 Err(err) if is_event_conflict(&err) => continue,
256 Err(err) => return Err(err),
257 }
258 match self.drive(run_id).await {
259 Ok(_) => return Ok(()),
260 Err(err) if is_event_conflict(&err) => continue,
261 Err(err) => return Err(err),
262 }
263 }
264 Some(_) => match self.drive(run_id).await {
265 Ok(_) => return Ok(()),
266 Err(err) if is_event_conflict(&err) => continue,
267 Err(err) => return Err(err),
268 },
269 None => {
270 return Err(FlowError::InvalidTransition(format!(
271 "hook {hook_id} does not exist for run {run_id}"
272 )))
273 }
274 }
275 }
276
277 Err(FlowError::ReplayLimitExceeded(self.max_replay_iterations))
278 }
279
280 pub async fn dispose_hook(&self, run_id: &str, hook_id: &str) -> Result<()> {
288 for _ in 0..self.max_replay_iterations {
289 let snapshot = self.snapshot(run_id).await?;
290 if snapshot.status.is_terminal() {
291 return Err(FlowError::RunTerminal(run_id.to_string()));
292 }
293 match snapshot.hooks.get(hook_id) {
294 Some(hook) if hook.status == HookStatus::Active => {
295 match self
296 .record_event_at(
297 run_id,
298 snapshot.last_sequence,
299 FlowEvent::HookDisposed {
300 hook_id: hook_id.to_string(),
301 },
302 )
303 .await
304 {
305 Ok(_) => {}
306 Err(err) if is_event_conflict(&err) => continue,
307 Err(err) => return Err(err),
308 }
309 match self.drive(run_id).await {
310 Ok(_) => return Ok(()),
311 Err(err) if is_event_conflict(&err) => continue,
312 Err(err) => return Err(err),
313 }
314 }
315 Some(_) => match self.drive(run_id).await {
316 Ok(_) => return Ok(()),
317 Err(err) if is_event_conflict(&err) => continue,
318 Err(err) => return Err(err),
319 },
320 None => {
321 return Err(FlowError::InvalidTransition(format!(
322 "hook {hook_id} does not exist for run {run_id}"
323 )))
324 }
325 }
326 }
327
328 Err(FlowError::ReplayLimitExceeded(self.max_replay_iterations))
329 }
330
331 pub async fn resume_hook_by_token(
336 &self,
337 token: &str,
338 payload: serde_json::Value,
339 ) -> Result<(String, String)> {
340 let mut matches = self
341 .store
342 .find_active_hooks_by_token(token)
343 .await?
344 .into_iter()
345 .map(|active| (active.run_id, active.hook.hook_id))
346 .collect::<Vec<_>>();
347
348 match matches.len() {
349 0 => Err(FlowError::HookTokenNotFound(token.to_string())),
350 1 => {
351 let (run_id, hook_id) = matches.remove(0);
352 self.resume_hook(&run_id, &hook_id, payload).await?;
353 Ok((run_id, hook_id))
354 }
355 _ => Err(FlowError::InvalidTransition(
356 "hook token is active in multiple runs (value redacted)".to_string(),
357 )),
358 }
359 }
360
361 pub async fn dispose_hook_by_token(&self, token: &str) -> Result<(String, String)> {
366 let mut matches = self
367 .store
368 .find_active_hooks_by_token(token)
369 .await?
370 .into_iter()
371 .map(|active| (active.run_id, active.hook.hook_id))
372 .collect::<Vec<_>>();
373
374 match matches.len() {
375 0 => Err(FlowError::HookTokenNotFound(token.to_string())),
376 1 => {
377 let (run_id, hook_id) = matches.remove(0);
378 self.dispose_hook(&run_id, &hook_id).await?;
379 Ok((run_id, hook_id))
380 }
381 _ => Err(FlowError::InvalidTransition(
382 "hook token is active in multiple runs (value redacted)".to_string(),
383 )),
384 }
385 }
386
387 pub async fn list_due_waits(&self, now: DateTime<Utc>) -> Result<Vec<(String, String)>> {
392 let mut due = self
393 .list_due_wakeups(now)
394 .await?
395 .into_iter()
396 .filter(|wakeup| wakeup.kind == ScheduledWakeupKind::Wait)
397 .map(|wakeup| (wakeup.run_id, wakeup.subject_id))
398 .collect::<Vec<_>>();
399 due.sort();
400 Ok(due)
401 }
402
403 pub async fn resume_due_waits(&self, now: DateTime<Utc>) -> Result<Vec<(String, String)>> {
408 let due = self.list_due_waits(now).await?;
409 let mut resumed = Vec::with_capacity(due.len());
410 for (run_id, wait_id) in due {
411 self.resume_wait(&run_id, &wait_id).await?;
412 resumed.push((run_id, wait_id));
413 }
414 Ok(resumed)
415 }
416
417 pub async fn list_due_retries(&self, now: DateTime<Utc>) -> Result<Vec<(String, String)>> {
419 let mut due = self
420 .list_due_wakeups(now)
421 .await?
422 .into_iter()
423 .filter(|wakeup| wakeup.kind == ScheduledWakeupKind::Retry)
424 .map(|wakeup| (wakeup.run_id, wakeup.subject_id))
425 .collect::<Vec<_>>();
426 due.sort();
427 Ok(due)
428 }
429
430 pub async fn list_due_wakeups(&self, now: DateTime<Utc>) -> Result<Vec<ScheduledWakeup>> {
432 let mut wakeups = self.store.list_due_wakeups(now).await?;
433 wakeups.sort_by(|left, right| {
434 (left.kind, left.run_id.as_str(), left.subject_id.as_str()).cmp(&(
435 right.kind,
436 right.run_id.as_str(),
437 right.subject_id.as_str(),
438 ))
439 });
440 Ok(wakeups)
441 }
442
443 pub async fn resume_due_retries(&self, now: DateTime<Utc>) -> Result<Vec<(String, String)>> {
445 let due = self.list_due_retries(now).await?;
446 let mut run_ids = Vec::new();
447 for (run_id, _) in &due {
448 if !run_ids.contains(run_id) {
449 run_ids.push(run_id.clone());
450 }
451 }
452 for run_id in run_ids {
453 self.drive_at(&run_id, now).await?;
454 }
455 Ok(due)
456 }
457
458 pub async fn resume_scheduled_run(
465 &self,
466 run_id: &str,
467 now: DateTime<Utc>,
468 ) -> Result<Vec<ScheduledWakeup>> {
469 let history = self.store.list(run_id).await?;
470 let snapshot = project_run(run_id, &history)?;
471 let due = scheduled_wakeups_for_snapshot(&snapshot)
472 .into_iter()
473 .filter(|wakeup| wakeup.scheduled_at <= now)
474 .collect::<Vec<_>>();
475
476 let due_wait_ids = due
477 .iter()
478 .filter(|wakeup| wakeup.kind == ScheduledWakeupKind::Wait)
479 .map(|wakeup| wakeup.subject_id.clone())
480 .collect::<Vec<_>>();
481 let has_due_retries = due
482 .iter()
483 .any(|wakeup| wakeup.kind == ScheduledWakeupKind::Retry);
484
485 for wait_id in due_wait_ids {
486 self.resume_wait(run_id, &wait_id).await?;
487 }
488 if has_due_retries {
489 self.drive_at(run_id, now).await?;
490 }
491
492 Ok(due)
493 }
494
495 pub async fn snapshot(&self, run_id: &str) -> Result<WorkflowRunSnapshot> {
496 let history = self.store.list(run_id).await?;
497 project_run(run_id, &history)
498 }
499
500 pub async fn history(&self, run_id: &str) -> Result<Vec<FlowEventEnvelope>> {
501 self.store.list(run_id).await
502 }
503
504 pub async fn list_run_ids(&self) -> Result<Vec<String>> {
505 self.store.list_run_ids().await
506 }
507
508 pub async fn list_snapshots(&self) -> Result<Vec<WorkflowRunSnapshot>> {
509 let mut snapshots = Vec::new();
510 for run_id in self.store.list_run_ids().await? {
511 snapshots.push(self.snapshot(&run_id).await?);
512 }
513 Ok(snapshots)
514 }
515
516 pub async fn run_summary(&self) -> Result<WorkflowRunSummary> {
522 let snapshots = self.list_snapshots().await?;
523 Ok(WorkflowRunSummary::from_snapshots(&snapshots))
524 }
525
526 pub async fn list_open_suspensions(
532 &self,
533 now: DateTime<Utc>,
534 ) -> Result<Vec<WorkflowRunSuspension>> {
535 let mut suspensions = Vec::new();
536 for run_id in self.store.list_run_ids().await? {
537 let snapshot = self.snapshot(&run_id).await?;
538 if snapshot.status.is_terminal() {
539 continue;
540 }
541 for wait in snapshot.waits.values() {
542 if wait.status == WaitStatus::Waiting {
543 suspensions.push(WorkflowRunSuspension::Wait {
544 run_id: run_id.clone(),
545 wait: wait.clone(),
546 due: wait.resume_at <= now,
547 });
548 }
549 }
550 for hook in snapshot.hooks.values() {
551 if hook.status == HookStatus::Active {
552 suspensions.push(WorkflowRunSuspension::Hook {
553 run_id: run_id.clone(),
554 hook: hook.clone(),
555 });
556 }
557 }
558 for step in snapshot.steps.values() {
559 if step.status == StepStatus::Pending {
560 if let Some(retry_after) = step.retry_after {
561 suspensions.push(WorkflowRunSuspension::Retry {
562 run_id: run_id.clone(),
563 step: step.clone(),
564 due: retry_after <= now,
565 });
566 }
567 }
568 }
569 }
570 suspensions.sort_by(|left, right| {
571 (left.run_id(), left.kind_order(), left.subject_id()).cmp(&(
572 right.run_id(),
573 right.kind_order(),
574 right.subject_id(),
575 ))
576 });
577 Ok(suspensions)
578 }
579
580 pub async fn next_wakeup(&self, now: DateTime<Utc>) -> Result<Option<WorkflowRunSuspension>> {
586 for _ in 0..2 {
587 let Some(wakeup) = self.store.next_scheduled_wakeup().await? else {
588 return Ok(None);
589 };
590 match self.snapshot(&wakeup.run_id).await {
591 Ok(snapshot) => {
592 if let Some(suspension) = resolve_scheduled_wakeup(&snapshot, &wakeup, now) {
593 return Ok(Some(suspension));
594 }
595 }
596 Err(FlowError::RunNotFound(_)) => {}
597 Err(error) => return Err(error),
598 }
599 }
600
601 self.next_wakeup_by_replay(now).await
602 }
603
604 async fn next_wakeup_by_replay(
605 &self,
606 now: DateTime<Utc>,
607 ) -> Result<Option<WorkflowRunSuspension>> {
608 let mut wakeups = self.list_open_suspensions(now).await?;
609 wakeups.retain(|suspension| suspension.scheduled_at().is_some());
610 wakeups.sort_by(|left, right| {
611 (
612 left.scheduled_at(),
613 left.run_id(),
614 left.kind_order(),
615 left.subject_id(),
616 )
617 .cmp(&(
618 right.scheduled_at(),
619 right.run_id(),
620 right.kind_order(),
621 right.subject_id(),
622 ))
623 });
624 Ok(wakeups.into_iter().next())
625 }
626
627 pub async fn list_active_hooks(&self) -> Result<Vec<ActiveHookSnapshot>> {
633 self.store.list_active_hooks().await
634 }
635
636 pub async fn drive(&self, run_id: &str) -> Result<WorkflowRunSnapshot> {
639 self.drive_at(run_id, Utc::now()).await
640 }
641
642 async fn drive_at(&self, run_id: &str, now: DateTime<Utc>) -> Result<WorkflowRunSnapshot> {
643 'replay: for _ in 0..self.max_replay_iterations {
644 let history = self.store.list(run_id).await?;
645 let snapshot = project_run(run_id, &history)?;
646 if snapshot.status.is_terminal()
647 || snapshot
648 .waits
649 .values()
650 .any(|wait| wait.status == WaitStatus::Waiting)
651 || snapshot
652 .hooks
653 .values()
654 .any(|hook| hook.status == HookStatus::Active)
655 || (snapshot.has_future_retry(now) && snapshot.due_retries(now).is_empty())
656 {
657 return Ok(snapshot);
658 }
659
660 let command = self
661 .runtime
662 .run_workflow(WorkflowInvocation {
663 run_id: run_id.to_string(),
664 spec: snapshot.spec.clone(),
665 input: snapshot.input.clone(),
666 history,
667 })
668 .await?;
669
670 match command {
671 RuntimeCommand::Complete { output } => {
672 if snapshot.status == WorkflowRunStatus::Cancelling {
673 return Err(FlowError::InvalidTransition(format!(
674 "workflow run {run_id} completed after cancellation was requested; cleanup-aware cancellation must return cancel or fail"
675 )));
676 }
677 match self
678 .record_event_at(
679 run_id,
680 snapshot.last_sequence,
681 FlowEvent::RunCompleted { output },
682 )
683 .await
684 {
685 Ok(_) => {}
686 Err(err) if is_event_conflict(&err) => continue,
687 Err(err) => return Err(err),
688 }
689 return self.snapshot(run_id).await;
690 }
691 RuntimeCommand::Fail { error } => {
692 match self
693 .record_event_at(
694 run_id,
695 snapshot.last_sequence,
696 FlowEvent::RunFailed { error },
697 )
698 .await
699 {
700 Ok(_) => {}
701 Err(err) if is_event_conflict(&err) => continue,
702 Err(err) => return Err(err),
703 }
704 return self.snapshot(run_id).await;
705 }
706 RuntimeCommand::Cancel => {
707 let cancellation = snapshot.cancellation.as_ref().ok_or_else(|| {
708 FlowError::InvalidTransition(format!(
709 "workflow run {run_id} returned cancel without a durable cancellation request"
710 ))
711 })?;
712 match self
713 .record_event_at(
714 run_id,
715 snapshot.last_sequence,
716 FlowEvent::RunCancelled {
717 reason: cancellation.request.reason.clone(),
718 },
719 )
720 .await
721 {
722 Ok(_) => {}
723 Err(err) if is_event_conflict(&err) => continue,
724 Err(err) => return Err(err),
725 }
726 return self.snapshot(run_id).await;
727 }
728 RuntimeCommand::Timeout { deadline, reason } => {
729 match self
730 .record_event_at(
731 run_id,
732 snapshot.last_sequence,
733 FlowEvent::RunTimedOut { deadline, reason },
734 )
735 .await
736 {
737 Ok(_) => {}
738 Err(err) if is_event_conflict(&err) => continue,
739 Err(err) => return Err(err),
740 }
741 return self.snapshot(run_id).await;
742 }
743 RuntimeCommand::RecordProgress { progress } => {
744 progress.validate()?;
745 if let Some(existing) = snapshot.progress(&progress.progress_id) {
746 ensure_progress_matches(run_id, existing, &progress)?;
747 return Err(FlowError::InvalidTransition(format!(
748 "workflow rescheduled progress {} without progress",
749 progress.progress_id
750 )));
751 }
752 match self
753 .record_event_at(
754 run_id,
755 snapshot.last_sequence,
756 FlowEvent::RunProgressRecorded { progress },
757 )
758 .await
759 {
760 Ok(_) => {}
761 Err(err) if is_event_conflict(&err) => continue,
762 Err(err) => return Err(err),
763 }
764 }
765 RuntimeCommand::LinkChildOperation { child } => {
766 child.validate()?;
767 if let Some(existing) = snapshot.child_operation(&child.reference_id) {
768 ensure_child_operation_matches(run_id, existing, &child)?;
769 return Err(FlowError::InvalidTransition(format!(
770 "workflow rescheduled child operation {} without progress",
771 child.reference_id
772 )));
773 }
774 match self
775 .record_event_at(
776 run_id,
777 snapshot.last_sequence,
778 FlowEvent::ChildOperationLinked { child },
779 )
780 .await
781 {
782 Ok(_) => {}
783 Err(err) if is_event_conflict(&err) => continue,
784 Err(err) => return Err(err),
785 }
786 }
787 RuntimeCommand::ScheduleStep {
788 step_id,
789 step_name,
790 input,
791 retry,
792 } => {
793 if let Some(step) = snapshot.steps.get(&step_id) {
794 ensure_step_command_matches(run_id, step, &step_name, &input, retry)?;
795 if matches!(
796 step.status,
797 StepStatus::Completed | StepStatus::Failed | StepStatus::Cancelled
798 ) {
799 return Err(FlowError::InvalidTransition(format!(
800 "workflow rescheduled terminal step {step_id} without progress"
801 )));
802 }
803 }
804 match self
805 .execute_step(
806 run_id,
807 &snapshot,
808 StepExecutionContext {
809 step_id,
810 step_name,
811 input,
812 retry,
813 now,
814 },
815 )
816 .await
817 {
818 Ok(()) => {}
819 Err(err) if is_event_conflict(&err) => continue,
820 Err(err) => return Err(err),
821 }
822 }
823 RuntimeCommand::ScheduleSteps { steps } => {
824 ensure_step_batch_valid(&steps)?;
825 for step in &steps {
826 if let Some(existing) = snapshot.steps.get(&step.step_id) {
827 ensure_step_command_matches(
828 run_id,
829 existing,
830 &step.step_name,
831 &step.input,
832 step.retry,
833 )?;
834 }
835 }
836 if steps.iter().all(|step| {
837 snapshot.steps.get(&step.step_id).is_some_and(|existing| {
838 matches!(
839 existing.status,
840 StepStatus::Completed | StepStatus::Failed | StepStatus::Cancelled
841 )
842 })
843 }) {
844 let step_ids = steps
845 .iter()
846 .map(|step| step.step_id.as_str())
847 .collect::<Vec<_>>()
848 .join(", ");
849 return Err(FlowError::InvalidTransition(format!(
850 "workflow rescheduled only terminal steps without progress: {step_ids}"
851 )));
852 }
853 match self.execute_step_batch(run_id, &snapshot, steps, now).await {
854 Ok(()) => {}
855 Err(err) if is_event_conflict(&err) => continue 'replay,
856 Err(err) => return Err(err),
857 }
858 }
859 RuntimeCommand::WaitUntil { wait_id, resume_at } => {
860 match snapshot.waits.get(&wait_id) {
861 Some(wait) => {
862 ensure_wait_command_matches(run_id, wait, resume_at)?;
863 match wait.status {
864 WaitStatus::Completed => continue,
865 WaitStatus::Waiting => return self.snapshot(run_id).await,
866 WaitStatus::Cancelled => {
867 return Err(FlowError::InvalidTransition(format!(
868 "workflow rescheduled cancelled wait {wait_id}; cancellation cleanup must use a distinct stable identity"
869 )))
870 }
871 }
872 }
873 None => {
874 match self
875 .record_event_at(
876 run_id,
877 snapshot.last_sequence,
878 FlowEvent::WaitCreated { wait_id, resume_at },
879 )
880 .await
881 {
882 Ok(_) => {}
883 Err(err) if is_event_conflict(&err) => continue,
884 Err(err) => return Err(err),
885 }
886 return self.snapshot(run_id).await;
887 }
888 }
889 }
890 RuntimeCommand::CreateHook {
891 hook_id,
892 token,
893 metadata,
894 } => match snapshot.hooks.get(&hook_id) {
895 Some(hook) => {
896 ensure_hook_command_matches(run_id, hook, &token, &metadata)?;
897 match hook.status {
898 HookStatus::Received | HookStatus::Disposed => continue,
899 HookStatus::Active => return self.snapshot(run_id).await,
900 HookStatus::Cancelled => {
901 return Err(FlowError::InvalidTransition(format!(
902 "workflow rescheduled cancelled hook {hook_id}; cancellation cleanup must use a distinct stable identity"
903 )))
904 }
905 }
906 }
907 None => {
908 self.ensure_hook_token_available(run_id, &hook_id, &token)
909 .await?;
910 match self
911 .record_event_at(
912 run_id,
913 snapshot.last_sequence,
914 FlowEvent::HookCreated {
915 hook_id,
916 token,
917 metadata,
918 },
919 )
920 .await
921 {
922 Ok(_) => {}
923 Err(err) if is_event_conflict(&err) => continue,
924 Err(err) => return Err(err),
925 }
926 return self.snapshot(run_id).await;
927 }
928 },
929 }
930 }
931
932 Err(FlowError::ReplayLimitExceeded(self.max_replay_iterations))
933 }
934
935 async fn terminate_run(&self, run_id: &str, event: FlowEvent) -> Result<()> {
936 for _ in 0..self.max_replay_iterations {
937 let snapshot = self.snapshot(run_id).await?;
938 if snapshot.status.is_terminal() {
939 return Ok(());
940 }
941 match self
942 .record_event_at(run_id, snapshot.last_sequence, event.clone())
943 .await
944 {
945 Ok(_) => return Ok(()),
946 Err(err) if is_event_conflict(&err) => continue,
947 Err(err) => return Err(err),
948 }
949 }
950 Err(FlowError::ReplayLimitExceeded(self.max_replay_iterations))
951 }
952
953 async fn record_event_at(
954 &self,
955 run_id: &str,
956 expected_sequence: u64,
957 event: FlowEvent,
958 ) -> Result<FlowEventEnvelope> {
959 let envelope = self
960 .store
961 .append_if_sequence(run_id, expected_sequence, event)
962 .await?;
963 self.observer.observe(envelope.clone()).await;
964 Ok(envelope)
965 }
966
967 async fn ensure_hook_token_available(
968 &self,
969 run_id: &str,
970 hook_id: &str,
971 token: &str,
972 ) -> Result<()> {
973 for active in self.store.find_active_hooks_by_token(token).await? {
974 if active.run_id == run_id && active.hook.hook_id == hook_id {
975 continue;
976 }
977 return Err(FlowError::HookTokenConflict {
978 token: token.to_string(),
979 existing_run_id: active.run_id,
980 existing_hook_id: active.hook.hook_id,
981 });
982 }
983 Ok(())
984 }
985}
986
987fn resolve_scheduled_wakeup(
988 snapshot: &WorkflowRunSnapshot,
989 wakeup: &ScheduledWakeup,
990 now: DateTime<Utc>,
991) -> Option<WorkflowRunSuspension> {
992 if snapshot.run_id != wakeup.run_id || snapshot.status.is_terminal() {
993 return None;
994 }
995 match wakeup.kind {
996 ScheduledWakeupKind::Wait => {
997 let wait = snapshot.waits.get(&wakeup.subject_id)?;
998 if wait.status != WaitStatus::Waiting || wait.resume_at != wakeup.scheduled_at {
999 return None;
1000 }
1001 Some(WorkflowRunSuspension::Wait {
1002 run_id: wakeup.run_id.clone(),
1003 wait: wait.clone(),
1004 due: wakeup.scheduled_at <= now,
1005 })
1006 }
1007 ScheduledWakeupKind::Retry => {
1008 let step = snapshot.steps.get(&wakeup.subject_id)?;
1009 if step.status != StepStatus::Pending || step.retry_after != Some(wakeup.scheduled_at) {
1010 return None;
1011 }
1012 Some(WorkflowRunSuspension::Retry {
1013 run_id: wakeup.run_id.clone(),
1014 step: step.clone(),
1015 due: wakeup.scheduled_at <= now,
1016 })
1017 }
1018 }
1019}