1use std::collections::HashMap;
2use std::future::Future;
3use std::sync::Arc;
4use std::time::{Duration, SystemTime, UNIX_EPOCH};
5
6use futures_util::TryStreamExt;
7use taquba::object_store::ObjectStore;
8use taquba::{
9 Clock, EnqueueOptions, EnqueueRequest, EnqueueResult, JobRecord, JobStatus, Queue,
10 SettlementEffects, WaitOutcome, WorkerHandle,
11};
12use tokio_util::sync::CancellationToken;
13use tracing::{debug, instrument, warn};
14
15use crate::durable::{
16 self, DurableCurrentStep, DurableErrorKind, DurableRunOutcome, DurableRunRecord,
17 DurableRunResult, DurableStepOutcome, DurableStepOutcomeRecord, DurableTermination,
18};
19use crate::effects::StagedEffects;
20use crate::error::{Error, Result};
21use crate::group::{GroupStore, Membership, RunGroup, pending_member, terminated_member};
22use crate::keys::{
23 DEDUP_PREFIX, GROUP_TERMINAL_KV_PREFIX, HEADER_RUN_ID, HEADER_STEP, HEADER_TERMINAL,
24 RESERVED_HEADER_PREFIX, RESERVED_KV_PREFIX, RunId, TERMINAL_KV_PREFIX, hash_input,
25 outcome_kv_key, run_kv_key, step_kv_key,
26};
27use crate::memo::{MemoStore, RUN_RESULT_MEMO_KEY};
28use crate::runner::{StepErrorKind, StepOutcome, StepRunner, Trigger};
29use crate::sweep::{Clearable, Sweep, run_periodically};
30use crate::terminal::{RunOutcome, TerminalHook, TerminalStatus};
31use crate::view::WorkflowView;
32use crate::worker::{ClaimedStep, StepWorker};
33
34fn current_step_bytes(step_number: u32, job_id: &str) -> Vec<u8> {
36 durable::encode(&DurableCurrentStep {
37 step_number,
38 job_id: job_id.to_string(),
39 })
40}
41
42fn remaining_delay(stored_at_ms: u64, now_ms: u64, delay: Duration) -> Duration {
46 let elapsed = Duration::from_millis(now_ms.saturating_sub(stored_at_ms));
47 delay.saturating_sub(elapsed)
48}
49
50#[derive(Debug, Default)]
55pub(crate) struct StepEnqueueOpts {
56 pub(crate) run_at: Option<SystemTime>,
58 pub(crate) priority: Option<u32>,
60 pub(crate) max_attempts: Option<u32>,
62 pub(crate) reserved_headers: Vec<(&'static str, String)>,
64}
65
66#[derive(Debug, Clone, Default)]
72pub struct RunOptions {
73 pub headers: HashMap<String, String>,
77 pub priority: Option<u32>,
79 pub max_attempts_per_step: Option<u32>,
82 pub run_at: Option<SystemTime>,
86}
87
88#[derive(Debug, Clone, Default)]
90pub struct RunSpec {
91 pub run_id: Option<RunId>,
103 pub input: Vec<u8>,
105 pub options: RunOptions,
107 pub effects: SettlementEffects,
116}
117
118#[derive(Debug, Clone)]
123pub struct SubmitOutcome {
124 pub run_id: RunId,
126 pub newly_submitted: bool,
132 pub job_id: String,
136}
137
138#[derive(Debug, Clone)]
141pub struct RunStatus {
142 pub run_id: RunId,
144 pub state: RunState,
146 pub current_step: u32,
149}
150
151#[derive(Debug, Clone, PartialEq, Eq)]
153pub enum RunState {
154 Pending,
156 Running,
158 Cancelling,
170 Terminated(RunTermination),
175}
176
177#[derive(Debug, Clone)]
180pub(crate) struct RunResult {
181 pub(crate) termination: RunTermination,
182 pub(crate) outcome: RunOutcome,
183}
184
185#[derive(Debug, Clone, PartialEq, Eq)]
188pub struct RunTermination {
189 pub status: TerminalStatus,
191 pub error: Option<String>,
195 pub error_kind: Option<StepErrorKind>,
201 pub final_step: u32,
203 pub terminated_at_ms: u64,
206}
207
208impl From<DurableTermination> for RunTermination {
209 fn from(record: DurableTermination) -> Self {
210 Self {
211 status: record.status.into(),
212 error: record.error,
213 error_kind: record.error_kind.map(Into::into),
214 final_step: record.final_step,
215 terminated_at_ms: record.terminated_at_ms,
216 }
217 }
218}
219
220#[derive(Debug, Clone)]
222pub struct RunEnd {
223 pub termination: RunTermination,
225 pub outcome: Option<RunOutcome>,
228}
229
230pub struct WorkflowRuntimeBuilder<R, H> {
234 queue: Arc<Queue>,
235 object_store: Arc<dyn ObjectStore>,
236 queue_name: String,
237 memo_prefix: Option<String>,
238 runner: R,
239 terminal_hook: H,
240 max_concurrent_steps: usize,
241 poll_interval: Duration,
242 memo_retention: Option<Duration>,
243 group_retention: Option<Duration>,
244 step_output_replay: bool,
245 clock: Arc<dyn Clock>,
246}
247
248impl<R: StepRunner, H: TerminalHook> WorkflowRuntimeBuilder<R, H> {
249 pub fn queue_name(mut self, name: impl Into<String>) -> Self {
253 self.queue_name = name.into();
254 self
255 }
256
257 pub fn memo_prefix(mut self, prefix: impl Into<String>) -> Self {
262 self.memo_prefix = Some(prefix.into());
263 self
264 }
265
266 pub fn max_concurrent_steps(mut self, n: usize) -> Self {
269 assert!(n > 0, "max_concurrent_steps must be at least 1");
270 self.max_concurrent_steps = n;
271 self
272 }
273
274 pub fn poll_interval(mut self, interval: Duration) -> Self {
277 self.poll_interval = interval;
278 self
279 }
280
281 pub fn memo_retention(mut self, retention: Duration) -> Self {
288 self.memo_retention = Some(retention);
289 self
290 }
291
292 pub fn step_output_replay(mut self) -> Self {
312 self.step_output_replay = true;
313 self
314 }
315
316 pub fn clock(mut self, clock: Arc<dyn Clock>) -> Self {
320 self.clock = clock;
321 self
322 }
323
324 pub fn group_retention(mut self, retention: Duration) -> Self {
332 self.group_retention = Some(retention);
333 self
334 }
335
336 pub fn build(self) -> WorkflowRuntime<R, H>
338 where
339 H: 'static,
340 {
341 let terminal_hook = Arc::new(self.terminal_hook);
342 let observes: Arc<dyn Fn(&RunOutcome) -> bool + Send + Sync> = {
343 let hook = terminal_hook.clone();
344 Arc::new(move |outcome| hook.observes(outcome))
345 };
346 let memo_prefix = self
347 .memo_prefix
348 .unwrap_or_else(|| format!("{}-memo", self.queue_name));
349 let memo_store = MemoStore::new(self.object_store.clone(), memo_prefix.clone());
350 let group_store = GroupStore::new(
351 self.object_store,
352 memo_prefix,
353 memo_store.clone(),
354 self.queue.clone(),
355 );
356 let memo_sweep = self.memo_retention.map(|retention| {
357 Arc::new(Sweep::new(
358 TERMINAL_KV_PREFIX,
359 retention,
360 RunStore {
361 memo_store: memo_store.clone(),
362 },
363 ))
364 });
365 let group_sweep = self.group_retention.map(|retention| {
366 Arc::new(Sweep::new(
367 GROUP_TERMINAL_KV_PREFIX,
368 retention,
369 group_store.clone(),
370 ))
371 });
372 let view = WorkflowView::new(self.queue.view().clone(), memo_store.clone());
373 let core = RuntimeCore {
374 queue: self.queue,
375 view,
376 queue_name: self.queue_name,
377 max_concurrent_steps: self.max_concurrent_steps,
378 poll_interval: self.poll_interval,
379 memo_store,
380 group_store,
381 memo_sweep,
382 group_sweep,
383 step_output_replay: self.step_output_replay,
384 clock: self.clock,
385 observes,
386 };
387 let inner = RuntimeInner {
388 runner: self.runner,
389 terminal_hook,
390 core: Arc::new(core),
391 };
392 WorkflowRuntime {
393 inner: Arc::new(inner),
394 }
395 }
396}
397
398struct RunStore {
402 memo_store: MemoStore,
403}
404
405impl Clearable for RunStore {
406 type Error = Error;
407
408 async fn clear(&self, run_id: &RunId) -> Result<Vec<Vec<u8>>> {
409 self.memo_store.clear_memos_for_run(run_id).await?;
410 Ok(vec![outcome_kv_key(run_id)])
411 }
412}
413
414pub struct WorkflowRuntime<R, H> {
416 pub(crate) inner: Arc<RuntimeInner<R, H>>,
417}
418
419impl<R, H> Clone for WorkflowRuntime<R, H> {
420 fn clone(&self) -> Self {
421 Self {
422 inner: self.inner.clone(),
423 }
424 }
425}
426
427pub(crate) struct RuntimeInner<R, H> {
431 pub(crate) runner: R,
432 pub(crate) terminal_hook: Arc<H>,
433 pub(crate) core: Arc<RuntimeCore>,
434}
435
436pub(crate) struct RuntimeCore {
441 pub(crate) queue: Arc<Queue>,
442 pub(crate) view: WorkflowView,
443 queue_name: String,
444 max_concurrent_steps: usize,
445 poll_interval: Duration,
446 pub(crate) memo_store: MemoStore,
447 pub(crate) group_store: GroupStore,
448 pub(crate) memo_sweep: Option<Arc<Sweep>>,
453 pub(crate) group_sweep: Option<Arc<Sweep>>,
458 pub(crate) step_output_replay: bool,
461 pub(crate) clock: Arc<dyn Clock>,
464 observes: Arc<dyn Fn(&RunOutcome) -> bool + Send + Sync>,
467}
468
469impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H> {
470 pub fn builder(
485 queue: Arc<Queue>,
486 object_store: Arc<dyn ObjectStore>,
487 runner: R,
488 terminal_hook: H,
489 ) -> WorkflowRuntimeBuilder<R, H> {
490 let clock = queue.clock();
491 WorkflowRuntimeBuilder {
492 queue,
493 object_store,
494 queue_name: "workflow-steps".to_string(),
495 memo_prefix: None,
496 runner,
497 terminal_hook,
498 max_concurrent_steps: 16,
499 poll_interval: Duration::from_millis(250),
500 memo_retention: None,
501 group_retention: None,
502 step_output_replay: false,
503 clock,
504 }
505 }
506
507 pub async fn submit(&self, spec: RunSpec) -> Result<SubmitOutcome> {
518 self.inner.core.submit(spec).await
519 }
520
521 pub fn view(&self) -> &WorkflowView {
524 &self.inner.core.view
525 }
526
527 pub async fn status(&self, run_id: &RunId) -> Result<Option<RunStatus>> {
530 self.inner.core.view.status(run_id).await
531 }
532
533 pub async fn outcome(&self, run_id: &RunId) -> Result<Option<RunOutcome>> {
535 self.inner.core.view.outcome(run_id).await
536 }
537
538 pub async fn wait(&self, run_id: &RunId) -> Result<RunEnd> {
549 self.inner.core.wait(run_id).await
550 }
551
552 pub async fn wait_timeout(&self, run_id: &RunId, timeout: Duration) -> Result<Option<RunEnd>> {
555 self.inner.core.wait_timeout(run_id, timeout).await
556 }
557
558 pub async fn cancel(&self, run_id: &RunId) -> Result<bool> {
589 self.inner.core.cancel(run_id).await
590 }
591
592 pub fn group(&self, id: RunId) -> RunGroup {
594 RunGroup::new(self.inner.core.clone(), id)
595 }
596
597 pub fn new_group(&self) -> RunGroup {
599 RunGroup::new(self.inner.core.clone(), RunId::generate())
600 }
601
602 pub fn spawn<F>(&self, shutdown: F) -> RunnerHandle
607 where
608 F: Future<Output = ()> + Send + 'static,
609 R: 'static,
610 H: 'static,
611 {
612 let runtime = self.clone();
613 WorkerHandle::spawn(shutdown, |stop| async move { runtime.run_with(stop).await })
614 }
615
616 pub async fn run<F>(&self, shutdown: F) -> Result<()>
624 where
625 F: Future<Output = ()>,
626 R: 'static,
627 H: 'static,
628 {
629 let stop = CancellationToken::new();
630 let mut worker = std::pin::pin!(self.run_with(stop.clone()));
631 tokio::select! {
632 res = &mut worker => res,
633 () = shutdown => {
634 stop.cancel();
635 worker.await
636 }
637 }
638 }
639
640 async fn run_with(&self, stop: CancellationToken) -> Result<()>
644 where
645 R: 'static,
646 H: 'static,
647 {
648 let mut background: Vec<_> = self
649 .inner
650 .core
651 .sweeps()
652 .map(|sweep| {
653 let sweep = sweep.clone();
654 let core = self.inner.core.clone();
655 let token = stop.clone();
656 tokio::spawn(async move {
657 sweep
658 .run(&core.queue, &*core.clock, core.poll_interval, token)
659 .await;
660 })
661 })
662 .collect();
663 background.push({
664 let core = self.inner.core.clone();
665 let token = stop.clone();
666 tokio::spawn(async move { core.run_dead_step_reconciliation(token).await })
667 });
668
669 let worker = Arc::new(StepWorker {
670 inner: self.inner.clone(),
671 });
672 let result = taquba::run_worker_concurrent(
673 &self.inner.core.queue,
674 &self.inner.core.queue_name,
675 worker,
676 self.inner.core.max_concurrent_steps,
677 self.inner.core.poll_interval,
678 stop.clone().cancelled_owned(),
679 )
680 .await;
681 stop.cancel();
682
683 for handle in background {
684 let _ = handle.await;
685 }
686
687 result?;
688 Ok(())
689 }
690}
691
692pub type RunnerHandle = WorkerHandle<Result<()>>;
699
700impl RuntimeCore {
701 #[instrument(skip(self, spec), fields(run_id))]
703 pub(crate) async fn submit(&self, spec: RunSpec) -> Result<SubmitOutcome> {
704 let run_id = Self::validate_spec(&spec)?;
705 tracing::Span::current().record("run_id", run_id.as_str());
706 self.enqueue_run(&run_id, spec, None).await
707 }
708
709 pub(crate) async fn submit_member(
713 &self,
714 membership: &Membership,
715 spec: RunSpec,
716 ) -> Result<SubmitOutcome> {
717 let run_id = Self::validate_spec(&spec)?;
718 self.enqueue_run(&run_id, spec, Some(membership)).await
719 }
720
721 fn validate_spec(spec: &RunSpec) -> Result<RunId> {
724 for k in spec.options.headers.keys() {
725 if k.starts_with(RESERVED_HEADER_PREFIX) {
726 return Err(Error::ReservedHeaderInSubmit(k.clone()));
727 }
728 }
729 let deletes = spec.effects.kv_deletes.iter();
730 for key in spec.effects.kv_writes.keys().chain(deletes) {
731 if key.starts_with(RESERVED_KV_PREFIX.as_bytes()) {
732 return Err(Error::ReservedKvKey(
733 String::from_utf8_lossy(key).into_owned(),
734 ));
735 }
736 }
737 Ok(spec.run_id.clone().unwrap_or_else(RunId::generate))
738 }
739
740 async fn enqueue_run(
748 &self,
749 run_id: &RunId,
750 spec: RunSpec,
751 membership: Option<&Membership>,
752 ) -> Result<SubmitOutcome> {
753 let input_hash = hash_input(&spec.input);
754 let duplicate = |job_id: String| SubmitOutcome {
755 run_id: run_id.clone(),
756 newly_submitted: false,
757 job_id,
758 };
759 let check_input = |existing: DurableRunRecord| {
760 if existing.input_hash == input_hash {
761 Ok(())
762 } else {
763 Err(Error::InputMismatch(run_id.clone()))
764 }
765 };
766
767 if let Some(existing) = self.view.run_record(run_id).await? {
768 check_input(existing)?;
769 let current = self.current_step(run_id).await?;
770 return Ok(duplicate(current.job_id));
771 }
772
773 let opts = StepEnqueueOpts {
774 run_at: spec.options.run_at,
775 priority: spec.options.priority,
776 max_attempts: spec.options.max_attempts_per_step,
777 reserved_headers: membership
778 .map(Membership::reserved_headers)
779 .unwrap_or_default(),
780 };
781 let (request, job_id) =
782 self.step_enqueue_request(run_id, 0, spec.input, &spec.options.headers, opts);
783
784 let record_bytes = durable::encode(&DurableRunRecord {
785 run_id: run_id.clone(),
786 submitted_at_ms: self.clock.now_ms(),
787 input_hash,
788 cancel_requested: false,
789 });
790 let mut effects = spec
791 .effects
792 .kv_put(run_kv_key(run_id), record_bytes)
793 .kv_put(step_kv_key(run_id), current_step_bytes(0, &job_id));
794 if let Some(membership) = membership {
795 effects = effects.kv_put(
796 membership.kv_key(),
797 durable::encode(&pending_member(run_id)),
798 );
799 }
800
801 let job_id = match self
802 .queue
803 .enqueue_with_effects(&request.queue, request.payload, request.options, effects)
804 .await?
805 .0
806 {
807 EnqueueResult::New(id) => id,
808 EnqueueResult::AlreadyEnqueued(existing) => {
813 if let Some(record) = self.view.run_record(run_id).await? {
814 check_input(record)?;
815 }
816 return Ok(duplicate(existing));
817 }
818 };
819
820 debug!(run_id = %run_id, job_id = %job_id, "run submitted");
821 Ok(SubmitOutcome {
822 run_id: run_id.clone(),
823 newly_submitted: true,
824 job_id,
825 })
826 }
827
828 pub(crate) async fn wait(&self, run_id: &RunId) -> Result<RunEnd> {
830 self.wait_run(run_id)
831 .await?
832 .ok_or_else(|| Error::RunNotFound(run_id.clone()))
833 }
834
835 pub(crate) async fn wait_timeout(
837 &self,
838 run_id: &RunId,
839 timeout: Duration,
840 ) -> Result<Option<RunEnd>> {
841 match tokio::time::timeout(timeout, self.wait(run_id)).await {
842 Ok(end) => end.map(Some),
843 Err(_) => Ok(None),
844 }
845 }
846
847 pub(crate) async fn cancel(&self, run_id: &RunId) -> Result<bool> {
849 let Some(input_hash) = self.request_cancel(run_id).await? else {
850 return Ok(false);
851 };
852 loop {
858 let Some((_, job)) = self.view.current_job(run_id).await? else {
859 return Ok(false);
861 };
862 if job.status == JobStatus::Dead {
863 return Ok(false);
867 }
868 let claimed = ClaimedStep::parse(&job)?;
869 let outcome = claimed.cancelled(None);
873 let termination = self.termination(&outcome, None, input_hash);
874 let effects = self.terminate_collecting_effects(&outcome, &claimed, termination);
875 match self.queue.cancel_with(&job.id, effects).await?.0 {
876 taquba::CancelOutcome::Removed | taquba::CancelOutcome::Requested => {
877 return Ok(true);
878 }
879 taquba::CancelOutcome::NotFound => continue,
880 }
881 }
882 }
883
884 pub(crate) fn terminate_collecting_effects(
902 &self,
903 outcome: &RunOutcome,
904 terminal_step: &ClaimedStep<'_>,
905 termination: DurableTermination,
906 ) -> SettlementEffects {
907 let terminated_at_ms = termination.terminated_at_ms;
908 let kv_deletes = vec![run_kv_key(&outcome.run_id), step_kv_key(&outcome.run_id)];
909 let mut kv_writes = HashMap::new();
910 kv_writes.insert(
911 outcome_kv_key(&outcome.run_id),
912 durable::encode(&termination),
913 );
914 if let Some(membership) = &terminal_step.membership {
915 kv_writes.insert(
916 membership.kv_key(),
917 durable::encode(&terminated_member(&outcome.run_id, termination)),
918 );
919 }
920 let enqueues = if (self.observes)(outcome) {
921 vec![self.notification_enqueue_request(outcome, Some(terminal_step.job))]
922 } else {
923 Vec::new()
924 };
925 let effects = SettlementEffects::default()
926 .enqueues(enqueues)
927 .kv_writes(kv_writes)
928 .kv_deletes(kv_deletes);
929 match &self.memo_sweep {
930 Some(sweep) => sweep.mark(effects, &outcome.run_id, terminated_at_ms),
931 None => effects,
932 }
933 }
934
935 pub(crate) async fn reconcile_dead_steps(&self) -> Result<usize> {
950 const PAGE: usize = 256;
951 let mut terminated = 0usize;
952 let mut dead = std::pin::pin!(self.queue.view().jobs(
953 &self.queue_name,
954 JobStatus::Dead,
955 PAGE
956 ));
957 while let Some(job) = dead.try_next().await? {
958 if job.headers.contains_key(HEADER_TERMINAL) {
959 continue;
960 }
961 let Ok(claimed) = ClaimedStep::parse(&job) else {
962 continue;
963 };
964 let run_id = &claimed.run_id;
965 let current = self.view.current_step_if_active(run_id).await?;
966 if current.is_none_or(|current| current.job_id != job.id) {
967 continue;
968 }
969 let Some(record) = self.view.run_record(run_id).await? else {
973 warn!(run_id = %run_id, job_id = %job.id, "dead step has a current-step pointer but no run record");
974 continue;
975 };
976 let error = job
977 .last_error
978 .clone()
979 .unwrap_or_else(|| "step dead-lettered outside the worker".to_string());
980 let outcome = claimed.failed(error);
981 let termination = self.termination(&outcome, None, record.input_hash);
982 let effects = self.terminate_collecting_effects(&outcome, &claimed, termination);
983 self.queue.commit_effects(effects).await?;
984 warn!(run_id = %run_id, step_number = claimed.step_number, job_id = %job.id, "terminated a run whose step was dead-lettered outside the worker");
985 terminated += 1;
986 }
987 Ok(terminated)
988 }
989
990 async fn run_dead_step_reconciliation(&self, stop: CancellationToken) {
995 run_periodically(
996 self.poll_interval,
997 &stop,
998 None,
999 |reconciled_at: Option<i64>| async move {
1000 match self.queue.view().stats(&self.queue_name).await {
1001 Ok(stats) if reconciled_at != Some(stats.dead) => {
1002 match self.reconcile_dead_steps().await {
1003 Ok(_) => Some(stats.dead),
1004 Err(err) => {
1005 warn!("dead-step reconciliation failed: {err}");
1006 reconciled_at
1007 }
1008 }
1009 }
1010 Ok(_) => reconciled_at,
1011 Err(err) => {
1012 warn!("dead-step reconciliation could not read queue stats: {err}");
1013 reconciled_at
1014 }
1015 }
1016 },
1017 )
1018 .await;
1019 }
1020
1021 fn sweeps(&self) -> impl Iterator<Item = &Arc<Sweep>> {
1023 self.memo_sweep.iter().chain(self.group_sweep.iter())
1024 }
1025
1026 #[cfg(test)]
1029 pub(crate) async fn sweep_once(&self) -> Result<usize> {
1030 let mut removed = 0;
1031 for sweep in self.sweeps() {
1032 removed += sweep.pass(&self.queue, &*self.clock).await?;
1033 }
1034 Ok(removed)
1035 }
1036
1037 pub(crate) async fn current_step(&self, run_id: &RunId) -> Result<DurableCurrentStep> {
1039 self.view
1040 .current_step_if_active(run_id)
1041 .await?
1042 .ok_or_else(|| Error::InconsistentRunState(run_id.clone()))
1043 }
1044
1045 pub(crate) async fn wait_run(&self, run_id: &RunId) -> Result<Option<RunEnd>> {
1051 loop {
1052 let Some((current, _)) = self.view.current_job(run_id).await? else {
1053 return self.run_end(run_id).await;
1054 };
1055 match self.queue.wait_for_completion(¤t.job_id).await? {
1056 WaitOutcome::Done(_) | WaitOutcome::Cancelled | WaitOutcome::NotFound => {}
1059 WaitOutcome::Dead(_) => {
1060 let unreconciled = self
1065 .view
1066 .current_step_if_active(run_id)
1067 .await?
1068 .is_some_and(|step| step.job_id == current.job_id);
1069 if unreconciled {
1070 tokio::time::sleep(self.poll_interval).await;
1071 }
1072 }
1073 }
1074 }
1075 }
1076
1077 async fn run_end(&self, run_id: &RunId) -> Result<Option<RunEnd>> {
1081 let Some(termination) = self.view.terminal_record(run_id).await? else {
1082 return Ok(None);
1083 };
1084 let termination = RunTermination::from(termination);
1085 let outcome = self
1086 .view
1087 .run_result_of(run_id, &termination)
1088 .await?
1089 .map(|result| result.outcome);
1090 Ok(Some(RunEnd {
1091 termination,
1092 outcome,
1093 }))
1094 }
1095
1096 pub(crate) fn termination(
1099 &self,
1100 outcome: &RunOutcome,
1101 error_kind: Option<StepErrorKind>,
1102 input_hash: [u8; 32],
1103 ) -> DurableTermination {
1104 DurableTermination {
1105 status: outcome.status.into(),
1106 error: outcome.error.clone(),
1107 error_kind: error_kind.map(DurableErrorKind::from),
1108 final_step: outcome.final_step,
1109 terminated_at_ms: self.clock.now_ms(),
1110 input_hash,
1111 }
1112 }
1113
1114 pub(crate) async fn store_run_result(
1117 &self,
1118 outcome: &RunOutcome,
1119 termination: &DurableTermination,
1120 ) -> Result<()> {
1121 let record = DurableRunResult {
1122 termination: termination.clone(),
1123 outcome: DurableRunOutcome::from(outcome),
1124 };
1125 self.memo_store
1126 .new_run_memo(&outcome.run_id)
1127 .put(RUN_RESULT_MEMO_KEY, &durable::encode(&record))
1128 .await
1129 }
1130
1131 async fn request_cancel(&self, run_id: &RunId) -> Result<Option<[u8; 32]>> {
1135 let key = run_kv_key(run_id);
1136 loop {
1137 let Some(current) = self.queue.view().kv_get(&key).await? else {
1138 return Ok(None);
1139 };
1140 let mut record: DurableRunRecord = durable::decode(¤t)?;
1141 if record.cancel_requested {
1142 return Ok(Some(record.input_hash));
1143 }
1144 record.cancel_requested = true;
1145 if self
1146 .queue
1147 .kv_compare_put(&key, Some(¤t), &durable::encode(&record))
1148 .await?
1149 {
1150 return Ok(Some(record.input_hash));
1151 }
1152 }
1153 }
1154
1155 fn step_enqueue_request(
1160 &self,
1161 run_id: &RunId,
1162 step_number: u32,
1163 payload: Vec<u8>,
1164 user_headers: &HashMap<String, String>,
1165 opts: StepEnqueueOpts,
1166 ) -> (EnqueueRequest, String) {
1167 let job_id = self.queue.next_job_id();
1168 let mut headers = user_headers.clone();
1169 headers.insert(HEADER_RUN_ID.to_string(), run_id.to_string());
1170 headers.insert(HEADER_STEP.to_string(), step_number.to_string());
1171 for (key, value) in &opts.reserved_headers {
1172 headers.insert((*key).to_string(), value.clone());
1173 }
1174
1175 let request = EnqueueRequest {
1176 queue: self.queue_name.clone(),
1177 payload,
1178 options: EnqueueOptions::default()
1179 .headers(headers)
1180 .run_at(opts.run_at)
1181 .priority(opts.priority)
1182 .max_attempts(opts.max_attempts)
1183 .dedup_key(Some(format!("{DEDUP_PREFIX}{run_id}:{step_number}")))
1184 .id_override(Some(job_id.clone())),
1185 };
1186 (request, job_id)
1187 }
1188
1189 fn notification_enqueue_request(
1193 &self,
1194 outcome: &RunOutcome,
1195 terminal_step: Option<&JobRecord>,
1196 ) -> EnqueueRequest {
1197 let payload = durable::encode(&DurableRunOutcome::from(outcome));
1198 let mut headers = HashMap::new();
1199 headers.insert(HEADER_RUN_ID.to_string(), outcome.run_id.to_string());
1200 headers.insert(HEADER_TERMINAL.to_string(), "1".to_string());
1201 EnqueueRequest {
1202 queue: self.queue_name.clone(),
1203 payload,
1204 options: EnqueueOptions::default()
1205 .headers(headers)
1206 .priority(terminal_step.map(|job| job.priority))
1207 .max_attempts(terminal_step.map(|job| job.max_attempts))
1208 .dedup_key(Some(format!("{DEDUP_PREFIX}{}:terminal", outcome.run_id))),
1209 }
1210 }
1211
1212 pub(crate) fn run_at_after(&self, delay: Duration) -> SystemTime {
1215 UNIX_EPOCH + Duration::from_millis(self.clock.now_ms()) + delay
1216 }
1217
1218 pub(crate) async fn load_step_output(
1219 &self,
1220 run_id: &RunId,
1221 step_number: u32,
1222 step_payload: &[u8],
1223 ) -> Result<Option<(StepOutcome, StagedEffects)>> {
1224 let Some(bytes) = self
1225 .memo_store
1226 .get_step_output(run_id, step_number, step_payload)
1227 .await?
1228 else {
1229 return Ok(None);
1230 };
1231 let Some(record) = durable::decode_or_absent::<DurableStepOutcomeRecord>(
1232 &bytes,
1233 "step-output replay record",
1234 &format_args!("{run_id}/{step_number}"),
1235 ) else {
1236 return Ok(None);
1237 };
1238 let mut outcome = StepOutcome::from(record.outcome);
1239 match &mut outcome {
1240 StepOutcome::Continue {
1241 when: Trigger::After(delay),
1242 ..
1243 } => {
1244 *delay = remaining_delay(record.stored_at_ms, self.clock.now_ms(), *delay);
1245 }
1246 StepOutcome::Continue {
1247 when: Trigger::OnSignal { timeout, .. },
1248 ..
1249 } => {
1250 *timeout = remaining_delay(record.stored_at_ms, self.clock.now_ms(), *timeout);
1251 }
1252 _ => {}
1253 }
1254 Ok(Some((outcome, record.effects)))
1255 }
1256
1257 pub(crate) async fn store_step_output(
1258 &self,
1259 run_id: &RunId,
1260 step_number: u32,
1261 step_payload: &[u8],
1262 outcome: &StepOutcome,
1263 effects: &StagedEffects,
1264 ) -> Result<()> {
1265 let record = DurableStepOutcomeRecord {
1266 stored_at_ms: self.clock.now_ms(),
1267 outcome: DurableStepOutcome::from(outcome),
1268 effects: effects.clone(),
1269 };
1270 let bytes = rmp_serde::to_vec_named(&record)?;
1271 self.memo_store
1272 .put_step_output(run_id, step_number, step_payload, &bytes)
1273 .await
1274 }
1275
1276 pub(crate) async fn advance(
1280 &self,
1281 claimed: &ClaimedStep<'_>,
1282 payload: Vec<u8>,
1283 opts: StepEnqueueOpts,
1284 ) -> SettlementEffects {
1285 self.advance_with_kv(claimed, payload, opts, |_| HashMap::new())
1286 .await
1287 }
1288
1289 pub(crate) async fn advance_with_kv(
1293 &self,
1294 claimed: &ClaimedStep<'_>,
1295 payload: Vec<u8>,
1296 opts: StepEnqueueOpts,
1297 kv_writes: impl FnOnce(&str) -> HashMap<Vec<u8>, Vec<u8>>,
1298 ) -> SettlementEffects {
1299 let run_id = &claimed.run_id;
1300 let next_step = claimed.step_number + 1;
1301 let (request, next_job_id) =
1302 self.step_enqueue_request(run_id, next_step, payload, &claimed.headers, opts);
1303 let mut kv_writes = kv_writes(&next_job_id);
1304 kv_writes.insert(
1305 step_kv_key(run_id),
1306 current_step_bytes(next_step, &next_job_id),
1307 );
1308 SettlementEffects::default()
1309 .enqueues(vec![request])
1310 .kv_writes(kv_writes)
1311 }
1312}
1313
1314#[cfg(test)]
1315mod tests {
1316 use super::*;
1317 use crate::durable::DurableMember;
1318 use crate::effects::{EffectsHandle, TerminalEffects};
1319 use crate::group::GroupMember;
1320 use crate::keys::group_member_kv_key;
1321 use crate::keys::{TERMINAL_KV_PREFIX, signal_buf_kv_key, signal_wait_kv_key};
1322 use crate::runner::{Step, StepError};
1323 use crate::signal::SignalOutcome;
1324 use crate::terminal::NoopTerminalHook;
1325 use crate::terminal::TerminalStatus;
1326 use crate::test_util::{
1327 advance, fast_options, open_queue, open_queue_at, open_queue_at_with, open_queue_with, rid,
1328 };
1329 use crate::view::WorkflowView;
1330 use std::sync::Mutex as StdMutex;
1331 use std::sync::atomic::{AtomicU32, Ordering};
1332 use taquba::object_store::ObjectStoreExt;
1333 use taquba::object_store::memory::InMemory;
1334 use taquba::{Expired, ExpiryIndex};
1335 use taquba::{LeaseHandle, MockClock, OpenOptions, QueueConfig, QueueReader};
1336 use tokio::sync::oneshot;
1337
1338 struct ChannelHook {
1340 tx: tokio::sync::mpsc::UnboundedSender<RunOutcome>,
1341 }
1342
1343 impl TerminalHook for ChannelHook {
1344 async fn on_termination(
1345 &self,
1346 outcome: &RunOutcome,
1347 _effects: &TerminalEffects,
1348 ) -> std::result::Result<(), StepError> {
1349 let _ = self.tx.send(outcome.clone());
1350 Ok(())
1351 }
1352 }
1353
1354 struct ScriptedRunner {
1356 script: Arc<StdMutex<Vec<StepOutcome>>>,
1357 }
1358
1359 impl ScriptedRunner {
1360 fn new(steps: Vec<StepOutcome>) -> Self {
1361 Self {
1362 script: Arc::new(StdMutex::new(steps)),
1363 }
1364 }
1365 }
1366
1367 impl StepRunner for ScriptedRunner {
1368 async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1369 let next = self.script.lock().unwrap().remove(0);
1370 Ok(next)
1371 }
1372 }
1373
1374 struct FixedRunner {
1377 result: std::result::Result<StepOutcome, StepError>,
1378 calls: Arc<AtomicU32>,
1379 }
1380
1381 impl FixedRunner {
1382 fn new(result: std::result::Result<StepOutcome, StepError>) -> Self {
1383 Self {
1384 result,
1385 calls: Arc::new(AtomicU32::new(0)),
1386 }
1387 }
1388 }
1389
1390 impl StepRunner for FixedRunner {
1391 async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1392 self.calls.fetch_add(1, Ordering::SeqCst);
1393 self.result.clone()
1394 }
1395 }
1396
1397 struct PauseRunner;
1399
1400 impl StepRunner for PauseRunner {
1401 async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1402 std::future::pending().await
1403 }
1404 }
1405
1406 struct UnreachableRunner;
1408
1409 impl StepRunner for UnreachableRunner {
1410 async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1411 unreachable!("worker must not claim the step");
1412 }
1413 }
1414
1415 struct GatedRunner {
1418 claimed: Arc<tokio::sync::Notify>,
1419 release: tokio::sync::Mutex<Option<oneshot::Receiver<()>>>,
1420 result: std::result::Result<StepOutcome, StepError>,
1421 calls: Arc<AtomicU32>,
1422 }
1423
1424 struct Gate {
1426 claimed: Arc<tokio::sync::Notify>,
1427 release: StdMutex<Option<oneshot::Sender<()>>>,
1428 calls: Arc<AtomicU32>,
1429 }
1430
1431 impl GatedRunner {
1432 fn new(result: std::result::Result<StepOutcome, StepError>) -> (Self, Gate) {
1433 let claimed = Arc::new(tokio::sync::Notify::new());
1434 let calls = Arc::new(AtomicU32::new(0));
1435 let (release_tx, release_rx) = oneshot::channel();
1436 let runner = Self {
1437 claimed: claimed.clone(),
1438 release: tokio::sync::Mutex::new(Some(release_rx)),
1439 result,
1440 calls: calls.clone(),
1441 };
1442 let gate = Gate {
1443 claimed,
1444 release: StdMutex::new(Some(release_tx)),
1445 calls,
1446 };
1447 (runner, gate)
1448 }
1449 }
1450
1451 impl StepRunner for GatedRunner {
1452 async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1453 self.calls.fetch_add(1, Ordering::SeqCst);
1454 self.claimed.notify_one();
1455 let rx = self
1456 .release
1457 .lock()
1458 .await
1459 .take()
1460 .expect("gate consumed twice");
1461 let _ = rx.await;
1462 self.result.clone()
1463 }
1464 }
1465
1466 impl Gate {
1467 async fn claimed(&self) {
1469 tokio::time::timeout(Duration::from_secs(2), self.claimed.notified())
1470 .await
1471 .expect("runner reached gate");
1472 }
1473
1474 fn release(&self) {
1475 if let Some(tx) = self.release.lock().unwrap().take() {
1476 let _ = tx.send(());
1477 }
1478 }
1479 }
1480
1481 async fn terminal_markers(queue: &Queue) -> Vec<(RunId, u64)> {
1484 let page = queue
1485 .view()
1486 .kv_scan(TERMINAL_KV_PREFIX, .., 1_000)
1487 .await
1488 .unwrap();
1489 let index = ExpiryIndex::new(TERMINAL_KV_PREFIX);
1490 page.entries
1491 .iter()
1492 .map(|(key, _)| {
1493 let (at_ms, suffix) = index.parse(key).expect("well-formed marker key");
1494 let id = std::str::from_utf8(suffix).expect("run id");
1495 (RunId::new(id).expect("run id"), at_ms)
1496 })
1497 .collect()
1498 }
1499
1500 async fn terminal_status_of<R: StepRunner, H: TerminalHook>(
1503 runtime: &WorkflowRuntime<R, H>,
1504 run_id: &RunId,
1505 ) -> Option<TerminalStatus> {
1506 match runtime.status(run_id).await.unwrap().map(|s| s.state) {
1507 Some(RunState::Terminated(termination)) => Some(termination.status),
1508 _ => None,
1509 }
1510 }
1511
1512 fn spawn_runtime<R, H>(runtime: WorkflowRuntime<R, H>) -> oneshot::Sender<()>
1513 where
1514 R: StepRunner + 'static,
1515 H: TerminalHook + 'static,
1516 {
1517 let (tx, rx) = oneshot::channel::<()>();
1518 tokio::spawn(async move {
1519 let _ = runtime
1520 .run(async move {
1521 let _ = rx.await;
1522 })
1523 .await;
1524 });
1525 tx
1526 }
1527
1528 struct RenewingRunner {
1531 queue: Arc<Queue>,
1532 tx: tokio::sync::mpsc::UnboundedSender<u64>,
1533 }
1534
1535 impl StepRunner for RenewingRunner {
1536 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
1537 step.lease
1538 .ensure_at_least(Duration::from_secs(600))
1539 .map_err(|e| StepError::transient(e.to_string()))?;
1540 let expiry = self
1541 .queue
1542 .lease_expiry("workflow-steps", &step.job_id)
1543 .expect("a running step holds a lease");
1544 let _ = self.tx.send(expiry);
1545 Ok(StepOutcome::Succeed { result: Vec::new() })
1546 }
1547 }
1548
1549 #[tokio::test(start_paused = true)]
1550 async fn a_step_runner_extends_its_lease_through_the_step() {
1551 let base = 1_700_000_000_000;
1552 let (queue, store, _clock) = open_queue_at(base).await;
1553 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1554 let (hook_tx, mut hook_rx) = tokio::sync::mpsc::unbounded_channel();
1555 let runtime = WorkflowRuntime::builder(
1556 queue.clone(),
1557 store.clone(),
1558 RenewingRunner { queue, tx },
1559 ChannelHook { tx: hook_tx },
1560 )
1561 .build();
1562 let shutdown = spawn_runtime(runtime.clone());
1563
1564 runtime
1565 .submit(RunSpec {
1566 input: Vec::new(),
1567 ..Default::default()
1568 })
1569 .await
1570 .unwrap();
1571 let expiry = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1572 .await
1573 .unwrap()
1574 .unwrap();
1575 assert!(
1576 expiry >= base + 600_000,
1577 "the extension must reach the lease registry",
1578 );
1579 let outcome = tokio::time::timeout(Duration::from_secs(2), hook_rx.recv())
1580 .await
1581 .unwrap()
1582 .unwrap();
1583 assert_eq!(outcome.status, TerminalStatus::Succeeded);
1584
1585 let _ = shutdown.send(());
1586 }
1587
1588 #[tokio::test(start_paused = true)]
1589 async fn single_step_succeeds_and_writes_no_marker_without_retention() {
1590 let (queue, store) = open_queue().await;
1591 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1592 let runtime = WorkflowRuntime::builder(
1593 queue,
1594 store.clone(),
1595 ScriptedRunner::new(vec![StepOutcome::Succeed {
1596 result: b"done".to_vec(),
1597 }]),
1598 ChannelHook { tx },
1599 )
1600 .build();
1601 let shutdown = spawn_runtime(runtime.clone());
1602
1603 let handle = runtime
1604 .submit(RunSpec {
1605 input: b"in".to_vec(),
1606 ..Default::default()
1607 })
1608 .await
1609 .unwrap();
1610 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1611 .await
1612 .unwrap()
1613 .unwrap();
1614
1615 assert_eq!(outcome.run_id, handle.run_id);
1616 assert_eq!(outcome.status, TerminalStatus::Succeeded);
1617 assert_eq!(outcome.result.as_deref(), Some(b"done".as_slice()));
1618 assert_eq!(outcome.final_step, 0);
1619 assert_eq!(
1620 terminal_status_of(&runtime, &handle.run_id).await,
1621 Some(TerminalStatus::Succeeded),
1622 "the terminal record is written without retention",
1623 );
1624 assert!(terminal_markers(&runtime.inner.core.queue).await.is_empty());
1625
1626 let _ = shutdown.send(());
1627 }
1628
1629 #[tokio::test(start_paused = true)]
1630 async fn multi_step_run_advances_through_continue_with_its_headers() {
1631 let (queue, store) = open_queue().await;
1632 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1633 let runtime = WorkflowRuntime::builder(
1634 queue,
1635 store.clone(),
1636 ScriptedRunner::new(vec![
1637 StepOutcome::continue_now(b"step1".to_vec()),
1638 StepOutcome::continue_now(b"step2".to_vec()),
1639 StepOutcome::Succeed {
1640 result: b"final".to_vec(),
1641 },
1642 ]),
1643 ChannelHook { tx },
1644 )
1645 .build();
1646 let shutdown = spawn_runtime(runtime.clone());
1647
1648 let handle = runtime
1649 .submit(RunSpec {
1650 input: b"start".to_vec(),
1651 options: RunOptions {
1652 headers: HashMap::from([
1653 ("trace_id".to_string(), "abc-123".to_string()),
1654 ("tenant".to_string(), "acme".to_string()),
1655 ]),
1656 ..Default::default()
1657 },
1658 ..Default::default()
1659 })
1660 .await
1661 .unwrap();
1662 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1663 .await
1664 .unwrap()
1665 .unwrap();
1666
1667 assert_eq!(outcome.run_id, handle.run_id);
1668 assert_eq!(outcome.final_step, 2);
1669 assert_eq!(outcome.status, TerminalStatus::Succeeded);
1670 assert_eq!(outcome.result.as_deref(), Some(b"final".as_slice()));
1671 assert_eq!(outcome.headers.get("trace_id").unwrap(), "abc-123");
1672 assert_eq!(outcome.headers.get("tenant").unwrap(), "acme");
1673 assert!(!outcome.headers.contains_key(HEADER_RUN_ID));
1674 assert!(!outcome.headers.contains_key(HEADER_STEP));
1675
1676 let recorded =
1677 runtime.outcome(&handle.run_id).await.unwrap().expect(
1678 "the worker writes the run result record before the terminating settlement",
1679 );
1680 assert_eq!(recorded.status, TerminalStatus::Succeeded);
1681 assert_eq!(recorded.final_step, 2);
1682 assert_eq!(recorded.result.as_deref(), Some(b"final".as_slice()));
1683 assert_eq!(recorded.headers, outcome.headers);
1684
1685 let _ = shutdown.send(());
1686 }
1687
1688 #[tokio::test(start_paused = true)]
1689 async fn continue_after_delays_next_step_until_promotion() {
1690 let initial = 1_700_000_000_000u64;
1691 let (queue, store, clock) = open_queue_at(initial).await;
1692 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1693 let runtime = WorkflowRuntime::builder(
1694 queue.clone(),
1695 store.clone(),
1696 ScriptedRunner::new(vec![
1697 StepOutcome::continue_after(b"step1".to_vec(), Duration::from_secs(60)),
1698 StepOutcome::Succeed {
1699 result: b"final".to_vec(),
1700 },
1701 ]),
1702 ChannelHook { tx },
1703 )
1704 .build();
1705 let shutdown = spawn_runtime(runtime.clone());
1706
1707 let handle = runtime
1708 .submit(RunSpec {
1709 input: b"start".to_vec(),
1710 ..Default::default()
1711 })
1712 .await
1713 .unwrap();
1714
1715 assert!(
1718 tokio::time::timeout(Duration::from_millis(500), rx.recv())
1719 .await
1720 .is_err()
1721 );
1722 let stats = queue.view().stats("workflow-steps").await.unwrap();
1723 assert_eq!(stats.scheduled, 1);
1724
1725 advance(&clock, Duration::from_secs(61)).await;
1726 queue.promote_scheduled_now().await.unwrap();
1727
1728 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1729 .await
1730 .unwrap()
1731 .unwrap();
1732 assert_eq!(outcome.run_id, handle.run_id);
1733 assert_eq!(outcome.final_step, 1);
1734 assert_eq!(outcome.status, TerminalStatus::Succeeded);
1735
1736 let _ = shutdown.send(());
1737 }
1738
1739 type ObservedSignals = Arc<StdMutex<Vec<Option<Vec<u8>>>>>;
1740
1741 struct SignalProbe {
1744 correlation_key: String,
1745 timeout: Duration,
1746 observed: ObservedSignals,
1747 }
1748
1749 impl StepRunner for SignalProbe {
1750 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
1751 if step.step_number == 0 {
1752 Ok(StepOutcome::continue_on_signal(
1753 Vec::new(),
1754 self.correlation_key.clone(),
1755 self.timeout,
1756 ))
1757 } else {
1758 self.observed.lock().unwrap().push(step.signal.clone());
1759 Ok(StepOutcome::Succeed { result: Vec::new() })
1760 }
1761 }
1762 }
1763
1764 async fn wait_for_scheduled(queue: &Queue, count: i64) {
1765 for _ in 0..200 {
1766 if queue
1767 .view()
1768 .stats("workflow-steps")
1769 .await
1770 .unwrap()
1771 .scheduled
1772 == count
1773 {
1774 return;
1775 }
1776 tokio::time::sleep(Duration::from_millis(10)).await;
1777 }
1778 panic!("scheduled count never reached {count}");
1779 }
1780
1781 fn signal_probe_runtime(
1782 queue: Arc<Queue>,
1783 store: Arc<dyn taquba::object_store::ObjectStore>,
1784 correlation_key: &str,
1785 timeout: Duration,
1786 ) -> (
1787 WorkflowRuntime<SignalProbe, ChannelHook>,
1788 ObservedSignals,
1789 tokio::sync::mpsc::UnboundedReceiver<RunOutcome>,
1790 ) {
1791 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1792 let observed = Arc::new(StdMutex::new(Vec::new()));
1793 let runtime = WorkflowRuntime::builder(
1794 queue,
1795 store,
1796 SignalProbe {
1797 correlation_key: correlation_key.to_string(),
1798 timeout,
1799 observed: observed.clone(),
1800 },
1801 ChannelHook { tx },
1802 )
1803 .build();
1804 (runtime, observed, rx)
1805 }
1806
1807 #[tokio::test(start_paused = true)]
1808 async fn signal_wakes_waiting_run_early_with_payload() {
1809 let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
1810 let (runtime, observed, mut rx) =
1811 signal_probe_runtime(queue.clone(), store, "order-1", Duration::from_secs(3600));
1812 let shutdown = spawn_runtime(runtime.clone());
1813
1814 runtime
1815 .submit(RunSpec {
1816 input: Vec::new(),
1817 ..Default::default()
1818 })
1819 .await
1820 .unwrap();
1821 wait_for_scheduled(&queue, 1).await;
1822
1823 let outcome = runtime.signal("order-1", b"paid".to_vec()).await.unwrap();
1824 assert_eq!(outcome, SignalOutcome::Delivered);
1825
1826 let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1828 .await
1829 .unwrap()
1830 .unwrap();
1831 assert_eq!(terminal.status, TerminalStatus::Succeeded);
1832 assert_eq!(
1833 observed.lock().unwrap().as_slice(),
1834 &[Some(b"paid".to_vec())]
1835 );
1836
1837 assert!(
1839 queue
1840 .view()
1841 .kv_get(&signal_wait_kv_key("order-1"))
1842 .await
1843 .unwrap()
1844 .is_none()
1845 );
1846 assert!(
1847 queue
1848 .view()
1849 .kv_get(&signal_buf_kv_key("order-1"))
1850 .await
1851 .unwrap()
1852 .is_none()
1853 );
1854
1855 let _ = shutdown.send(());
1856 }
1857
1858 #[tokio::test(start_paused = true)]
1859 async fn a_run_waiting_on_a_signal_survives_a_close_and_reopen() {
1860 let store: Arc<dyn taquba::object_store::ObjectStore> = Arc::new(InMemory::new());
1861 let open = |store: Arc<dyn taquba::object_store::ObjectStore>| async move {
1862 Arc::new(
1863 Queue::open_with_options(
1864 store,
1865 "test",
1866 OpenOptions::default().clock(Arc::new(MockClock::new(1_700_000_000_000))),
1867 )
1868 .await
1869 .unwrap(),
1870 )
1871 };
1872
1873 let queue = open(store.clone()).await;
1874 let (runtime, _observed, _rx) = signal_probe_runtime(
1875 queue.clone(),
1876 store.clone(),
1877 "approval",
1878 Duration::from_secs(3600),
1879 );
1880 let (stop_tx, stop_rx) = oneshot::channel::<()>();
1881 let worker = tokio::spawn({
1882 let runtime = runtime.clone();
1883 async move {
1884 runtime
1885 .run(async move {
1886 let _ = stop_rx.await;
1887 })
1888 .await
1889 }
1890 });
1891 runtime
1892 .submit(RunSpec {
1893 input: Vec::new(),
1894 ..Default::default()
1895 })
1896 .await
1897 .unwrap();
1898 wait_for_scheduled(&queue, 1).await;
1899 let _ = stop_tx.send(());
1900 worker.await.unwrap().unwrap();
1901 drop(runtime);
1902 Arc::into_inner(queue)
1903 .expect("no other queue references at close")
1904 .close()
1905 .await
1906 .unwrap();
1907
1908 let queue = open(store.clone()).await;
1909 assert_eq!(
1910 queue
1911 .view()
1912 .stats("workflow-steps")
1913 .await
1914 .unwrap()
1915 .scheduled,
1916 1
1917 );
1918
1919 let (runtime, observed, mut rx) =
1920 signal_probe_runtime(queue.clone(), store, "approval", Duration::from_secs(3600));
1921 let shutdown = spawn_runtime(runtime.clone());
1922
1923 let delivery = runtime
1924 .signal("approval", b"approved".to_vec())
1925 .await
1926 .unwrap();
1927 assert_eq!(delivery, SignalOutcome::Delivered);
1928
1929 let terminal = tokio::time::timeout(Duration::from_secs(5), rx.recv())
1930 .await
1931 .unwrap()
1932 .unwrap();
1933 assert_eq!(terminal.status, TerminalStatus::Succeeded);
1934 assert_eq!(
1935 observed.lock().unwrap().as_slice(),
1936 &[Some(b"approved".to_vec())]
1937 );
1938 let _ = shutdown.send(());
1939 }
1940
1941 #[tokio::test(start_paused = true)]
1942 async fn signal_timeout_delivers_none() {
1943 let (queue, store, clock) = open_queue_at(1_700_000_000_000).await;
1944 let (runtime, observed, mut rx) =
1945 signal_probe_runtime(queue.clone(), store, "order-2", Duration::from_secs(60));
1946 let shutdown = spawn_runtime(runtime.clone());
1947
1948 runtime
1949 .submit(RunSpec {
1950 input: Vec::new(),
1951 ..Default::default()
1952 })
1953 .await
1954 .unwrap();
1955 wait_for_scheduled(&queue, 1).await;
1956
1957 advance(&clock, Duration::from_secs(61)).await;
1958 queue.promote_scheduled_now().await.unwrap();
1959
1960 let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1961 .await
1962 .unwrap()
1963 .unwrap();
1964 assert_eq!(terminal.status, TerminalStatus::Succeeded);
1965 assert_eq!(observed.lock().unwrap().as_slice(), &[None]);
1966
1967 assert!(
1968 queue
1969 .view()
1970 .kv_get(&signal_wait_kv_key("order-2"))
1971 .await
1972 .unwrap()
1973 .is_none()
1974 );
1975
1976 let _ = shutdown.send(());
1977 }
1978
1979 #[tokio::test(start_paused = true)]
1980 async fn a_buffered_signal_is_consumed_at_registration_and_a_later_one_replaces_it() {
1981 let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
1982 let (runtime, observed, mut rx) =
1983 signal_probe_runtime(queue.clone(), store, "order-3", Duration::from_secs(3600));
1984 let shutdown = spawn_runtime(runtime.clone());
1985
1986 assert_eq!(
1987 runtime.signal("order-3", b"first".to_vec()).await.unwrap(),
1988 SignalOutcome::Buffered
1989 );
1990 assert_eq!(
1991 runtime.signal("order-3", b"second".to_vec()).await.unwrap(),
1992 SignalOutcome::Buffered
1993 );
1994
1995 runtime
1996 .submit(RunSpec {
1997 input: Vec::new(),
1998 ..Default::default()
1999 })
2000 .await
2001 .unwrap();
2002
2003 let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2006 .await
2007 .unwrap()
2008 .unwrap();
2009 assert_eq!(terminal.status, TerminalStatus::Succeeded);
2010 assert_eq!(
2011 observed.lock().unwrap().as_slice(),
2012 &[Some(b"second".to_vec())]
2013 );
2014
2015 assert!(
2016 queue
2017 .view()
2018 .kv_get(&signal_buf_kv_key("order-3"))
2019 .await
2020 .unwrap()
2021 .is_none()
2022 );
2023
2024 let _ = shutdown.send(());
2025 }
2026
2027 #[tokio::test(start_paused = true)]
2028 async fn clear_signal_discards_buffered_signal() {
2029 let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
2030 let (runtime, _observed, _rx) =
2031 signal_probe_runtime(queue.clone(), store, "order-5", Duration::from_secs(60));
2032
2033 assert_eq!(
2034 runtime.signal("order-5", b"stale".to_vec()).await.unwrap(),
2035 SignalOutcome::Buffered
2036 );
2037 assert!(runtime.clear_signal("order-5").await.unwrap());
2038 assert!(!runtime.clear_signal("order-5").await.unwrap());
2039 assert!(
2040 queue
2041 .view()
2042 .kv_get(&signal_buf_kv_key("order-5"))
2043 .await
2044 .unwrap()
2045 .is_none()
2046 );
2047 }
2048
2049 #[tokio::test(start_paused = true)]
2050 async fn duplicate_waiter_registration_fails_the_run() {
2051 let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
2052 let (runtime, _observed, mut rx) =
2053 signal_probe_runtime(queue.clone(), store, "order-6", Duration::from_secs(3600));
2054 let shutdown = spawn_runtime(runtime.clone());
2055
2056 runtime
2057 .submit(RunSpec {
2058 run_id: Some(rid("run-a")),
2059 input: Vec::new(),
2060 ..Default::default()
2061 })
2062 .await
2063 .unwrap();
2064 wait_for_scheduled(&queue, 1).await;
2065
2066 runtime
2067 .submit(RunSpec {
2068 run_id: Some(rid("run-b")),
2069 input: Vec::new(),
2070 ..Default::default()
2071 })
2072 .await
2073 .unwrap();
2074
2075 let terminal = tokio::time::timeout(Duration::from_secs(5), rx.recv())
2076 .await
2077 .unwrap()
2078 .unwrap();
2079 assert_eq!(terminal.run_id, "run-b");
2080 assert_eq!(terminal.status, TerminalStatus::Failed);
2081 assert!(
2082 terminal
2083 .error
2084 .as_deref()
2085 .is_some_and(|e| e.contains("already registered"))
2086 );
2087 assert_eq!(
2088 queue.view().stats("workflow-steps").await.unwrap().dead,
2089 1,
2090 "the rejected registration dead-letters run-b's step",
2091 );
2092 assert!(
2093 queue
2094 .view()
2095 .kv_get(&run_kv_key(&rid("run-b")))
2096 .await
2097 .unwrap()
2098 .is_none(),
2099 "the run record delete rides the dead-letter",
2100 );
2101 assert!(
2102 queue
2103 .view()
2104 .kv_get(&run_kv_key(&rid("run-a")))
2105 .await
2106 .unwrap()
2107 .is_some(),
2108 "the waiting run keeps its record",
2109 );
2110
2111 let _ = shutdown.send(());
2112 }
2113
2114 #[tokio::test(start_paused = true)]
2115 async fn buffered_signal_missed_by_the_wake_is_delivered_at_timeout() {
2116 let (queue, store, clock) = open_queue_at(1_700_000_000_000).await;
2117 let (runtime, observed, mut rx) =
2118 signal_probe_runtime(queue.clone(), store, "order-7", Duration::from_secs(60));
2119 let shutdown = spawn_runtime(runtime.clone());
2120
2121 runtime
2122 .submit(RunSpec {
2123 input: Vec::new(),
2124 ..Default::default()
2125 })
2126 .await
2127 .unwrap();
2128 wait_for_scheduled(&queue, 1).await;
2129
2130 queue
2133 .kv_put(&signal_buf_kv_key("order-7"), b"late")
2134 .await
2135 .unwrap();
2136
2137 advance(&clock, Duration::from_secs(61)).await;
2138 queue.promote_scheduled_now().await.unwrap();
2139
2140 let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2141 .await
2142 .unwrap()
2143 .unwrap();
2144 assert_eq!(terminal.status, TerminalStatus::Succeeded);
2145 assert_eq!(
2146 observed.lock().unwrap().as_slice(),
2147 &[Some(b"late".to_vec())]
2148 );
2149
2150 let _ = shutdown.send(());
2151 }
2152
2153 #[tokio::test(start_paused = true)]
2154 async fn cancelled_waiter_leaves_no_live_index_for_the_next_signal() {
2155 let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
2156 let (runtime, _observed, mut rx) =
2157 signal_probe_runtime(queue.clone(), store, "order-8", Duration::from_secs(3600));
2158 let shutdown = spawn_runtime(runtime.clone());
2159
2160 let handle = runtime
2161 .submit(RunSpec {
2162 input: Vec::new(),
2163 ..Default::default()
2164 })
2165 .await
2166 .unwrap();
2167 wait_for_scheduled(&queue, 1).await;
2168
2169 assert!(runtime.cancel(&handle.run_id).await.unwrap());
2170 let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2171 .await
2172 .unwrap()
2173 .unwrap();
2174 assert_eq!(terminal.status, TerminalStatus::Cancelled);
2175
2176 assert_eq!(
2178 runtime.signal("order-8", b"orphan".to_vec()).await.unwrap(),
2179 SignalOutcome::Buffered
2180 );
2181 assert!(
2182 queue
2183 .view()
2184 .kv_get(&signal_wait_kv_key("order-8"))
2185 .await
2186 .unwrap()
2187 .is_none()
2188 );
2189
2190 let _ = shutdown.send(());
2191 }
2192
2193 #[tokio::test(start_paused = true)]
2194 async fn a_failure_notification_inherits_the_step_limits() {
2195 let (queue, store) = open_queue().await;
2196 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2197 let runtime = WorkflowRuntime::builder(
2198 queue.clone(),
2199 store,
2200 FixedRunner::new(Err(StepError::permanent("nope"))),
2201 ChannelHook { tx },
2202 )
2203 .build();
2204 runtime
2205 .submit(RunSpec {
2206 input: b"x".to_vec(),
2207 options: RunOptions {
2208 priority: Some(3),
2209 max_attempts_per_step: Some(5),
2210 ..Default::default()
2211 },
2212 ..Default::default()
2213 })
2214 .await
2215 .unwrap();
2216
2217 let job = queue
2218 .claim("workflow-steps", Duration::from_secs(30))
2219 .await
2220 .unwrap()
2221 .unwrap();
2222 let err = runtime
2223 .inner
2224 .process_step(&job, &LeaseHandle::detached())
2225 .await
2226 .unwrap_err();
2227 let failure = err
2228 .downcast_ref::<taquba::FailWith>()
2229 .expect("a terminating failure carries its effects");
2230 let notification = &failure.effects.enqueues[0];
2231 assert_eq!(notification.options.priority, Some(3));
2232 assert_eq!(notification.options.max_attempts, Some(5));
2233 }
2234
2235 #[tokio::test(start_paused = true)]
2236 async fn a_duplicate_submit_is_idempotent_drops_its_effects_and_rejects_a_changed_input() {
2237 let (queue, store) = open_queue().await;
2238 let runtime = WorkflowRuntime::builder(
2239 queue.clone(),
2240 store,
2241 ScriptedRunner::new(vec![]),
2242 NoopTerminalHook,
2243 )
2244 .build();
2245 let spec = |input: &[u8], key: &[u8]| RunSpec {
2248 run_id: Some(rid("fixed-id")),
2249 input: input.to_vec(),
2250 effects: SettlementEffects::default().kv_put(key, b"1"),
2251 ..Default::default()
2252 };
2253
2254 let first = runtime.submit(spec(b"x", b"app/first")).await.unwrap();
2255 assert!(first.newly_submitted);
2256 assert!(runtime.status(&rid("fixed-id")).await.unwrap().is_some());
2257 assert_eq!(
2258 queue.view().kv_get(b"app/first").await.unwrap().as_deref(),
2259 Some(b"1".as_slice())
2260 );
2261
2262 let duplicate = runtime.submit(spec(b"x", b"app/second")).await.unwrap();
2263 assert_eq!(duplicate.run_id, "fixed-id");
2264 assert!(!duplicate.newly_submitted);
2265 assert_eq!(duplicate.job_id, first.job_id);
2266 assert!(queue.view().kv_get(b"app/second").await.unwrap().is_none());
2267
2268 let err = runtime.submit(spec(b"y", b"app/third")).await.unwrap_err();
2269 assert!(matches!(&err, Error::InputMismatch(id) if id == "fixed-id"));
2270 assert!(err.is_permanent());
2271 assert!(queue.view().kv_get(b"app/third").await.unwrap().is_none());
2272 }
2273
2274 #[tokio::test(start_paused = true)]
2275 async fn a_duplicate_known_only_from_the_durable_record_reports_the_current_job() {
2276 let (queue, store) = open_queue().await;
2277 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2278 let first = WorkflowRuntime::builder(
2281 queue.clone(),
2282 store.clone(),
2283 ScriptedRunner::new(vec![StepOutcome::continue_after(
2284 b"next".to_vec(),
2285 Duration::from_secs(3600),
2286 )]),
2287 ChannelHook { tx },
2288 )
2289 .build();
2290 let shutdown = spawn_runtime(first.clone());
2291 let submitted = first
2292 .submit(RunSpec {
2293 run_id: Some(rid("durable")),
2294 input: b"x".to_vec(),
2295 ..Default::default()
2296 })
2297 .await
2298 .unwrap();
2299 for _ in 0..200 {
2300 if queue
2301 .view()
2302 .stats("workflow-steps")
2303 .await
2304 .unwrap()
2305 .scheduled
2306 == 1
2307 {
2308 break;
2309 }
2310 tokio::time::sleep(Duration::from_millis(10)).await;
2311 }
2312 let _ = shutdown.send(());
2313
2314 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2315 let second = WorkflowRuntime::builder(
2316 queue.clone(),
2317 store,
2318 ScriptedRunner::new(vec![]),
2319 ChannelHook { tx },
2320 )
2321 .build();
2322 let duplicate = second
2323 .submit(RunSpec {
2324 run_id: Some(rid("durable")),
2325 input: b"x".to_vec(),
2326 ..Default::default()
2327 })
2328 .await
2329 .unwrap();
2330 assert!(!duplicate.newly_submitted);
2331 assert_ne!(
2332 duplicate.job_id, submitted.job_id,
2333 "the pointer moved to step 1"
2334 );
2335 let step_1 = queue
2336 .view()
2337 .get_job(&duplicate.job_id)
2338 .await
2339 .unwrap()
2340 .unwrap();
2341 assert_eq!(step_1.status, taquba::JobStatus::Scheduled);
2342 assert_eq!(
2343 step_1.headers.get(HEADER_STEP).map(String::as_str),
2344 Some("1")
2345 );
2346 }
2347
2348 #[tokio::test(start_paused = true)]
2349 async fn a_run_submitted_with_run_at_stays_scheduled_until_then() {
2350 struct Echo;
2351 impl StepRunner for Echo {
2352 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
2353 Ok(StepOutcome::Succeed {
2354 result: step.payload.clone(),
2355 })
2356 }
2357 }
2358
2359 let t0 = 1_700_000_000_000;
2360 let (queue, store, clock) = open_queue_at(t0).await;
2361 let runtime =
2362 WorkflowRuntime::builder(queue.clone(), store, Echo, NoopTerminalHook).build();
2363 let shutdown = spawn_runtime(runtime.clone());
2364
2365 runtime
2366 .submit(RunSpec {
2367 input: b"x".to_vec(),
2368 options: RunOptions {
2369 run_at: Some(UNIX_EPOCH + Duration::from_millis(t0 + 60_000)),
2370 ..Default::default()
2371 },
2372 ..Default::default()
2373 })
2374 .await
2375 .unwrap();
2376 let scheduled = queue
2377 .view()
2378 .list_jobs("workflow-steps", taquba::JobStatus::Scheduled, None, 10)
2379 .await
2380 .unwrap()
2381 .jobs;
2382 assert_eq!(scheduled.len(), 1);
2383 let job_id = scheduled[0].id.clone();
2384
2385 let waiter = tokio::spawn({
2386 let queue = queue.clone();
2387 async move { queue.wait_for_completion(&job_id).await }
2388 });
2389 advance(&clock, Duration::from_secs(120)).await;
2390 assert!(matches!(
2391 waiter.await.unwrap().unwrap(),
2392 taquba::WaitOutcome::Done(_),
2393 ));
2394
2395 let _ = shutdown.send(());
2396 }
2397
2398 #[tokio::test(start_paused = true)]
2399 async fn a_step_reports_its_attempt_limit() {
2400 struct Recording(Arc<std::sync::Mutex<Option<u32>>>);
2401 impl StepRunner for Recording {
2402 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
2403 *self.0.lock().unwrap() = Some(step.max_attempts);
2404 Ok(StepOutcome::Succeed { result: Vec::new() })
2405 }
2406 }
2407
2408 let seen = Arc::new(std::sync::Mutex::new(None));
2409 let (queue, store) = open_queue().await;
2410 let runtime = WorkflowRuntime::builder(
2411 queue.clone(),
2412 store,
2413 Recording(seen.clone()),
2414 NoopTerminalHook,
2415 )
2416 .build();
2417 let shutdown = spawn_runtime(runtime.clone());
2418
2419 let outcome = runtime
2420 .submit(RunSpec {
2421 input: b"x".to_vec(),
2422 options: RunOptions {
2423 max_attempts_per_step: Some(7),
2424 ..Default::default()
2425 },
2426 ..Default::default()
2427 })
2428 .await
2429 .unwrap();
2430 let job_id = outcome.job_id;
2431 queue.wait_for_completion(&job_id).await.unwrap();
2432 assert_eq!(*seen.lock().unwrap(), Some(7));
2433
2434 let _ = shutdown.send(());
2435 }
2436
2437 #[tokio::test(start_paused = true)]
2438 async fn concurrent_submits_of_one_run_admit_one_and_reject_a_changed_input() {
2439 let (queue, store) = open_queue().await;
2440 let runtime =
2441 WorkflowRuntime::builder(queue, store.clone(), PauseRunner, NoopTerminalHook).build();
2442 let spec = |input: &[u8]| RunSpec {
2443 run_id: Some(rid("raced")),
2444 input: input.to_vec(),
2445 ..Default::default()
2446 };
2447
2448 let (first, same, changed) = tokio::join!(
2449 runtime.submit(spec(b"x")),
2450 runtime.submit(spec(b"x")),
2451 runtime.submit(spec(b"y")),
2452 );
2453 let first = first.unwrap();
2454 let same = same.unwrap();
2455 assert!(first.newly_submitted);
2456 assert!(!same.newly_submitted);
2457 assert_eq!(same.job_id, first.job_id);
2458 assert!(matches!(changed, Err(Error::InputMismatch(id)) if id == "raced"));
2459 }
2460
2461 #[tokio::test(start_paused = true)]
2462 async fn restart_resumes_at_next_step() {
2463 struct CompleteOnStep1;
2473 impl StepRunner for CompleteOnStep1 {
2474 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
2475 assert_eq!(step.step_number, 1, "runtime B should only ever see step 1");
2476 assert_eq!(step.payload.as_slice(), b"step1-payload");
2477 Ok(StepOutcome::Succeed {
2478 result: b"resumed".to_vec(),
2479 })
2480 }
2481 }
2482
2483 let (queue, store) = open_queue().await;
2484
2485 let (runner, gate) =
2486 GatedRunner::new(Ok(StepOutcome::continue_now(b"step1-payload".to_vec())));
2487 let runtime_a =
2488 WorkflowRuntime::builder(queue.clone(), store.clone(), runner, NoopTerminalHook)
2489 .max_concurrent_steps(1)
2490 .build();
2491
2492 let (shutdown_a_tx, shutdown_a_rx) = oneshot::channel::<()>();
2493 let worker_a = {
2494 let runtime_a = runtime_a.clone();
2495 tokio::spawn(async move {
2496 let _ = runtime_a
2497 .run(async move {
2498 let _ = shutdown_a_rx.await;
2499 })
2500 .await;
2501 })
2502 };
2503
2504 let handle = runtime_a
2505 .submit(RunSpec {
2506 input: b"input".to_vec(),
2507 ..Default::default()
2508 })
2509 .await
2510 .unwrap();
2511
2512 gate.claimed().await;
2513 let s = runtime_a
2514 .status(&handle.run_id)
2515 .await
2516 .unwrap()
2517 .expect("status");
2518 assert_eq!(s.state, RunState::Running);
2519 assert_eq!(s.current_step, 0);
2520
2521 let _ = shutdown_a_tx.send(());
2525 gate.release();
2526
2527 worker_a.await.expect("runtime A drained cleanly");
2528
2529 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2532 let runtime_b =
2533 WorkflowRuntime::builder(queue, store.clone(), CompleteOnStep1, ChannelHook { tx })
2534 .build();
2535 let shutdown_b = spawn_runtime(runtime_b.clone());
2536
2537 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2538 .await
2539 .expect("hook fired in time")
2540 .expect("hook channel open");
2541
2542 assert_eq!(outcome.run_id, handle.run_id);
2543 assert_eq!(outcome.status, TerminalStatus::Succeeded);
2544 assert_eq!(outcome.result.as_deref(), Some(b"resumed".as_slice()));
2545 assert_eq!(outcome.final_step, 1);
2546
2547 let _ = shutdown_b.send(());
2548 }
2549
2550 #[test]
2551 fn remaining_delay_measures_from_stored_timestamp() {
2552 let delay = Duration::from_secs(10);
2553 assert_eq!(remaining_delay(1_000, 4_000, delay), Duration::from_secs(7));
2554 assert_eq!(remaining_delay(1_000, 20_000, delay), Duration::ZERO);
2555 assert_eq!(remaining_delay(5_000, 1_000, delay), delay);
2556 }
2557
2558 #[tokio::test(start_paused = true)]
2559 async fn step_output_replay_skips_runner_after_crash_before_ack() {
2560 let (queue, store) = open_queue().await;
2561 let calls = Arc::new(AtomicU32::new(0));
2562 let runtime = WorkflowRuntime::builder(
2563 queue.clone(),
2564 store.clone(),
2565 FixedRunner {
2566 result: Ok(StepOutcome::continue_now(b"step1-payload".to_vec())),
2567 calls: calls.clone(),
2568 },
2569 NoopTerminalHook,
2570 )
2571 .step_output_replay()
2572 .build();
2573
2574 runtime
2575 .submit(RunSpec {
2576 run_id: Some(rid("replay-run")),
2577 input: b"input".to_vec(),
2578 ..Default::default()
2579 })
2580 .await
2581 .unwrap();
2582
2583 let job = queue
2584 .claim("workflow-steps", Duration::from_secs(30))
2585 .await
2586 .unwrap()
2587 .unwrap();
2588
2589 let _ = runtime
2593 .inner
2594 .process_step(&job, &LeaseHandle::detached())
2595 .await
2596 .unwrap();
2597 assert_eq!(calls.load(Ordering::SeqCst), 1);
2598
2599 let effects = runtime
2602 .inner
2603 .process_step(&job, &LeaseHandle::detached())
2604 .await
2605 .unwrap();
2606 assert_eq!(calls.load(Ordering::SeqCst), 1);
2607 queue.ack_with(&job, effects).await.unwrap();
2608
2609 let next = queue
2610 .claim("workflow-steps", Duration::from_secs(30))
2611 .await
2612 .unwrap()
2613 .unwrap();
2614 assert_eq!(next.payload.as_slice(), b"step1-payload");
2615 assert_eq!(next.headers.get(HEADER_RUN_ID).unwrap(), "replay-run");
2616 assert_eq!(next.headers.get(HEADER_STEP).unwrap(), "1");
2617 assert!(
2618 queue
2619 .claim("workflow-steps", Duration::from_secs(30))
2620 .await
2621 .unwrap()
2622 .is_none(),
2623 "the replayed continue must enqueue step 1 exactly once",
2624 );
2625 }
2626
2627 #[tokio::test(start_paused = true)]
2628 async fn corrupt_step_output_replay_entry_falls_back_to_runner() {
2629 let (queue, store) = open_queue().await;
2630 let calls = Arc::new(AtomicU32::new(0));
2631 let runtime = WorkflowRuntime::builder(
2632 queue.clone(),
2633 store.clone(),
2634 FixedRunner {
2635 result: Ok(StepOutcome::continue_now(b"step1-payload".to_vec())),
2636 calls: calls.clone(),
2637 },
2638 NoopTerminalHook,
2639 )
2640 .step_output_replay()
2641 .build();
2642
2643 runtime
2644 .submit(RunSpec {
2645 run_id: Some(rid("corrupt-run")),
2646 input: b"input".to_vec(),
2647 ..Default::default()
2648 })
2649 .await
2650 .unwrap();
2651
2652 let job = queue
2653 .claim("workflow-steps", Duration::from_secs(30))
2654 .await
2655 .unwrap()
2656 .unwrap();
2657 runtime
2658 .inner
2659 .core
2660 .memo_store
2661 .put_step_output(&rid("corrupt-run"), 0, &job.payload, b"not msgpack")
2662 .await
2663 .unwrap();
2664
2665 runtime
2666 .inner
2667 .process_step(&job, &LeaseHandle::detached())
2668 .await
2669 .unwrap();
2670 assert_eq!(
2671 calls.load(Ordering::SeqCst),
2672 1,
2673 "corrupt entry is treated as a miss",
2674 );
2675
2676 runtime
2679 .inner
2680 .process_step(&job, &LeaseHandle::detached())
2681 .await
2682 .unwrap();
2683 assert_eq!(calls.load(Ordering::SeqCst), 1);
2684 }
2685
2686 #[tokio::test(start_paused = true)]
2687 async fn step_output_replay_of_terminal_outcome_skips_runner() {
2688 let (queue, store) = open_queue().await;
2689 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2690 let calls = Arc::new(AtomicU32::new(0));
2691 let runtime = WorkflowRuntime::builder(
2692 queue.clone(),
2693 store.clone(),
2694 FixedRunner {
2695 result: Ok(StepOutcome::Succeed {
2696 result: b"final".to_vec(),
2697 }),
2698 calls: calls.clone(),
2699 },
2700 ChannelHook { tx },
2701 )
2702 .step_output_replay()
2703 .build();
2704
2705 runtime
2706 .submit(RunSpec {
2707 run_id: Some(rid("terminal-replay")),
2708 input: b"input".to_vec(),
2709 ..Default::default()
2710 })
2711 .await
2712 .unwrap();
2713 let job = queue
2714 .claim("workflow-steps", Duration::from_secs(30))
2715 .await
2716 .unwrap()
2717 .unwrap();
2718
2719 runtime
2720 .inner
2721 .process_step(&job, &LeaseHandle::detached())
2722 .await
2723 .unwrap();
2724 assert_eq!(calls.load(Ordering::SeqCst), 1);
2725
2726 let effects = runtime
2729 .inner
2730 .process_step(&job, &LeaseHandle::detached())
2731 .await
2732 .unwrap();
2733 assert_eq!(calls.load(Ordering::SeqCst), 1);
2734 queue.ack_with(&job, effects).await.unwrap();
2735
2736 let notification = queue
2739 .claim("workflow-steps", Duration::from_secs(30))
2740 .await
2741 .unwrap()
2742 .unwrap();
2743 let effects = runtime
2744 .inner
2745 .process_step(¬ification, &LeaseHandle::detached())
2746 .await
2747 .unwrap();
2748 queue.ack_with(¬ification, effects).await.unwrap();
2749 let outcome = rx.recv().await.unwrap();
2750 assert_eq!(outcome.status, TerminalStatus::Succeeded);
2751 assert_eq!(outcome.result.as_deref(), Some(b"final".as_slice()));
2752 assert_eq!(calls.load(Ordering::SeqCst), 1);
2753 }
2754
2755 async fn assert_transient_retries_until_max(max_attempts: u32) {
2761 let (queue, store) = open_queue_with(fast_options()).await;
2762 let calls = Arc::new(AtomicU32::new(0));
2763 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2764 let runtime = WorkflowRuntime::builder(
2765 queue.clone(),
2766 store.clone(),
2767 FixedRunner {
2768 result: Err(StepError::transient("flaky")),
2769 calls: calls.clone(),
2770 },
2771 ChannelHook { tx },
2772 )
2773 .memo_retention(Duration::from_secs(60))
2774 .build();
2775 let shutdown = spawn_runtime(runtime.clone());
2776
2777 let handle = runtime
2778 .submit(RunSpec {
2779 input: b"x".to_vec(),
2780 options: RunOptions {
2781 max_attempts_per_step: Some(max_attempts),
2782 ..Default::default()
2783 },
2784 ..Default::default()
2785 })
2786 .await
2787 .unwrap();
2788
2789 let outcome = tokio::time::timeout(Duration::from_secs(3), rx.recv())
2790 .await
2791 .expect("hook fired in time")
2792 .expect("hook channel open");
2793
2794 assert_eq!(outcome.status, TerminalStatus::Failed);
2795 assert_eq!(outcome.error.as_deref(), Some("flaky"));
2796 assert_eq!(
2797 calls.load(Ordering::SeqCst),
2798 max_attempts,
2799 "runner called once per attempt up to max_attempts"
2800 );
2801
2802 tokio::time::sleep(Duration::from_millis(50)).await;
2804 assert!(rx.try_recv().is_err(), "hook fired more than once");
2805
2806 assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 1);
2809 assert!(
2810 queue
2811 .view()
2812 .kv_get(&run_kv_key(&handle.run_id))
2813 .await
2814 .unwrap()
2815 .is_none(),
2816 "the run record delete rides the exhausted nack",
2817 );
2818 assert_eq!(
2819 terminal_markers(&queue)
2820 .await
2821 .iter()
2822 .filter(|(run_id, _)| *run_id == handle.run_id)
2823 .count(),
2824 1,
2825 "the terminal marker rides the exhausted nack",
2826 );
2827
2828 let _ = shutdown.send(());
2829 }
2830
2831 #[tokio::test(start_paused = true)]
2832 async fn a_cancellation_survives_a_restart() {
2833 let (queue, store, _clock) = open_queue_at(10_000).await;
2838 let (tx_a, _rx_a) = tokio::sync::mpsc::unbounded_channel();
2839 let before = WorkflowRuntime::builder(
2840 queue.clone(),
2841 store.clone(),
2842 ScriptedRunner::new(vec![StepOutcome::Succeed {
2843 result: b"done".to_vec(),
2844 }]),
2845 ChannelHook { tx: tx_a },
2846 )
2847 .build();
2848
2849 let handle = before
2850 .submit(RunSpec {
2851 input: b"x".to_vec(),
2852 ..Default::default()
2853 })
2854 .await
2855 .unwrap();
2856 let claim = queue
2857 .claim("workflow-steps", Duration::from_secs(30))
2858 .await
2859 .unwrap()
2860 .expect("step 0 is claimable");
2861 assert!(before.cancel(&handle.run_id).await.unwrap());
2862
2863 let (tx_b, mut rx_b) = tokio::sync::mpsc::unbounded_channel();
2864 let after = WorkflowRuntime::builder(
2865 queue.clone(),
2866 store.clone(),
2867 ScriptedRunner::new(vec![StepOutcome::Succeed {
2868 result: b"done".to_vec(),
2869 }]),
2870 ChannelHook { tx: tx_b },
2871 )
2872 .build();
2873 assert_eq!(
2874 after.status(&handle.run_id).await.unwrap().map(|s| s.state),
2875 Some(RunState::Cancelling),
2876 "the fresh runtime reads the request from the run record",
2877 );
2878
2879 let effects = after
2880 .inner
2881 .process_step(&claim, &queue.lease_handle(&claim))
2882 .await
2883 .unwrap();
2884 queue.ack_with(&claim, effects).await.unwrap();
2885
2886 let notification = queue
2887 .claim("workflow-steps", Duration::from_secs(30))
2888 .await
2889 .unwrap()
2890 .expect("the terminal notification is claimable");
2891 let effects = after
2892 .inner
2893 .process_step(¬ification, &LeaseHandle::detached())
2894 .await
2895 .unwrap();
2896 queue.ack_with(¬ification, effects).await.unwrap();
2897
2898 let outcome = rx_b.recv().await.unwrap();
2899 assert_eq!(outcome.status, TerminalStatus::Cancelled);
2900 assert!(outcome.result.is_none(), "the succeed payload is discarded");
2901 }
2902
2903 #[tokio::test(start_paused = true)]
2904 async fn a_cancellation_after_the_settlement_read_reaches_the_next_step() {
2905 let (queue, store, _clock) = open_queue_at(10_000).await;
2910 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2911 let runtime = WorkflowRuntime::builder(
2912 queue.clone(),
2913 store.clone(),
2914 ScriptedRunner::new(vec![
2915 StepOutcome::Continue {
2916 payload: b"next".to_vec(),
2917 when: Trigger::Immediate,
2918 },
2919 StepOutcome::Succeed {
2920 result: b"done".to_vec(),
2921 },
2922 ]),
2923 ChannelHook { tx },
2924 )
2925 .build();
2926
2927 let handle = runtime
2928 .submit(RunSpec {
2929 input: b"x".to_vec(),
2930 ..Default::default()
2931 })
2932 .await
2933 .unwrap();
2934 let step0 = queue
2935 .claim("workflow-steps", Duration::from_secs(30))
2936 .await
2937 .unwrap()
2938 .expect("step 0 is claimable");
2939 let effects = runtime
2940 .inner
2941 .process_step(&step0, &queue.lease_handle(&step0))
2942 .await
2943 .unwrap();
2944 assert!(runtime.cancel(&handle.run_id).await.unwrap());
2945 queue.ack_with(&step0, effects).await.unwrap();
2946
2947 let step1 = queue
2948 .claim("workflow-steps", Duration::from_secs(30))
2949 .await
2950 .unwrap()
2951 .expect("step 1 is claimable");
2952 let effects = runtime
2953 .inner
2954 .process_step(&step1, &queue.lease_handle(&step1))
2955 .await
2956 .unwrap();
2957 queue.ack_with(&step1, effects).await.unwrap();
2958
2959 let notification = queue
2960 .claim("workflow-steps", Duration::from_secs(30))
2961 .await
2962 .unwrap()
2963 .expect("the terminal notification is claimable");
2964 let effects = runtime
2965 .inner
2966 .process_step(¬ification, &LeaseHandle::detached())
2967 .await
2968 .unwrap();
2969 queue.ack_with(¬ification, effects).await.unwrap();
2970
2971 let outcome = rx.recv().await.unwrap();
2972 assert_eq!(outcome.status, TerminalStatus::Cancelled);
2973 assert_eq!(outcome.final_step, 1);
2974 assert_eq!(
2975 terminal_status_of(&runtime, &handle.run_id).await,
2976 Some(outcome.status)
2977 );
2978 }
2979
2980 #[tokio::test(start_paused = true)]
2981 async fn cancelling_a_pending_run_commits_its_marker_and_fires_the_hook_once() {
2982 let (queue, store, _clock) = open_queue_at(10_000).await;
2987 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2988 let runtime = WorkflowRuntime::builder(
2989 queue.clone(),
2990 store.clone(),
2991 UnreachableRunner,
2992 ChannelHook { tx },
2993 )
2994 .memo_retention(Duration::from_secs(60))
2995 .build();
2996 let mut headers = HashMap::new();
3000 headers.insert("tenant".to_string(), "acme".to_string());
3001
3002 let handle = runtime
3003 .submit(RunSpec {
3004 input: b"x".to_vec(),
3005 options: RunOptions {
3006 headers,
3007 ..Default::default()
3008 },
3009 ..Default::default()
3010 })
3011 .await
3012 .unwrap();
3013 let status = runtime
3014 .status(&handle.run_id)
3015 .await
3016 .unwrap()
3017 .expect("active");
3018 assert_eq!(status.state, RunState::Pending);
3019
3020 let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3021 assert!(was_cancelled);
3022 let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3023 assert_eq!(
3024 status.state,
3025 RunState::Terminated(RunTermination {
3026 status: TerminalStatus::Cancelled,
3027 error: None,
3028 error_kind: None,
3029 final_step: 0,
3030 terminated_at_ms: 10_000,
3031 }),
3032 "the terminal record commits with the removal",
3033 );
3034 assert_eq!(status.current_step, 0);
3035 assert!(
3036 runtime.outcome(&handle.run_id).await.unwrap().is_none(),
3037 "no worker terminated the run, so no run result record exists",
3038 );
3039 assert!(
3040 !runtime.cancel(&handle.run_id).await.unwrap(),
3041 "a second cancel finds no run record",
3042 );
3043
3044 let markers = terminal_markers(&queue).await;
3047 assert_eq!(markers, vec![(handle.run_id.clone(), 10_000)]);
3048 assert_eq!(
3049 queue
3050 .view()
3051 .kv_get(&run_kv_key(&handle.run_id))
3052 .await
3053 .unwrap(),
3054 None,
3055 );
3056
3057 let notification = queue
3059 .claim("workflow-steps", Duration::from_secs(30))
3060 .await
3061 .unwrap()
3062 .unwrap();
3063 let effects = runtime
3064 .inner
3065 .process_step(¬ification, &LeaseHandle::detached())
3066 .await
3067 .unwrap();
3068 queue.ack_with(¬ification, effects).await.unwrap();
3069 let outcome = rx.recv().await.unwrap();
3070 assert_eq!(outcome.run_id, handle.run_id);
3071 assert_eq!(outcome.status, TerminalStatus::Cancelled);
3072 assert!(outcome.error.is_none());
3074 assert_eq!(outcome.headers.get("tenant").unwrap(), "acme");
3075 assert!(
3076 queue
3077 .claim("workflow-steps", Duration::from_secs(30))
3078 .await
3079 .unwrap()
3080 .is_none(),
3081 "the cancel enqueues one notification",
3082 );
3083 assert!(rx.try_recv().is_err());
3084
3085 let stats = queue.view().stats("workflow-steps").await.unwrap();
3086 assert_eq!(stats.dead, 0, "cancel must not dead-letter");
3087 assert_eq!(stats.pending, 0, "cancelled job must be removed");
3088 }
3089
3090 #[tokio::test(start_paused = true)]
3091 async fn a_view_reads_the_same_state_through_the_runtime_and_a_reader() {
3092 let (queue, store, _clock) = open_queue_at(10_000).await;
3093 let runtime = WorkflowRuntime::builder(
3094 queue.clone(),
3095 store.clone(),
3096 UnreachableRunner,
3097 NoopTerminalHook,
3098 )
3099 .memo_prefix("memo")
3100 .build();
3101 let handle = runtime.submit(RunSpec::default()).await.unwrap();
3103
3104 let reader = QueueReader::open(store.clone(), "test").await.unwrap();
3105 let view = WorkflowView::new(reader.view().clone(), MemoStore::new(store.clone(), "memo"));
3106 let through_reader = view.status(&handle.run_id).await.unwrap().expect("active");
3107 let through_runtime = runtime.status(&handle.run_id).await.unwrap().unwrap();
3108 assert_eq!(through_reader.run_id, handle.run_id);
3109 assert_eq!(through_reader.state, RunState::Pending);
3110 assert_eq!(through_reader.state, through_runtime.state);
3111 assert_eq!(through_reader.current_step, through_runtime.current_step);
3112 assert!(view.outcome(&handle.run_id).await.unwrap().is_none());
3113 assert!(view.status(&rid("unknown")).await.unwrap().is_none());
3114
3115 assert!(runtime.cancel(&handle.run_id).await.unwrap());
3116 let reader = QueueReader::open(store.clone(), "test").await.unwrap();
3117 let view = WorkflowView::new(reader.view().clone(), MemoStore::new(store, "memo"));
3118 let status = view
3119 .status(&handle.run_id)
3120 .await
3121 .unwrap()
3122 .expect("terminal record");
3123 assert_eq!(
3124 status.state,
3125 RunState::Terminated(RunTermination {
3126 status: TerminalStatus::Cancelled,
3127 error: None,
3128 error_kind: None,
3129 final_step: 0,
3130 terminated_at_ms: 10_000,
3131 }),
3132 );
3133 }
3134
3135 #[tokio::test(start_paused = true)]
3136 async fn the_status_of_a_run_does_not_read_its_offloaded_step_payload() {
3137 let payloads: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
3138 let (queue, store, _clock) = open_queue_at_with(
3139 10_000,
3140 OpenOptions::default()
3141 .payload_offload_threshold(64)
3142 .payload_store(payloads.clone()),
3143 )
3144 .await;
3145 let runtime = WorkflowRuntime::builder(queue, store, UnreachableRunner, NoopTerminalHook)
3146 .memo_prefix("memo")
3147 .build();
3148 let handle = runtime
3149 .submit(RunSpec {
3150 input: vec![7u8; 512],
3151 ..Default::default()
3152 })
3153 .await
3154 .unwrap();
3155
3156 let objects: Vec<_> = payloads.list(None).try_collect().await.unwrap();
3159 assert!(!objects.is_empty(), "the step payload is offloaded");
3160 for object in objects {
3161 payloads.delete(&object.location).await.unwrap();
3162 }
3163 let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3164 assert_eq!(status.state, RunState::Pending);
3165 }
3166
3167 async fn assert_cancel_suppresses_runner_error(error: StepError) {
3175 let (queue, store) = open_queue_with(fast_options()).await;
3176 let (runner, gate) = GatedRunner::new(Err(error));
3177 let (hook_tx, mut hook_rx) = tokio::sync::mpsc::unbounded_channel();
3178 let runtime = WorkflowRuntime::builder(
3179 queue.clone(),
3180 store.clone(),
3181 runner,
3182 ChannelHook { tx: hook_tx },
3183 )
3184 .build();
3185 let shutdown = spawn_runtime(runtime.clone());
3186
3187 let handle = runtime
3188 .submit(RunSpec {
3189 input: b"x".to_vec(),
3190 ..Default::default()
3191 })
3192 .await
3193 .unwrap();
3194 gate.claimed().await;
3195
3196 let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3197 assert!(was_cancelled);
3198
3199 gate.release();
3203
3204 let outcome = tokio::time::timeout(Duration::from_secs(2), hook_rx.recv())
3205 .await
3206 .expect("hook fired")
3207 .expect("hook channel open");
3208 assert_eq!(outcome.status, TerminalStatus::Cancelled);
3209 assert!(
3210 outcome.error.is_none(),
3211 "external cancel must carry no reason (Some(_) would imply runner-issued StepOutcome::Cancel)",
3212 );
3213 assert_eq!(
3214 terminal_status_of(&runtime, &handle.run_id).await,
3215 Some(outcome.status)
3216 );
3217
3218 tokio::time::sleep(Duration::from_millis(100)).await;
3221 assert_eq!(
3222 gate.calls.load(Ordering::SeqCst),
3223 1,
3224 "cancellation must suppress retries",
3225 );
3226 let stats = queue.view().stats("workflow-steps").await.unwrap();
3227 assert_eq!(stats.dead, 0, "cancellation must suppress dead-letter");
3228 assert!(
3229 hook_rx.try_recv().is_err(),
3230 "hook must fire exactly once for the cancelled run",
3231 );
3232
3233 let _ = shutdown.send(());
3234 }
3235
3236 #[tokio::test(start_paused = true)]
3237 async fn cancel_suppresses_a_runner_error() {
3238 assert_cancel_suppresses_runner_error(StepError::permanent("would-dead-letter")).await;
3243 assert_cancel_suppresses_runner_error(StepError::transient("would-retry")).await;
3244 }
3245
3246 #[tokio::test(start_paused = true)]
3247 async fn cancel_signals_step_token_for_cooperative_short_circuit() {
3248 struct CooperativeRunner {
3256 claimed: Arc<tokio::sync::Notify>,
3257 }
3258 impl StepRunner for CooperativeRunner {
3259 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
3260 self.claimed.notify_one();
3261 tokio::select! {
3262 _ = tokio::time::sleep(Duration::from_secs(30)) => {
3263 Ok(StepOutcome::Succeed { result: b"slow".to_vec() })
3264 }
3265 _ = step.cancel_token.cancelled() => {
3266 Ok(StepOutcome::Cancel { reason: "cooperative".to_string() })
3267 }
3268 }
3269 }
3270 }
3271
3272 let (queue, store) = open_queue().await;
3273 let claimed = Arc::new(tokio::sync::Notify::new());
3274 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3275 let runtime = WorkflowRuntime::builder(
3276 queue.clone(),
3277 store.clone(),
3278 CooperativeRunner {
3279 claimed: claimed.clone(),
3280 },
3281 ChannelHook { tx },
3282 )
3283 .build();
3284 let shutdown = spawn_runtime(runtime.clone());
3285
3286 let handle = runtime
3287 .submit(RunSpec {
3288 input: b"x".to_vec(),
3289 ..Default::default()
3290 })
3291 .await
3292 .unwrap();
3293 tokio::time::timeout(Duration::from_secs(2), claimed.notified())
3294 .await
3295 .expect("runner observed token");
3296
3297 let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3298 assert!(was_cancelled);
3299
3300 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3301 .await
3302 .expect("hook fired well before the 30s sleep would have")
3303 .expect("hook channel open");
3304
3305 assert_eq!(outcome.status, TerminalStatus::Cancelled);
3306 assert_eq!(outcome.error.as_deref(), Some("cooperative"));
3309 assert_eq!(
3310 terminal_status_of(&runtime, &handle.run_id).await,
3311 Some(outcome.status)
3312 );
3313
3314 let stats = queue.view().stats("workflow-steps").await.unwrap();
3315 assert_eq!(stats.dead, 0);
3316
3317 let _ = shutdown.send(());
3318 }
3319
3320 #[tokio::test(start_paused = true)]
3321 async fn cancel_returns_false_for_a_terminated_or_unknown_run() {
3322 let (queue, store) = open_queue().await;
3327 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3328 let runtime = WorkflowRuntime::builder(
3329 queue,
3330 store.clone(),
3331 ScriptedRunner::new(vec![StepOutcome::Succeed {
3332 result: b"done".to_vec(),
3333 }]),
3334 ChannelHook { tx },
3335 )
3336 .build();
3337 let shutdown = spawn_runtime(runtime.clone());
3338
3339 let handle = runtime
3340 .submit(RunSpec {
3341 input: b"x".to_vec(),
3342 ..Default::default()
3343 })
3344 .await
3345 .unwrap();
3346
3347 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3348 .await
3349 .expect("Succeeded hook fired")
3350 .expect("hook channel open");
3351 assert_eq!(outcome.status, TerminalStatus::Succeeded);
3352 assert_eq!(
3353 terminal_status_of(&runtime, &handle.run_id).await,
3354 Some(outcome.status)
3355 );
3356
3357 let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3358 assert!(
3359 !was_cancelled,
3360 "cancel on an already-terminated run must report Ok(false)",
3361 );
3362
3363 tokio::time::sleep(Duration::from_millis(50)).await;
3364 assert!(
3365 rx.try_recv().is_err(),
3366 "no Cancelled hook may fire after the run already terminated as Succeeded",
3367 );
3368
3369 assert!(!runtime.cancel(&rid("never-submitted")).await.unwrap());
3370
3371 let _ = shutdown.send(());
3372 }
3373
3374 #[tokio::test(start_paused = true)]
3375 async fn transient_retries_until_max_attempts() {
3376 assert_transient_retries_until_max(1).await;
3377 assert_transient_retries_until_max(3).await;
3378 }
3379
3380 #[tokio::test(start_paused = true)]
3381 async fn step_memo_survives_across_attempts_of_the_same_step() {
3382 struct MemoRetryRunner;
3388 impl StepRunner for MemoRetryRunner {
3389 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
3390 if step.attempts == 1 {
3391 step.memo
3392 .put("cached", b"first-attempt-value")
3393 .await
3394 .map_err(|e| StepError::transient(e.to_string()))?;
3395 return Err(StepError::transient("force a retry"));
3396 }
3397 let got = step
3398 .memo
3399 .get("cached")
3400 .await
3401 .map_err(|e| StepError::transient(e.to_string()))?;
3402 assert_eq!(got, Some(b"first-attempt-value".to_vec()));
3403 Ok(StepOutcome::Succeed {
3404 result: got.unwrap_or_default(),
3405 })
3406 }
3407 }
3408
3409 let (queue, store) = open_queue_with(fast_options()).await;
3410 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3411 let runtime =
3412 WorkflowRuntime::builder(queue, store, MemoRetryRunner, ChannelHook { tx }).build();
3413 let shutdown = spawn_runtime(runtime.clone());
3414
3415 runtime
3416 .submit(RunSpec {
3417 input: b"start".to_vec(),
3418 options: RunOptions {
3419 max_attempts_per_step: Some(3),
3420 ..Default::default()
3421 },
3422 ..Default::default()
3423 })
3424 .await
3425 .unwrap();
3426 let outcome = tokio::time::timeout(Duration::from_secs(3), rx.recv())
3427 .await
3428 .expect("hook fired in time")
3429 .expect("hook channel open");
3430 assert_eq!(outcome.status, TerminalStatus::Succeeded);
3431 assert_eq!(
3432 outcome.result.as_deref(),
3433 Some(b"first-attempt-value".as_slice())
3434 );
3435
3436 let _ = shutdown.send(());
3437 }
3438
3439 #[tokio::test(start_paused = true)]
3440 async fn terminal_marker_is_written_at_the_runtime_clock() {
3441 let (queue, store, clock) = open_queue_at(10_000).await;
3445 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3446 let runtime = WorkflowRuntime::builder(
3447 queue,
3448 store.clone(),
3449 ScriptedRunner::new(vec![StepOutcome::Succeed {
3450 result: b"done".to_vec(),
3451 }]),
3452 ChannelHook { tx },
3453 )
3454 .memo_retention(Duration::from_secs(60))
3455 .build();
3456 let shutdown = spawn_runtime(runtime.clone());
3457
3458 let handle = runtime
3459 .submit(RunSpec {
3460 input: b"in".to_vec(),
3461 ..Default::default()
3462 })
3463 .await
3464 .unwrap();
3465 advance(&clock, Duration::from_secs(30)).await;
3466 let _ = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3467 .await
3468 .unwrap()
3469 .unwrap();
3470
3471 let markers = terminal_markers(&runtime.inner.core.queue).await;
3472 assert_eq!(markers.len(), 1);
3473 assert_eq!(markers[0].0, handle.run_id);
3474 assert_eq!(markers[0].1, 10_000 + 30_000);
3477
3478 let _ = shutdown.send(());
3479 }
3480
3481 #[tokio::test(start_paused = true)]
3482 async fn submit_rejects_reserved_headers_and_reserved_kv_keys() {
3483 let (queue, store, _clock) = open_queue_at(10_000).await;
3484 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3485 let runtime = WorkflowRuntime::builder(
3486 queue,
3487 store,
3488 ScriptedRunner::new(vec![]),
3489 ChannelHook { tx },
3490 )
3491 .build();
3492
3493 let err = runtime
3494 .submit(RunSpec {
3495 input: b"x".to_vec(),
3496 options: RunOptions {
3497 headers: HashMap::from([("workflow.run_id".to_string(), "evil".to_string())]),
3498 ..Default::default()
3499 },
3500 ..Default::default()
3501 })
3502 .await
3503 .unwrap_err();
3504 assert!(
3505 matches!(&err, Error::ReservedHeaderInSubmit(k) if k == "workflow.run_id"),
3506 "got: {err:?}"
3507 );
3508
3509 let err = runtime
3510 .submit(RunSpec {
3511 input: Vec::new(),
3512 effects: SettlementEffects::default().kv_put(b"workflow/x", b"v"),
3513 ..Default::default()
3514 })
3515 .await
3516 .unwrap_err();
3517 assert!(matches!(err, Error::ReservedKvKey(_)));
3518
3519 let err = runtime
3520 .submit(RunSpec {
3521 input: Vec::new(),
3522 effects: SettlementEffects::default().kv_delete(b"workflow/x"),
3523 ..Default::default()
3524 })
3525 .await
3526 .unwrap_err();
3527 assert!(matches!(err, Error::ReservedKvKey(_)));
3528 }
3529
3530 #[tokio::test(start_paused = true)]
3531 async fn submit_applies_the_deletes_and_expiry_entries_of_the_spec() {
3532 let (queue, store, _clock) = open_queue_at(10_000).await;
3533 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3534 let runtime = WorkflowRuntime::builder(
3535 queue.clone(),
3536 store,
3537 ScriptedRunner::new(vec![]),
3538 ChannelHook { tx },
3539 )
3540 .build();
3541 let index = ExpiryIndex::new(b"app/expiry/".to_vec());
3542 queue.kv_put(b"app/stale", b"1").await.unwrap();
3543 queue
3546 .commit_effects(SettlementEffects::default().expiry_entry(&index, 9_000, b"later"))
3547 .await
3548 .unwrap();
3549
3550 runtime
3551 .submit(RunSpec {
3552 input: b"x".to_vec(),
3553 effects: SettlementEffects::default()
3554 .kv_delete(b"app/stale")
3555 .expiry_entry(&index, 8_000, b"run"),
3556 ..Default::default()
3557 })
3558 .await
3559 .unwrap();
3560
3561 assert!(queue.view().kv_get(b"app/stale").await.unwrap().is_none());
3562 let mut seen = Vec::new();
3563 index
3564 .pass(&queue, 10_000, Duration::ZERO, |at_ms, suffix| {
3565 seen.push((at_ms, suffix));
3566 std::future::ready(Expired::Delete(SettlementEffects::default()))
3567 })
3568 .await
3569 .unwrap();
3570 assert_eq!(seen, [(8_000, b"run".to_vec()), (9_000, b"later".to_vec())]);
3571 }
3572
3573 #[tokio::test(start_paused = true)]
3574 async fn a_malformed_terminal_marker_is_deleted_without_clearing_memos() {
3575 let (queue, store, clock) = open_queue_at(10_000).await;
3576 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3577 let runtime = WorkflowRuntime::builder(
3578 queue.clone(),
3579 store.clone(),
3580 ScriptedRunner::new(vec![]),
3581 ChannelHook { tx },
3582 )
3583 .memo_retention(Duration::from_secs(60))
3584 .build();
3585
3586 let memos = MemoStore::new(store, "workflow-steps-memo");
3587 memos
3588 .new_memo(&rid("bystander"), 0)
3589 .put("k", b"expensive")
3590 .await
3591 .unwrap();
3592 let marker = ExpiryIndex::new(TERMINAL_KV_PREFIX).entry_key(0, b"");
3594 queue.kv_put(&marker, b"").await.unwrap();
3595 let unparseable = [TERMINAL_KV_PREFIX, b"short"].concat();
3597 queue.kv_put(&unparseable, b"").await.unwrap();
3598
3599 advance(&clock, Duration::from_secs(3_600)).await;
3600 runtime.inner.core.sweep_once().await.unwrap();
3601 assert_eq!(
3602 memos.new_memo(&rid("bystander"), 0).get("k").await.unwrap(),
3603 Some(b"expensive".to_vec()),
3604 "an unrelated run's memo entries must survive",
3605 );
3606 assert!(
3607 queue.view().kv_get(&marker).await.unwrap().is_none(),
3608 "the marker is removed and not retried on every sweep",
3609 );
3610 assert!(queue.view().kv_get(&unparseable).await.unwrap().is_none());
3611 }
3612
3613 #[tokio::test(start_paused = true)]
3614 async fn cancelling_a_running_step_overrides_its_outcome_and_writes_its_marker_at_settlement() {
3615 let (queue, store, _clock) = open_queue_at(10_000).await;
3620 let (runner, gate) = GatedRunner::new(Ok(StepOutcome::Succeed {
3621 result: b"would-have-succeeded".to_vec(),
3622 }));
3623 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3624 let runtime =
3625 WorkflowRuntime::builder(queue.clone(), store.clone(), runner, ChannelHook { tx })
3626 .memo_retention(Duration::from_secs(60))
3627 .build();
3628 let shutdown = spawn_runtime(runtime.clone());
3629
3630 let handle = runtime
3631 .submit(RunSpec {
3632 input: b"x".to_vec(),
3633 ..Default::default()
3634 })
3635 .await
3636 .unwrap();
3637 gate.claimed().await;
3638 assert_eq!(
3639 runtime
3640 .status(&handle.run_id)
3641 .await
3642 .unwrap()
3643 .expect("active")
3644 .state,
3645 RunState::Running
3646 );
3647
3648 assert!(runtime.cancel(&handle.run_id).await.unwrap());
3649 assert_eq!(
3650 runtime
3651 .status(&handle.run_id)
3652 .await
3653 .unwrap()
3654 .expect("entry retained while termination is in flight")
3655 .state,
3656 RunState::Cancelling
3657 );
3658 assert!(
3659 terminal_markers(&queue).await.is_empty(),
3660 "a run still executing its step must have no terminal marker",
3661 );
3662 assert!(
3663 queue
3664 .view()
3665 .kv_get(&run_kv_key(&handle.run_id))
3666 .await
3667 .unwrap()
3668 .is_some(),
3669 "the run record must survive a cancel the worker has to finish",
3670 );
3671
3672 gate.release();
3673 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3674 .await
3675 .expect("hook fired")
3676 .expect("hook channel open");
3677 assert_eq!(outcome.status, TerminalStatus::Cancelled);
3678 assert!(
3679 outcome.result.is_none(),
3680 "succeed payload must be discarded"
3681 );
3682 let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3683 assert!(matches!(
3684 status.state,
3685 RunState::Terminated(RunTermination {
3686 status: TerminalStatus::Cancelled,
3687 error: None,
3688 ..
3689 })
3690 ));
3691 assert_eq!(
3692 terminal_markers(&queue).await,
3693 vec![(handle.run_id.clone(), 10_000)],
3694 "the worker's settlement writes it",
3695 );
3696 assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 0);
3697
3698 let _ = shutdown.send(());
3699 }
3700
3701 #[tokio::test(start_paused = true)]
3702 async fn a_permanent_step_error_dead_letters_with_its_marker_and_no_staged_effects() {
3703 let (queue, store, _clock) = open_queue_at(10_000).await;
3704 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3705 let runtime = WorkflowRuntime::builder(
3706 queue.clone(),
3707 store,
3708 EffectStagingRunner::new(vec![Err(StepError::permanent("nope"))]),
3709 ChannelHook { tx },
3710 )
3711 .memo_retention(Duration::from_secs(60))
3712 .build();
3713 let shutdown = spawn_runtime(runtime.clone());
3714
3715 let handle = runtime
3716 .submit(RunSpec {
3717 input: b"x".to_vec(),
3718 ..Default::default()
3719 })
3720 .await
3721 .unwrap();
3722 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3723 .await
3724 .unwrap()
3725 .unwrap();
3726 assert_eq!(outcome.run_id, handle.run_id);
3727 assert_eq!(outcome.status, TerminalStatus::Failed);
3728 assert_eq!(outcome.error.as_deref(), Some("nope"));
3729 let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3730 assert!(
3731 matches!(
3732 status.state,
3733 RunState::Terminated(RunTermination {
3734 status: TerminalStatus::Failed,
3735 error: Some(ref error),
3736 error_kind: Some(StepErrorKind::Permanent),
3737 ..
3738 }) if error == "nope"
3739 ),
3740 "the terminal record commits with the dead-letter and carries the error kind",
3741 );
3742 let recorded = runtime.outcome(&handle.run_id).await.unwrap().unwrap();
3743 assert_eq!(recorded.status, TerminalStatus::Failed);
3744 assert_eq!(recorded.error.as_deref(), Some("nope"));
3745
3746 assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 1);
3750 assert_eq!(
3751 terminal_markers(&queue).await,
3752 vec![(handle.run_id.clone(), 10_000)],
3753 );
3754 assert!(queue.view().kv_get(b"app/step-0").await.unwrap().is_none());
3755
3756 let _ = shutdown.send(());
3757 }
3758
3759 #[tokio::test(start_paused = true)]
3760 async fn a_retrying_step_error_commits_no_terminal_marker() {
3761 struct FlakyRunner {
3764 attempts: Arc<std::sync::atomic::AtomicUsize>,
3765 clock: MockClock,
3766 }
3767 impl StepRunner for FlakyRunner {
3768 async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
3769 if self
3770 .attempts
3771 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
3772 == 0
3773 {
3774 return Err(StepError::transient("flaky"));
3775 }
3776 self.clock.advance(Duration::from_secs(1));
3780 Ok(StepOutcome::Succeed {
3781 result: b"done".to_vec(),
3782 })
3783 }
3784 }
3785
3786 let (queue, store, clock) = open_queue_at_with(10_000, fast_options()).await;
3787 let attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
3788 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3789 let runtime = WorkflowRuntime::builder(
3790 queue.clone(),
3791 store.clone(),
3792 FlakyRunner {
3793 attempts: attempts.clone(),
3794 clock: clock.clone(),
3795 },
3796 ChannelHook { tx },
3797 )
3798 .memo_retention(Duration::from_secs(60))
3799 .build();
3800 let shutdown = spawn_runtime(runtime.clone());
3801
3802 let handle = runtime
3803 .submit(RunSpec {
3804 input: b"x".to_vec(),
3805 options: RunOptions {
3806 max_attempts_per_step: Some(3),
3807 ..Default::default()
3808 },
3809 ..Default::default()
3810 })
3811 .await
3812 .unwrap();
3813 let outcome = tokio::time::timeout(Duration::from_secs(5), rx.recv())
3814 .await
3815 .unwrap()
3816 .unwrap();
3817 assert_eq!(outcome.status, TerminalStatus::Succeeded);
3818 assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 2);
3819
3820 let markers = terminal_markers(&queue).await;
3823 assert_eq!(markers.len(), 1);
3824 assert_eq!(markers[0].0, handle.run_id);
3825 assert_eq!(markers[0].1, 11_000);
3826
3827 let _ = shutdown.send(());
3828 }
3829
3830 #[tokio::test(start_paused = true)]
3831 async fn the_sweep_clears_only_markers_older_than_the_cutoff() {
3832 let (queue, store, _clock) = open_queue_at(10_000).await;
3835 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3836 let runtime = WorkflowRuntime::builder(
3837 queue.clone(),
3838 store.clone(),
3839 ScriptedRunner::new(vec![]),
3840 ChannelHook { tx },
3841 )
3842 .memo_retention(Duration::from_secs(1))
3843 .build();
3844
3845 let memos = MemoStore::new(store, "workflow-steps-memo");
3846 let sweep = runtime.inner.core.memo_sweep.as_ref().unwrap();
3847 for (run_id, at_ms) in [("old", 1_000u64), ("young", 9_500u64)] {
3848 let run_id = rid(run_id);
3849 memos.new_memo(&run_id, 0).put("k", b"v").await.unwrap();
3850 queue
3851 .commit_effects(sweep.mark(SettlementEffects::default(), &run_id, at_ms))
3852 .await
3853 .unwrap();
3854 }
3855
3856 let cleared = runtime.inner.core.sweep_once().await.unwrap();
3858 assert_eq!(cleared, 1);
3859
3860 let remaining = terminal_markers(&queue).await;
3861 assert_eq!(remaining.len(), 1);
3862 assert_eq!(remaining[0].0, "young");
3863 assert_eq!(memos.new_memo(&rid("old"), 0).get("k").await.unwrap(), None);
3864 assert_eq!(
3865 memos.new_memo(&rid("young"), 0).get("k").await.unwrap(),
3866 Some(b"v".to_vec()),
3867 );
3868 }
3869
3870 #[tokio::test(start_paused = true)]
3871 async fn a_re_submitted_run_shares_its_entries_until_the_first_run_expires() {
3872 let (queue, store, clock) = open_queue_at(10_000).await;
3873 let runtime = WorkflowRuntime::builder(
3874 queue.clone(),
3875 store.clone(),
3876 FixedRunner::new(Ok(StepOutcome::Succeed {
3877 result: b"done".to_vec(),
3878 })),
3879 NoopTerminalHook,
3880 )
3881 .memo_retention(Duration::from_secs(1))
3882 .build();
3883 let spec = RunSpec {
3884 run_id: Some(rid("shared")),
3885 input: b"x".to_vec(),
3886 ..Default::default()
3887 };
3888 runtime.submit(spec.clone()).await.unwrap();
3889 let claim = queue
3890 .claim("workflow-steps", Duration::from_secs(30))
3891 .await
3892 .unwrap()
3893 .unwrap();
3894 let effects = runtime
3895 .inner
3896 .process_step(&claim, &LeaseHandle::detached())
3897 .await
3898 .unwrap();
3899 queue.ack_with(&claim, effects).await.unwrap();
3900 let memos = MemoStore::new(store, "workflow-steps-memo");
3901 memos
3902 .new_memo(&rid("shared"), 0)
3903 .put("k", b"v")
3904 .await
3905 .unwrap();
3906
3907 assert!(runtime.submit(spec).await.unwrap().newly_submitted);
3908 assert_eq!(
3909 memos.new_memo(&rid("shared"), 0).get("k").await.unwrap(),
3910 Some(b"v".to_vec()),
3911 "the second run reads the first run's entry",
3912 );
3913
3914 clock.advance(Duration::from_secs(2));
3916 assert_eq!(runtime.inner.core.sweep_once().await.unwrap(), 1);
3917 assert_eq!(
3918 memos.new_memo(&rid("shared"), 0).get("k").await.unwrap(),
3919 None
3920 );
3921 assert_eq!(
3922 runtime
3923 .status(&rid("shared"))
3924 .await
3925 .unwrap()
3926 .map(|s| s.state),
3927 Some(RunState::Pending),
3928 "the second run is still active",
3929 );
3930 }
3931
3932 async fn yield_until<F, Fut>(iters: usize, mut cond: F) -> bool
3937 where
3938 F: FnMut() -> Fut,
3939 Fut: Future<Output = bool>,
3940 {
3941 for _ in 0..iters {
3942 if cond().await {
3943 return true;
3944 }
3945 tokio::task::yield_now().await;
3946 }
3947 false
3948 }
3949
3950 #[tokio::test(start_paused = true)]
3951 async fn the_sweeper_clears_a_marker_only_after_retention_elapses() {
3952 let (queue, store, clock) = open_queue_at(10_000).await;
3957 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3958 let runtime = WorkflowRuntime::builder(
3959 queue,
3960 store.clone(),
3961 ScriptedRunner::new(vec![StepOutcome::Succeed {
3962 result: b"done".to_vec(),
3963 }]),
3964 ChannelHook { tx },
3965 )
3966 .memo_retention(Duration::from_millis(200))
3967 .poll_interval(Duration::from_millis(10))
3968 .build();
3969 let shutdown = spawn_runtime(runtime.clone());
3970
3971 let handle = runtime
3972 .submit(RunSpec {
3973 input: b"in".to_vec(),
3974 ..Default::default()
3975 })
3976 .await
3977 .unwrap();
3978 let _ = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3979 .await
3980 .unwrap()
3981 .unwrap();
3982 let memos = MemoStore::new(store.clone(), "workflow-steps-memo");
3983 memos
3984 .new_memo(&handle.run_id, 0)
3985 .put("k", b"cached")
3986 .await
3987 .unwrap();
3988
3989 advance(&clock, Duration::from_millis(199)).await;
3990 assert_eq!(runtime.inner.core.sweep_once().await.unwrap(), 0);
3991 let markers = terminal_markers(&runtime.inner.core.queue).await;
3992 assert_eq!(
3993 markers.len(),
3994 1,
3995 "a marker within the window must not be swept"
3996 );
3997
3998 advance(&clock, Duration::from_millis(1)).await;
3999 advance(&clock, Duration::from_millis(10)).await;
4000 let cleared = yield_until(50, || async {
4001 terminal_markers(&runtime.inner.core.queue).await.is_empty()
4002 })
4003 .await;
4004 assert!(cleared, "sweeper did not clear the expired marker");
4005 assert_eq!(
4006 memos.new_memo(&handle.run_id, 0).get("k").await.unwrap(),
4007 None,
4008 "sweeper did not clear the run's memo entries",
4009 );
4010 assert!(
4011 runtime.status(&handle.run_id).await.unwrap().is_none(),
4012 "sweeper did not clear the run's terminal record",
4013 );
4014
4015 let _ = shutdown.send(());
4016 }
4017
4018 #[tokio::test(start_paused = true)]
4019 async fn sweeper_keeps_memos_of_runs_without_a_terminal_marker() {
4020 let (queue, store, clock) = open_queue_at(10_000).await;
4027 let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
4028 let runtime = WorkflowRuntime::builder(
4029 queue,
4030 store.clone(),
4031 ScriptedRunner::new(vec![]),
4032 ChannelHook { tx },
4033 )
4034 .memo_retention(Duration::from_millis(100))
4035 .build();
4036 let shutdown = spawn_runtime(runtime.clone());
4037
4038 let memos = MemoStore::new(store.clone(), "workflow-steps-memo");
4039 memos
4040 .new_memo(&rid("in-flight-run"), 0)
4041 .put("k", b"cached")
4042 .await
4043 .unwrap();
4044
4045 advance(&clock, Duration::from_millis(500)).await;
4046 for _ in 0..50 {
4048 tokio::task::yield_now().await;
4049 }
4050
4051 assert_eq!(
4052 memos
4053 .new_memo(&rid("in-flight-run"), 0)
4054 .get("k")
4055 .await
4056 .unwrap(),
4057 Some(b"cached".to_vec()),
4058 "sweep must not remove memos of a run with no terminal marker",
4059 );
4060
4061 let _ = shutdown.send(());
4062 }
4063
4064 async fn wait_for_kv(queue: &Queue, key: &[u8]) -> Vec<u8> {
4065 for _ in 0..200 {
4066 if let Some(v) = queue.view().kv_get(key).await.unwrap() {
4067 return v.to_vec();
4068 }
4069 tokio::time::sleep(Duration::from_millis(10)).await;
4070 }
4071 panic!(
4072 "kv key `{}` was never written",
4073 String::from_utf8_lossy(key)
4074 );
4075 }
4076
4077 async fn wait_for_drained(queue: &Queue) {
4078 for _ in 0..200 {
4079 let stats = queue.view().stats("workflow-steps").await.unwrap();
4080 if stats.pending == 0 && stats.claimed == 0 && stats.scheduled == 0 {
4081 return;
4082 }
4083 tokio::time::sleep(Duration::from_millis(10)).await;
4084 }
4085 panic!("the queue never drained");
4086 }
4087
4088 struct EffectStagingRunner {
4091 script: Arc<StdMutex<Vec<std::result::Result<StepOutcome, StepError>>>>,
4092 }
4093
4094 impl EffectStagingRunner {
4095 fn new(script: Vec<std::result::Result<StepOutcome, StepError>>) -> Self {
4096 Self {
4097 script: Arc::new(StdMutex::new(script)),
4098 }
4099 }
4100 }
4101
4102 impl StepRunner for EffectStagingRunner {
4103 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4104 step.effects
4105 .put(format!("app/step-{}", step.step_number), b"done".to_vec())
4106 .map_err(|e| StepError::permanent(e.to_string()))?;
4107 self.script.lock().unwrap().remove(0)
4108 }
4109 }
4110
4111 #[tokio::test(start_paused = true)]
4112 async fn step_effects_commit_with_the_acking_settlement() {
4113 let (queue, store) = open_queue().await;
4114 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4115 let runtime = WorkflowRuntime::builder(
4116 queue.clone(),
4117 store,
4118 EffectStagingRunner::new(vec![
4119 Ok(StepOutcome::continue_now(b"next".to_vec())),
4120 Ok(StepOutcome::Succeed { result: Vec::new() }),
4121 ]),
4122 ChannelHook { tx },
4123 )
4124 .build();
4125 let shutdown = spawn_runtime(runtime.clone());
4126
4127 runtime
4128 .submit(RunSpec {
4129 input: Vec::new(),
4130 ..Default::default()
4131 })
4132 .await
4133 .unwrap();
4134 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4135 .await
4136 .unwrap()
4137 .unwrap();
4138 assert_eq!(outcome.status, TerminalStatus::Succeeded);
4139 assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4140 assert_eq!(wait_for_kv(&queue, b"app/step-1").await, b"done");
4141
4142 let _ = shutdown.send(());
4143 }
4144
4145 #[tokio::test(start_paused = true)]
4146 async fn a_step_effect_is_readable_by_the_next_step() {
4147 struct ReadingRunner {
4148 read_under_staging: Arc<StdMutex<Option<Option<Vec<u8>>>>>,
4149 }
4150 impl StepRunner for ReadingRunner {
4151 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4152 let read = step
4153 .kv
4154 .get(b"app/marker")
4155 .await
4156 .map_err(|e| StepError::permanent(e.to_string()))?
4157 .map(|b| b.to_vec());
4158 if step.step_number == 0 {
4159 step.effects
4160 .put("app/marker", b"v".to_vec())
4161 .map_err(|e| StepError::permanent(e.to_string()))?;
4162 let staged_read = step
4163 .kv
4164 .get(b"app/marker")
4165 .await
4166 .map_err(|e| StepError::permanent(e.to_string()))?
4167 .map(|b| b.to_vec());
4168 *self.read_under_staging.lock().unwrap() = Some(staged_read);
4169 return Ok(StepOutcome::continue_now(Vec::new()));
4170 }
4171 Ok(StepOutcome::Succeed {
4172 result: read.unwrap_or_default(),
4173 })
4174 }
4175 }
4176
4177 let (queue, store) = open_queue().await;
4178 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4179 let read_under_staging = Arc::new(StdMutex::new(None));
4180 let runtime = WorkflowRuntime::builder(
4181 queue.clone(),
4182 store,
4183 ReadingRunner {
4184 read_under_staging: read_under_staging.clone(),
4185 },
4186 ChannelHook { tx },
4187 )
4188 .build();
4189 let shutdown = spawn_runtime(runtime.clone());
4190
4191 runtime
4192 .submit(RunSpec {
4193 input: Vec::new(),
4194 ..Default::default()
4195 })
4196 .await
4197 .unwrap();
4198 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4199 .await
4200 .unwrap()
4201 .unwrap();
4202 assert_eq!(outcome.status, TerminalStatus::Succeeded);
4203 assert_eq!(outcome.result.as_deref(), Some(b"v".as_slice()));
4204 assert_eq!(*read_under_staging.lock().unwrap(), Some(None));
4205
4206 let _ = shutdown.send(());
4207 }
4208
4209 #[tokio::test(start_paused = true)]
4210 async fn a_run_memo_written_in_one_step_is_readable_in_the_next() {
4211 struct JournalRunner;
4212 impl StepRunner for JournalRunner {
4213 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4214 if step.step_number == 0 {
4215 step.run_memo.put("journal", b"entry").await?;
4216 return Ok(StepOutcome::continue_now(Vec::new()));
4217 }
4218 let value = step.run_memo.get("journal").await?;
4219 Ok(StepOutcome::Succeed {
4220 result: value.unwrap_or_default(),
4221 })
4222 }
4223 }
4224
4225 let (queue, store) = open_queue().await;
4226 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4227 let runtime =
4228 WorkflowRuntime::builder(queue, store, JournalRunner, ChannelHook { tx }).build();
4229 let shutdown = spawn_runtime(runtime.clone());
4230
4231 runtime
4232 .submit(RunSpec {
4233 input: Vec::new(),
4234 ..Default::default()
4235 })
4236 .await
4237 .unwrap();
4238 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4239 .await
4240 .unwrap()
4241 .unwrap();
4242 assert_eq!(outcome.status, TerminalStatus::Succeeded);
4243 assert_eq!(outcome.result.as_deref(), Some(b"entry".as_slice()));
4244
4245 let _ = shutdown.send(());
4246 }
4247
4248 #[tokio::test(start_paused = true)]
4249 async fn a_fail_verdict_acks_with_its_effects_and_no_dead_letter() {
4250 let (queue, store) = open_queue().await;
4251 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4252 let runtime = WorkflowRuntime::builder(
4253 queue.clone(),
4254 store,
4255 EffectStagingRunner::new(vec![Ok(StepOutcome::Fail {
4256 reason: "denied".to_string(),
4257 })]),
4258 ChannelHook { tx },
4259 )
4260 .build();
4261 let shutdown = spawn_runtime(runtime.clone());
4262
4263 let handle = runtime
4264 .submit(RunSpec {
4265 input: Vec::new(),
4266 ..Default::default()
4267 })
4268 .await
4269 .unwrap();
4270 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4271 .await
4272 .unwrap()
4273 .unwrap();
4274 assert_eq!(outcome.run_id, handle.run_id);
4275 assert_eq!(outcome.status, TerminalStatus::Failed);
4276 assert_eq!(outcome.error.as_deref(), Some("denied"));
4277 assert_eq!(
4278 terminal_status_of(&runtime, &handle.run_id).await,
4279 Some(outcome.status)
4280 );
4281 assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4282 assert_eq!(
4283 queue.view().stats("workflow-steps").await.unwrap().dead,
4284 0,
4285 "a Fail verdict must not dead-letter"
4286 );
4287
4288 let _ = shutdown.send(());
4289 }
4290
4291 #[tokio::test(start_paused = true)]
4292 async fn a_runner_cancelling_its_own_token_is_not_an_external_cancel() {
4293 struct SelfCancellingRunner;
4297 impl StepRunner for SelfCancellingRunner {
4298 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4299 step.cancel_token.cancel();
4300 step.effects
4301 .put("app/step-0", b"done")
4302 .map_err(|e| StepError::permanent(e.to_string()))?;
4303 Ok(StepOutcome::Succeed {
4304 result: b"finished".to_vec(),
4305 })
4306 }
4307 }
4308
4309 let (queue, store, _clock) = open_queue_at(10_000).await;
4310 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4311 let runtime = WorkflowRuntime::builder(
4312 queue.clone(),
4313 store,
4314 SelfCancellingRunner,
4315 ChannelHook { tx },
4316 )
4317 .build();
4318 let shutdown = spawn_runtime(runtime.clone());
4319
4320 runtime
4321 .submit(RunSpec {
4322 input: Vec::new(),
4323 ..Default::default()
4324 })
4325 .await
4326 .unwrap();
4327 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4328 .await
4329 .unwrap()
4330 .unwrap();
4331 assert_eq!(outcome.status, TerminalStatus::Succeeded);
4332 assert_eq!(outcome.result.as_deref(), Some(b"finished".as_slice()));
4333 assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4334
4335 let _ = shutdown.send(());
4336 }
4337
4338 #[tokio::test(start_paused = true)]
4339 async fn a_cancel_verdict_acks_with_its_effects_and_no_dead_letter() {
4340 let (queue, store) = open_queue().await;
4341 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4342 let runtime = WorkflowRuntime::builder(
4343 queue.clone(),
4344 store,
4345 EffectStagingRunner::new(vec![Ok(StepOutcome::Cancel {
4346 reason: "obsolete".to_string(),
4347 })]),
4348 ChannelHook { tx },
4349 )
4350 .build();
4351 let shutdown = spawn_runtime(runtime.clone());
4352
4353 let handle = runtime
4354 .submit(RunSpec {
4355 input: Vec::new(),
4356 ..Default::default()
4357 })
4358 .await
4359 .unwrap();
4360 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4361 .await
4362 .unwrap()
4363 .unwrap();
4364 assert_eq!(outcome.run_id, handle.run_id);
4365 assert_eq!(outcome.status, TerminalStatus::Cancelled);
4366 assert_eq!(outcome.error.as_deref(), Some("obsolete"));
4367 assert_eq!(
4368 terminal_status_of(&runtime, &handle.run_id).await,
4369 Some(outcome.status)
4370 );
4371 assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4372 assert_eq!(
4373 queue.view().stats("workflow-steps").await.unwrap().dead,
4374 0,
4375 "a Cancel verdict must not dead-letter"
4376 );
4377
4378 let _ = shutdown.send(());
4379 }
4380
4381 #[tokio::test(start_paused = true)]
4382 async fn an_external_cancel_discards_staged_effects() {
4383 struct StageThenAwaitCancel {
4384 started: tokio::sync::mpsc::UnboundedSender<()>,
4385 }
4386
4387 impl StepRunner for StageThenAwaitCancel {
4388 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4389 step.effects
4390 .put(b"app/override".to_vec(), b"staged".to_vec())
4391 .map_err(|e| StepError::permanent(e.to_string()))?;
4392 let _ = self.started.send(());
4393 step.cancel_token.cancelled().await;
4394 Ok(StepOutcome::continue_now(Vec::new()))
4395 }
4396 }
4397
4398 let (queue, store) = open_queue().await;
4399 let (started_tx, mut started_rx) = tokio::sync::mpsc::unbounded_channel();
4400 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4401 let runtime = WorkflowRuntime::builder(
4402 queue.clone(),
4403 store,
4404 StageThenAwaitCancel {
4405 started: started_tx,
4406 },
4407 ChannelHook { tx },
4408 )
4409 .build();
4410 let shutdown = spawn_runtime(runtime.clone());
4411
4412 let handle = runtime
4413 .submit(RunSpec {
4414 input: Vec::new(),
4415 ..Default::default()
4416 })
4417 .await
4418 .unwrap();
4419 tokio::time::timeout(Duration::from_secs(2), started_rx.recv())
4420 .await
4421 .unwrap()
4422 .unwrap();
4423 assert!(runtime.cancel(&handle.run_id).await.unwrap());
4424
4425 let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4426 .await
4427 .unwrap()
4428 .unwrap();
4429 assert_eq!(outcome.status, TerminalStatus::Cancelled);
4430 assert_eq!(outcome.error, None);
4431
4432 wait_for_drained(&queue).await;
4433 assert!(
4434 queue
4435 .view()
4436 .kv_get(b"app/override")
4437 .await
4438 .unwrap()
4439 .is_none()
4440 );
4441
4442 let _ = shutdown.send(());
4443 }
4444
4445 #[tokio::test(start_paused = true)]
4446 async fn a_run_dead_lettered_by_the_reaper_is_terminated_by_reconciliation() {
4447 let (queue, store, clock) = open_queue_at_with(
4448 1_700_000_000_000,
4449 fast_options().default_queue_config(
4450 QueueConfig::default()
4451 .retry_backoff_base(Duration::ZERO)
4452 .lease_duration(Duration::from_secs(1)),
4453 ),
4454 )
4455 .await;
4456 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4457 let runtime =
4458 WorkflowRuntime::builder(queue.clone(), store, PauseRunner, ChannelHook { tx })
4459 .poll_interval(Duration::from_millis(10))
4460 .build();
4461 let shutdown = spawn_runtime(runtime.clone());
4462
4463 let submitted = runtime
4464 .submit(RunSpec {
4465 run_id: Some(rid("hung")),
4466 input: Vec::new(),
4467 options: RunOptions {
4468 max_attempts_per_step: Some(1),
4469 ..Default::default()
4470 },
4471 ..Default::default()
4472 })
4473 .await
4474 .unwrap();
4475 for _ in 0..200 {
4476 if queue.view().stats("workflow-steps").await.unwrap().claimed == 1 {
4477 break;
4478 }
4479 tokio::time::sleep(Duration::from_millis(10)).await;
4480 }
4481 assert_eq!(
4482 runtime.status(&rid("hung")).await.unwrap().map(|s| s.state),
4483 Some(RunState::Running)
4484 );
4485
4486 advance(&clock, Duration::from_secs(2)).await;
4489 let outcome = tokio::time::timeout(Duration::from_secs(5), rx.recv())
4490 .await
4491 .unwrap()
4492 .unwrap();
4493 assert_eq!(outcome.run_id, "hung");
4494 assert_eq!(outcome.status, TerminalStatus::Failed);
4495 assert_eq!(outcome.final_step, 0);
4496 assert_eq!(
4497 queue
4498 .view()
4499 .get_job(&submitted.job_id)
4500 .await
4501 .unwrap()
4502 .unwrap()
4503 .status,
4504 JobStatus::Dead
4505 );
4506 assert!(
4507 queue
4508 .view()
4509 .kv_get(&run_kv_key(&rid("hung")))
4510 .await
4511 .unwrap()
4512 .is_none()
4513 );
4514 assert!(
4515 queue
4516 .view()
4517 .kv_get(&step_kv_key(&rid("hung")))
4518 .await
4519 .unwrap()
4520 .is_none()
4521 );
4522 assert_eq!(
4523 terminal_status_of(&runtime, &rid("hung")).await,
4524 Some(TerminalStatus::Failed)
4525 );
4526 let _ = shutdown.send(());
4527
4528 let again = runtime
4532 .submit(RunSpec {
4533 run_id: Some(rid("hung")),
4534 input: Vec::new(),
4535 ..Default::default()
4536 })
4537 .await
4538 .unwrap();
4539 assert!(again.newly_submitted);
4540 assert_eq!(runtime.inner.core.reconcile_dead_steps().await.unwrap(), 0);
4541 assert!(
4542 runtime.status(&rid("hung")).await.unwrap().is_some(),
4543 "the re-submitted run is active"
4544 );
4545 }
4546
4547 #[tokio::test(start_paused = true)]
4548 async fn a_wait_follows_the_run_across_its_steps() {
4549 let (queue, store) = open_queue().await;
4550 let runtime = WorkflowRuntime::builder(
4551 queue,
4552 store,
4553 ScriptedRunner::new(vec![
4554 StepOutcome::continue_now(b"next".to_vec()),
4555 StepOutcome::Succeed {
4556 result: b"done".to_vec(),
4557 },
4558 ]),
4559 NoopTerminalHook,
4560 )
4561 .build();
4562 assert!(matches!(
4563 runtime.wait(&rid("absent")).await,
4564 Err(Error::RunNotFound(id)) if id == "absent"
4565 ));
4566 let shutdown = spawn_runtime(runtime.clone());
4567
4568 runtime
4569 .submit(RunSpec {
4570 run_id: Some(rid("two")),
4571 input: b"x".to_vec(),
4572 ..Default::default()
4573 })
4574 .await
4575 .unwrap();
4576 let end = tokio::time::timeout(Duration::from_secs(5), runtime.wait(&rid("two")))
4577 .await
4578 .expect("the wait resolved")
4579 .unwrap();
4580 let outcome = end.outcome.expect("the worker recorded the outcome");
4581 assert_eq!(
4582 (outcome.final_step, outcome.result.as_deref()),
4583 (1, Some(b"done".as_slice()))
4584 );
4585 assert_eq!(
4586 (end.termination.status, end.termination.final_step),
4587 (TerminalStatus::Succeeded, 1)
4588 );
4589
4590 let again = runtime.wait(&rid("two")).await.unwrap();
4592 assert_eq!(again.termination, end.termination);
4593 let _ = shutdown.send(());
4594 }
4595
4596 #[tokio::test(start_paused = true)]
4597 async fn a_result_record_of_an_earlier_run_is_not_reported_for_a_re_submitted_run_id() {
4598 let (queue, store, clock) = open_queue_at(10_000).await;
4599 let runtime = WorkflowRuntime::builder(
4600 queue.clone(),
4601 store,
4602 FixedRunner::new(Ok(StepOutcome::Succeed {
4603 result: b"done".to_vec(),
4604 })),
4605 NoopTerminalHook,
4606 )
4607 .build();
4608 let spec = RunSpec {
4609 run_id: Some(rid("again")),
4610 input: b"x".to_vec(),
4611 ..Default::default()
4612 };
4613 runtime.submit(spec.clone()).await.unwrap();
4614 let claim = queue
4615 .claim("workflow-steps", Duration::from_secs(30))
4616 .await
4617 .unwrap()
4618 .unwrap();
4619 let effects = runtime
4620 .inner
4621 .process_step(&claim, &LeaseHandle::detached())
4622 .await
4623 .unwrap();
4624 queue.ack_with(&claim, effects).await.unwrap();
4625 let first = runtime.wait(&rid("again")).await.unwrap();
4626 assert_eq!(first.termination.status, TerminalStatus::Succeeded);
4627 assert!(first.outcome.is_some());
4628
4629 clock.advance(Duration::from_secs(1));
4632 assert!(runtime.submit(spec).await.unwrap().newly_submitted);
4633 assert!(
4634 runtime.outcome(&rid("again")).await.unwrap().is_none(),
4635 "the run is active"
4636 );
4637 assert!(runtime.cancel(&rid("again")).await.unwrap());
4638 let end = runtime.wait(&rid("again")).await.unwrap();
4639 assert_eq!(
4640 (end.termination.status, end.termination.terminated_at_ms),
4641 (TerminalStatus::Cancelled, 11_000)
4642 );
4643 assert!(
4644 end.outcome.is_none(),
4645 "no worker terminated the new run, so the earlier run's record is not its outcome"
4646 );
4647 assert!(runtime.outcome(&rid("again")).await.unwrap().is_none());
4648 }
4649
4650 #[tokio::test(start_paused = true)]
4651 async fn a_cancel_request_does_not_reach_a_re_submission_of_the_run_id() {
4652 let (queue, store, _clock) = open_queue_at(10_000).await;
4653 let runtime =
4654 WorkflowRuntime::builder(queue.clone(), store, UnreachableRunner, NoopTerminalHook)
4655 .build();
4656 let spec = RunSpec {
4657 run_id: Some(rid("again")),
4658 input: b"x".to_vec(),
4659 ..Default::default()
4660 };
4661 runtime.submit(spec.clone()).await.unwrap();
4662 assert!(runtime.cancel(&rid("again")).await.unwrap());
4663 let end = runtime.wait(&rid("again")).await.unwrap();
4664 assert_eq!(end.termination.status, TerminalStatus::Cancelled);
4665
4666 assert!(runtime.submit(spec).await.unwrap().newly_submitted);
4667 let record = runtime
4668 .view()
4669 .run_record(&rid("again"))
4670 .await
4671 .unwrap()
4672 .unwrap();
4673 assert!(!record.cancel_requested);
4674 let status = runtime.status(&rid("again")).await.unwrap().unwrap();
4675 assert_eq!(status.state, RunState::Pending);
4676 }
4677
4678 #[tokio::test(start_paused = true)]
4679 async fn a_step_dead_lettered_outside_the_worker_is_waited_for_until_reconciliation() {
4680 let (queue, store, _clock) = open_queue_at(10_000).await;
4681 let runtime =
4682 WorkflowRuntime::builder(queue.clone(), store, UnreachableRunner, NoopTerminalHook)
4683 .poll_interval(Duration::from_millis(10))
4684 .build();
4685 runtime
4686 .submit(RunSpec {
4687 run_id: Some(rid("hung")),
4688 input: Vec::new(),
4689 ..Default::default()
4690 })
4691 .await
4692 .unwrap();
4693 let claim = queue
4694 .claim("workflow-steps", Duration::from_secs(60))
4695 .await
4696 .unwrap()
4697 .unwrap();
4698 queue.dead_letter(&claim, "hung").await.unwrap();
4699
4700 let waiting = tokio::spawn({
4701 let runtime = runtime.clone();
4702 async move { runtime.wait(&rid("hung")).await }
4703 });
4704 assert!(
4705 !runtime.cancel(&rid("hung")).await.unwrap(),
4706 "the request is not honoured"
4707 );
4708 assert!(
4709 runtime
4710 .wait_timeout(&rid("hung"), Duration::from_secs(1))
4711 .await
4712 .unwrap()
4713 .is_none(),
4714 "nothing terminates the run without a worker"
4715 );
4716 assert_eq!(runtime.inner.core.reconcile_dead_steps().await.unwrap(), 1);
4717 let end = tokio::time::timeout(Duration::from_secs(5), waiting)
4718 .await
4719 .expect("the wait resolved")
4720 .unwrap()
4721 .unwrap();
4722 assert_eq!(
4723 end.termination,
4724 RunTermination {
4725 status: TerminalStatus::Failed,
4726 error: Some("hung".into()),
4727 error_kind: None,
4728 final_step: 0,
4729 terminated_at_ms: 10_000,
4730 },
4731 "the termination is read from the terminal record",
4732 );
4733 assert!(end.outcome.is_none());
4734 assert_eq!(
4735 runtime.wait(&rid("hung")).await.unwrap().termination,
4736 end.termination,
4737 "the record is retained without memo retention",
4738 );
4739 }
4740
4741 #[tokio::test(start_paused = true)]
4742 async fn a_pointer_over_a_missing_job_is_an_inconsistent_run_state() {
4743 let (queue, store) = open_queue().await;
4744 let runtime =
4745 WorkflowRuntime::builder(queue.clone(), store, UnreachableRunner, NoopTerminalHook)
4746 .build();
4747 let submitted = runtime
4748 .submit(RunSpec {
4749 run_id: Some(rid("torn")),
4750 input: Vec::new(),
4751 ..Default::default()
4752 })
4753 .await
4754 .unwrap();
4755 queue.cancel(&submitted.job_id).await.unwrap();
4758 assert!(matches!(
4759 runtime.status(&rid("torn")).await,
4760 Err(Error::InconsistentRunState(id)) if id == "torn"
4761 ));
4762 assert!(matches!(
4763 runtime.wait(&rid("torn")).await,
4764 Err(Error::InconsistentRunState(id)) if id == "torn"
4765 ));
4766 assert!(matches!(
4767 runtime.cancel(&rid("torn")).await,
4768 Err(Error::InconsistentRunState(id)) if id == "torn"
4769 ));
4770 }
4771
4772 #[tokio::test(start_paused = true)]
4773 async fn a_member_record_is_rewritten_only_by_the_terminating_settlement() {
4774 struct RecordReadingRunner {
4775 pending_seen: Arc<StdMutex<Vec<bool>>>,
4776 }
4777
4778 impl StepRunner for RecordReadingRunner {
4779 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4780 let record = step
4781 .kv
4782 .get(&group_member_kv_key(&rid("g"), "m"))
4783 .await?
4784 .expect("the member record is written with the submission");
4785 let member: DurableMember = rmp_serde::from_slice(&record).unwrap();
4786 self.pending_seen
4787 .lock()
4788 .unwrap()
4789 .push(member.terminated.is_none());
4790 Err(StepError::transient("still failing"))
4791 }
4792 }
4793
4794 let (queue, store) = open_queue_with(fast_options()).await;
4795 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4796 let pending_seen = Arc::new(StdMutex::new(Vec::new()));
4797 let runtime = WorkflowRuntime::builder(
4798 queue.clone(),
4799 store,
4800 RecordReadingRunner {
4801 pending_seen: pending_seen.clone(),
4802 },
4803 ChannelHook { tx },
4804 )
4805 .build();
4806 let shutdown = spawn_runtime(runtime.clone());
4807
4808 let group = runtime.group(rid("g"));
4809 group
4810 .submit(
4811 vec![GroupMember {
4812 key: "m".to_string(),
4813 input: Vec::new(),
4814 }],
4815 &RunOptions {
4816 max_attempts_per_step: Some(2),
4817 ..RunOptions::default()
4818 },
4819 )
4820 .await
4821 .unwrap();
4822 let outcome = tokio::time::timeout(Duration::from_secs(5), rx.recv())
4823 .await
4824 .unwrap()
4825 .unwrap();
4826 assert_eq!(outcome.status, TerminalStatus::Failed);
4827 for _ in 0..200 {
4828 if queue.view().stats("workflow-steps").await.unwrap().dead == 1 {
4829 break;
4830 }
4831 tokio::time::sleep(Duration::from_millis(10)).await;
4832 }
4833 assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 1);
4834
4835 assert_eq!(*pending_seen.lock().unwrap(), vec![true, true]);
4839 let members = group.members().await.unwrap();
4840 assert_eq!(members.len(), 1);
4841 assert_eq!(members[0].key, "m");
4842 assert_eq!(members[0].status(), Some(TerminalStatus::Failed));
4843 assert_eq!(
4844 members[0]
4845 .record
4846 .terminated
4847 .as_ref()
4848 .unwrap()
4849 .error
4850 .as_deref(),
4851 Some("still failing")
4852 );
4853
4854 let _ = shutdown.send(());
4855 }
4856
4857 #[tokio::test(start_paused = true)]
4858 async fn a_replayed_step_outcome_restores_its_staged_effects() {
4859 struct StagingContinueRunner {
4860 calls: Arc<AtomicU32>,
4861 }
4862
4863 impl StepRunner for StagingContinueRunner {
4864 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4865 self.calls.fetch_add(1, Ordering::SeqCst);
4866 step.effects
4867 .put(b"app/replayed".to_vec(), b"v".to_vec())
4868 .map_err(|e| StepError::transient(e.to_string()))?;
4869 step.effects
4870 .delete(b"app/stale".to_vec())
4871 .map_err(|e| StepError::transient(e.to_string()))?;
4872 Ok(StepOutcome::continue_now(b"step1".to_vec()))
4873 }
4874 }
4875
4876 let (queue, store) = open_queue().await;
4877 let calls = Arc::new(AtomicU32::new(0));
4878 let runtime = WorkflowRuntime::builder(
4879 queue.clone(),
4880 store,
4881 StagingContinueRunner {
4882 calls: calls.clone(),
4883 },
4884 NoopTerminalHook,
4885 )
4886 .step_output_replay()
4887 .build();
4888
4889 queue.kv_put(b"app/stale", b"old").await.unwrap();
4890 runtime
4891 .submit(RunSpec {
4892 run_id: Some(rid("replay-effects")),
4893 input: b"input".to_vec(),
4894 ..Default::default()
4895 })
4896 .await
4897 .unwrap();
4898
4899 let job = queue
4900 .claim("workflow-steps", Duration::from_secs(30))
4901 .await
4902 .unwrap()
4903 .unwrap();
4904
4905 let _ = runtime
4908 .inner
4909 .process_step(&job, &LeaseHandle::detached())
4910 .await
4911 .unwrap();
4912 assert_eq!(calls.load(Ordering::SeqCst), 1);
4913 assert!(
4914 queue
4915 .view()
4916 .kv_get(b"app/replayed")
4917 .await
4918 .unwrap()
4919 .is_none()
4920 );
4921
4922 let effects = runtime
4925 .inner
4926 .process_step(&job, &LeaseHandle::detached())
4927 .await
4928 .unwrap();
4929 assert_eq!(calls.load(Ordering::SeqCst), 1);
4930 assert_eq!(
4931 effects.kv_writes.get(b"app/replayed".as_slice()),
4932 Some(&b"v".to_vec())
4933 );
4934 assert!(effects.kv_deletes.contains(&b"app/stale".to_vec()));
4935 queue.ack_with(&job, effects).await.unwrap();
4936 assert_eq!(
4937 queue
4938 .view()
4939 .kv_get(b"app/replayed")
4940 .await
4941 .unwrap()
4942 .as_deref(),
4943 Some(b"v".as_slice())
4944 );
4945 assert!(queue.view().kv_get(b"app/stale").await.unwrap().is_none());
4946 }
4947
4948 #[tokio::test(start_paused = true)]
4949 async fn only_the_committed_outcome_produces_a_notification() {
4950 struct GatedSecondAttempt {
4951 calls: Arc<AtomicU32>,
4952 running: tokio::sync::mpsc::UnboundedSender<()>,
4953 }
4954
4955 impl StepRunner for GatedSecondAttempt {
4956 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4957 if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
4958 return Ok(StepOutcome::Succeed {
4959 result: b"done".to_vec(),
4960 });
4961 }
4962 let _ = self.running.send(());
4963 step.cancel_token.cancelled().await;
4964 Ok(StepOutcome::Succeed {
4965 result: b"done".to_vec(),
4966 })
4967 }
4968 }
4969
4970 let (queue, store) = open_queue().await;
4971 let (running_tx, mut running_rx) = tokio::sync::mpsc::unbounded_channel();
4972 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4973 let runtime = WorkflowRuntime::builder(
4974 queue.clone(),
4975 store,
4976 GatedSecondAttempt {
4977 calls: Arc::new(AtomicU32::new(0)),
4978 running: running_tx,
4979 },
4980 ChannelHook { tx },
4981 )
4982 .build();
4983
4984 runtime
4985 .submit(RunSpec {
4986 run_id: Some(rid("phantom")),
4987 input: Vec::new(),
4988 ..Default::default()
4989 })
4990 .await
4991 .unwrap();
4992 let job = queue
4993 .claim("workflow-steps", Duration::from_secs(30))
4994 .await
4995 .unwrap()
4996 .unwrap();
4997
4998 let _ = runtime
5002 .inner
5003 .process_step(&job, &queue.lease_handle(&job))
5004 .await
5005 .unwrap();
5006
5007 let worker = {
5010 let inner = runtime.inner.clone();
5011 let queue = queue.clone();
5012 tokio::spawn(async move {
5013 let effects = inner
5014 .process_step(&job, &queue.lease_handle(&job))
5015 .await
5016 .unwrap();
5017 queue.ack_with(&job, effects).await.unwrap();
5018 })
5019 };
5020 running_rx.recv().await.unwrap();
5021 assert!(runtime.cancel(&rid("phantom")).await.unwrap());
5022 worker.await.unwrap();
5023
5024 let notification = queue
5025 .claim("workflow-steps", Duration::from_secs(30))
5026 .await
5027 .unwrap()
5028 .expect("the committed settlement enqueued its notification");
5029 let effects = runtime
5030 .inner
5031 .process_step(¬ification, &LeaseHandle::detached())
5032 .await
5033 .unwrap();
5034 queue.ack_with(¬ification, effects).await.unwrap();
5035 let outcome = rx.recv().await.unwrap();
5036 assert_eq!(outcome.status, TerminalStatus::Cancelled);
5037 assert!(
5038 queue
5039 .claim("workflow-steps", Duration::from_secs(30))
5040 .await
5041 .unwrap()
5042 .is_none(),
5043 "the outcome that never committed must produce no notification",
5044 );
5045 assert!(rx.try_recv().is_err());
5046 }
5047
5048 #[tokio::test(start_paused = true)]
5049 async fn hook_effects_commit_with_the_notification_ack() {
5050 struct EffectHook;
5051
5052 impl TerminalHook for EffectHook {
5053 async fn on_termination(
5054 &self,
5055 outcome: &RunOutcome,
5056 effects: &TerminalEffects,
5057 ) -> std::result::Result<(), StepError> {
5058 effects
5059 .put(
5060 format!("app/outcomes/{}", outcome.run_id),
5061 outcome.status.as_str(),
5062 )
5063 .map_err(|e| StepError::permanent(e.to_string()))?;
5064 effects
5065 .enqueue(EnqueueRequest {
5066 queue: "side-effects".to_string(),
5067 payload: outcome.run_id.to_string().into_bytes(),
5068 options: EnqueueOptions::default(),
5069 })
5070 .map_err(|e| StepError::permanent(e.to_string()))?;
5071 Ok(())
5072 }
5073 }
5074
5075 let (queue, store) = open_queue().await;
5076 let runtime = WorkflowRuntime::builder(
5077 queue.clone(),
5078 store,
5079 ScriptedRunner::new(vec![StepOutcome::Succeed { result: Vec::new() }]),
5080 EffectHook,
5081 )
5082 .build();
5083 let shutdown = spawn_runtime(runtime.clone());
5084
5085 runtime
5086 .submit(RunSpec {
5087 run_id: Some(rid("hooked")),
5088 input: Vec::new(),
5089 ..Default::default()
5090 })
5091 .await
5092 .unwrap();
5093
5094 assert_eq!(
5095 wait_for_kv(&queue, b"app/outcomes/hooked").await,
5096 b"succeeded"
5097 );
5098 let side = queue
5099 .claim("side-effects", Duration::from_secs(30))
5100 .await
5101 .unwrap()
5102 .expect("the staged enqueue committed with the notification ack");
5103 assert_eq!(side.payload.as_slice(), b"hooked");
5104
5105 let _ = shutdown.send(());
5106 }
5107
5108 #[tokio::test(start_paused = true)]
5109 async fn a_transiently_failing_hook_retries_the_notification() {
5110 struct FlakyHook {
5111 calls: Arc<AtomicU32>,
5112 }
5113
5114 impl TerminalHook for FlakyHook {
5115 async fn on_termination(
5116 &self,
5117 outcome: &RunOutcome,
5118 effects: &TerminalEffects,
5119 ) -> std::result::Result<(), StepError> {
5120 if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
5121 return Err(StepError::transient("first attempt fails"));
5122 }
5123 effects
5124 .put(format!("app/notified/{}", outcome.run_id), b"1".to_vec())
5125 .map_err(|e| StepError::permanent(e.to_string()))?;
5126 Ok(())
5127 }
5128 }
5129
5130 let (queue, store) = open_queue_with(fast_options()).await;
5131 let calls = Arc::new(AtomicU32::new(0));
5132 let runtime = WorkflowRuntime::builder(
5133 queue.clone(),
5134 store,
5135 ScriptedRunner::new(vec![StepOutcome::Succeed { result: Vec::new() }]),
5136 FlakyHook {
5137 calls: calls.clone(),
5138 },
5139 )
5140 .build();
5141 let shutdown = spawn_runtime(runtime.clone());
5142
5143 runtime
5144 .submit(RunSpec {
5145 run_id: Some(rid("flaky")),
5146 input: Vec::new(),
5147 ..Default::default()
5148 })
5149 .await
5150 .unwrap();
5151
5152 assert_eq!(wait_for_kv(&queue, b"app/notified/flaky").await, b"1");
5153 assert_eq!(calls.load(Ordering::SeqCst), 2);
5154
5155 let _ = shutdown.send(());
5156 }
5157
5158 #[tokio::test(start_paused = true)]
5159 async fn a_noop_hook_enqueues_no_notification() {
5160 let (queue, store) = open_queue().await;
5161 let runtime = WorkflowRuntime::builder(
5162 queue.clone(),
5163 store,
5164 ScriptedRunner::new(vec![StepOutcome::Succeed { result: Vec::new() }]),
5165 NoopTerminalHook,
5166 )
5167 .build();
5168 runtime
5169 .submit(RunSpec {
5170 input: Vec::new(),
5171 ..Default::default()
5172 })
5173 .await
5174 .unwrap();
5175 let job = queue
5176 .claim("workflow-steps", Duration::from_secs(30))
5177 .await
5178 .unwrap()
5179 .unwrap();
5180 let effects = runtime
5181 .inner
5182 .process_step(&job, &LeaseHandle::detached())
5183 .await
5184 .unwrap();
5185 assert!(effects.enqueues.is_empty());
5186 queue.ack_with(&job, effects).await.unwrap();
5187 assert!(
5188 queue
5189 .claim("workflow-steps", Duration::from_secs(30))
5190 .await
5191 .unwrap()
5192 .is_none()
5193 );
5194 }
5195
5196 #[cfg(feature = "webhooks")]
5197 #[tokio::test(start_paused = true)]
5198 async fn the_webhook_hook_stages_its_delivery_as_a_notification_effect() {
5199 use crate::terminal::WebhookTerminalHook;
5200
5201 let (queue, store) = open_queue().await;
5202 let runtime = WorkflowRuntime::builder(
5203 queue.clone(),
5204 store,
5205 ScriptedRunner::new(vec![
5206 StepOutcome::Succeed {
5207 result: b"payload".to_vec(),
5208 },
5209 StepOutcome::Succeed { result: Vec::new() },
5210 ]),
5211 WebhookTerminalHook::new("callbacks"),
5212 )
5213 .build();
5214 let shutdown = spawn_runtime(runtime.clone());
5215
5216 runtime
5217 .submit(RunSpec {
5218 run_id: Some(rid("with-callback")),
5219 input: Vec::new(),
5220 options: RunOptions {
5221 headers: HashMap::from([(
5222 "callback_url".to_string(),
5223 "https://example.com/done".to_string(),
5224 )]),
5225 ..Default::default()
5226 },
5227 ..Default::default()
5228 })
5229 .await
5230 .unwrap();
5231
5232 let webhook = loop {
5233 if let Some(job) = queue
5234 .claim("callbacks", Duration::from_secs(30))
5235 .await
5236 .unwrap()
5237 {
5238 break job;
5239 }
5240 tokio::time::sleep(Duration::from_millis(10)).await;
5241 };
5242 assert_eq!(webhook.payload.as_slice(), b"payload");
5243 assert_eq!(
5244 webhook.headers.get("webhook.url").unwrap(),
5245 "https://example.com/done"
5246 );
5247 assert_eq!(
5248 webhook.headers.get("http.Workflow-Run-Status").unwrap(),
5249 "succeeded"
5250 );
5251
5252 runtime
5254 .submit(RunSpec {
5255 run_id: Some(rid("without-callback")),
5256 input: Vec::new(),
5257 ..Default::default()
5258 })
5259 .await
5260 .unwrap();
5261 wait_for_drained(&queue).await;
5262 assert!(
5263 queue
5264 .claim("callbacks", Duration::from_secs(30))
5265 .await
5266 .unwrap()
5267 .is_none()
5268 );
5269
5270 let _ = shutdown.send(());
5271 }
5272
5273 #[tokio::test(start_paused = true)]
5274 async fn a_late_write_through_an_escaped_handle_is_refused() {
5275 struct EscapingRunner {
5276 escaped: Arc<StdMutex<Option<EffectsHandle>>>,
5277 }
5278
5279 impl StepRunner for EscapingRunner {
5280 async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
5281 *self.escaped.lock().unwrap() = Some(step.effects.clone());
5282 Ok(StepOutcome::Succeed { result: Vec::new() })
5283 }
5284 }
5285
5286 let (queue, store) = open_queue().await;
5287 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
5288 let escaped = Arc::new(StdMutex::new(None));
5289 let runtime = WorkflowRuntime::builder(
5290 queue,
5291 store,
5292 EscapingRunner {
5293 escaped: escaped.clone(),
5294 },
5295 ChannelHook { tx },
5296 )
5297 .build();
5298 let shutdown = spawn_runtime(runtime.clone());
5299
5300 runtime
5301 .submit(RunSpec {
5302 input: Vec::new(),
5303 ..Default::default()
5304 })
5305 .await
5306 .unwrap();
5307 tokio::time::timeout(Duration::from_secs(2), rx.recv())
5308 .await
5309 .unwrap()
5310 .unwrap();
5311
5312 let handle = escaped.lock().unwrap().take().unwrap();
5313 assert!(matches!(
5314 handle.put("app/late", "v"),
5315 Err(Error::EffectsSealed)
5316 ));
5317
5318 let _ = shutdown.send(());
5319 }
5320}