Skip to main content

taquba_workflow/
group.rs

1//! Run groups: a durable set of runs of one runtime, identified by a
2//! group id, whose membership is a manifest in the object store and
3//! whose per-member state is a record in the queue's KV namespace.
4//! [`RunGroup`] submits the members, yields their results, cancels them
5//! and removes the group's state; [`jobs::JobGroup`](crate::jobs::JobGroup)
6//! is its typed presentation.
7
8use std::collections::HashMap;
9use std::sync::Arc;
10
11use futures_util::stream::{self, FuturesUnordered, Stream, StreamExt, TryStreamExt};
12use serde::{Deserialize, Serialize};
13use taquba::object_store::{ObjectStore, path::Path};
14use taquba::{Queue, SettlementEffects};
15use tracing::warn;
16
17use crate::blob::ObjectPrefix;
18use crate::durable::{self, DurableMember, DurableTermination};
19use crate::error::{Error, Result};
20use crate::keys::{
21    HEADER_GROUP, HEADER_GROUP_KEY, RunId, group_member_kv_key, group_members_kv_prefix,
22    outcome_kv_key,
23};
24use crate::memo::MemoStore;
25use crate::runtime::{RunOptions, RunSpec, RunTermination, RuntimeCore};
26use crate::sweep::Clearable;
27use crate::terminal::{RunOutcome, TerminalStatus};
28
29/// Member submissions and cancellations in flight at once. Each blocks
30/// on a durable commit, and concurrent commits share WAL flushes.
31const SUBMIT_CONCURRENCY: usize = 32;
32/// Member records read per page, and deleted per transaction by
33/// [`GroupStore::forget`].
34const MEMBER_PAGE_SIZE: usize = 1000;
35
36/// The run id of the member `key` of group `group_id`: the hex SHA-256
37/// digest of `{group_id}/{key}`, so a key can contain characters a run
38/// id rejects and groups never share run state.
39pub(crate) fn member_run_id(group_id: &RunId, key: &str) -> RunId {
40    RunId::digest(&[group_id.as_bytes(), b"/", key.as_bytes()])
41}
42
43/// The group membership of a run, set on every step job of the run in
44/// the [`HEADER_GROUP`] and [`HEADER_GROUP_KEY`] headers.
45#[derive(Debug, Clone, PartialEq, Eq)]
46pub(crate) struct Membership {
47    pub(crate) group_id: RunId,
48    pub(crate) key: String,
49}
50
51impl Membership {
52    /// The membership named by `headers`, when both headers are present.
53    pub(crate) fn from_headers(headers: &HashMap<String, String>) -> Result<Option<Self>> {
54        let (Some(group_id), Some(key)) =
55            (headers.get(HEADER_GROUP), headers.get(HEADER_GROUP_KEY))
56        else {
57            return Ok(None);
58        };
59        Ok(Some(Self {
60            group_id: RunId::new(group_id.as_str())?,
61            key: key.clone(),
62        }))
63    }
64
65    /// The reserved headers that hold this membership on a step job.
66    pub(crate) fn reserved_headers(&self) -> Vec<(&'static str, String)> {
67        vec![
68            (HEADER_GROUP, self.group_id.to_string()),
69            (HEADER_GROUP_KEY, self.key.clone()),
70        ]
71    }
72
73    /// The KV key of this member's record.
74    pub(crate) fn kv_key(&self) -> Vec<u8> {
75        group_member_kv_key(&self.group_id, &self.key)
76    }
77
78    /// The run id of this member.
79    pub(crate) fn run_id(&self) -> RunId {
80        member_run_id(&self.group_id, &self.key)
81    }
82}
83
84/// One member of a [`RunGroup`]: its key, unique within the group, and
85/// the input of its run.
86#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
87pub struct GroupMember {
88    /// The member's key.
89    pub key: String,
90    /// The input of the member's step 0.
91    pub input: Vec<u8>,
92}
93
94/// The members of one group, in submission order.
95#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
96pub(crate) struct Manifest {
97    pub(crate) group_id: RunId,
98    pub(crate) members: Vec<GroupMember>,
99}
100
101/// The durable state of a group, read from its manifest and member
102/// records by [`RunGroup::status`].
103#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct GroupStatus {
105    /// The group id.
106    pub group_id: RunId,
107    /// Number of members in the group's manifest.
108    pub total: usize,
109    /// Members submitted and not yet terminated.
110    pub pending: usize,
111    /// Members whose last recorded termination is a success.
112    pub succeeded: usize,
113    /// Members whose last recorded termination is a failure.
114    pub failed: usize,
115    /// Members whose last recorded termination is a cancellation.
116    pub cancelled: usize,
117}
118
119/// A terminated member of a group, yielded by [`RunGroup::results`].
120#[derive(Debug, Clone)]
121pub struct MemberResult {
122    /// The member's key.
123    pub key: String,
124    /// The run id of the member's run.
125    pub run_id: RunId,
126    /// The member's last recorded termination.
127    pub termination: RunTermination,
128    /// The member's committed outcome, read as [`WorkflowRuntime::outcome`](crate::WorkflowRuntime::outcome)
129    /// reads it; `None` for a member terminated without a worker.
130    pub outcome: Option<RunOutcome>,
131}
132
133/// A member record as read back, with the key it is stored under.
134pub(crate) struct MemberState {
135    pub(crate) key: String,
136    pub(crate) record: DurableMember,
137}
138
139impl MemberState {
140    /// The member's terminal status, `None` while it is active.
141    pub(crate) fn status(&self) -> Option<TerminalStatus> {
142        self.record
143            .terminated
144            .as_ref()
145            .map(|termination| termination.status.into())
146    }
147}
148
149/// The durable state of groups: manifests under
150/// `<memo_prefix>/groups/<group_id>/manifest` in the object store and
151/// member records under `workflow/groups/<group_id>/` in the queue's KV
152/// namespace.
153#[derive(Clone)]
154pub(crate) struct GroupStore {
155    objects: ObjectPrefix,
156    memo_store: MemoStore,
157    queue: Arc<Queue>,
158}
159
160impl GroupStore {
161    pub(crate) fn new(
162        store: Arc<dyn ObjectStore>,
163        prefix: impl Into<String>,
164        memo_store: MemoStore,
165        queue: Arc<Queue>,
166    ) -> Self {
167        Self {
168            objects: ObjectPrefix::new(store, prefix),
169            memo_store,
170            queue,
171        }
172    }
173
174    fn manifest_path(&self, group_id: &RunId) -> Path {
175        self.objects.path(&format!("groups/{group_id}/manifest"))
176    }
177
178    pub(crate) async fn read_manifest(&self, group_id: &RunId) -> Result<Option<Manifest>> {
179        match self.objects.get(&self.manifest_path(group_id)).await? {
180            Some(bytes) => durable::decode(&bytes).map(Some),
181            None => Ok(None),
182        }
183    }
184
185    async fn write_manifest(&self, manifest: &Manifest) -> Result<()> {
186        self.objects
187            .put(
188                &self.manifest_path(&manifest.group_id),
189                &durable::encode(manifest),
190            )
191            .await
192    }
193
194    /// The member record of `key` in `group_id`, when one exists.
195    pub(crate) async fn member(
196        &self,
197        group_id: &RunId,
198        key: &str,
199    ) -> Result<Option<DurableMember>> {
200        durable::kv_record(self.queue.view(), &group_member_kv_key(group_id, key)).await
201    }
202
203    /// Every member record of `group_id`, in key order. A record that
204    /// fails to decode is skipped.
205    pub(crate) async fn members(&self, group_id: &RunId) -> Result<Vec<MemberState>> {
206        let prefix = group_members_kv_prefix(group_id);
207        let mut members = Vec::new();
208        let mut entries =
209            std::pin::pin!(self.queue.view().kv_entries(&prefix, .., MEMBER_PAGE_SIZE));
210        while let Some((kv_key, value)) = entries.try_next().await? {
211            let key = String::from_utf8_lossy(&kv_key[prefix.len()..]).into_owned();
212            if let Some(record) = durable::decode_or_absent(
213                &value,
214                "group member record",
215                &format_args!("{group_id}/{key}"),
216            ) {
217                members.push(MemberState { key, record });
218            }
219        }
220        Ok(members)
221    }
222
223    /// Remove the state of `group_id`: the memo entries and the terminal
224    /// record of every member in its manifest, its member records and
225    /// the manifest. A group without a manifest has its member records
226    /// removed and nothing else.
227    pub(crate) async fn forget(&self, group_id: &RunId) -> Result<()> {
228        let mut keys = Vec::new();
229        if let Some(manifest) = self.read_manifest(group_id).await? {
230            for member in &manifest.members {
231                let run_id = member_run_id(group_id, &member.key);
232                self.memo_store.clear_memos_for_run(&run_id).await?;
233                keys.push(outcome_kv_key(&run_id));
234                if keys.len() == MEMBER_PAGE_SIZE {
235                    self.delete_keys(std::mem::take(&mut keys)).await?;
236                }
237            }
238        }
239        let prefix = group_members_kv_prefix(group_id);
240        let mut entries =
241            std::pin::pin!(self.queue.view().kv_entries(&prefix, .., MEMBER_PAGE_SIZE));
242        while let Some((key, _)) = entries.try_next().await? {
243            keys.push(key);
244            if keys.len() == MEMBER_PAGE_SIZE {
245                self.delete_keys(std::mem::take(&mut keys)).await?;
246            }
247        }
248        if !keys.is_empty() {
249            self.delete_keys(keys).await?;
250        }
251        self.objects.delete(&self.manifest_path(group_id)).await?;
252        Ok(())
253    }
254
255    /// Delete the KV entries under `keys` in one transaction.
256    async fn delete_keys(&self, keys: Vec<Vec<u8>>) -> Result<()> {
257        self.queue
258            .commit_effects(SettlementEffects::default().kv_deletes(keys))
259            .await?;
260        Ok(())
261    }
262}
263
264impl Clearable for GroupStore {
265    type Error = Error;
266
267    async fn clear(&self, group_id: &RunId) -> Result<Vec<Vec<u8>>> {
268        self.forget(group_id).await.map(|()| Vec::new())
269    }
270}
271
272/// A group of runs of one runtime, identified by a group id. Obtained
273/// from [`WorkflowRuntime::group`](crate::WorkflowRuntime::group) or [`WorkflowRuntime::new_group`](crate::WorkflowRuntime::new_group);
274/// cheap to clone.
275///
276/// The group's members are identified by key. A member's run id is
277/// derived from the group id and its key, so the same input submitted
278/// to two groups runs twice, and a second [`submit`](Self::submit) of
279/// the group runs again every member that did not succeed. The group's
280/// membership is a durable manifest, so [`results`](Self::results),
281/// [`status`](Self::status), [`cancel`](Self::cancel) and
282/// [`forget`](Self::forget) answer after a restart and from any runtime
283/// over the same queue.
284#[derive(Clone)]
285pub struct RunGroup {
286    runtime: Arc<RuntimeCore>,
287    id: RunId,
288}
289
290impl RunGroup {
291    pub(crate) fn new(runtime: Arc<RuntimeCore>, id: RunId) -> Self {
292        Self { runtime, id }
293    }
294
295    /// The group id.
296    pub fn id(&self) -> &RunId {
297        &self.id
298    }
299
300    fn core(&self) -> &RuntimeCore {
301        &self.runtime
302    }
303
304    fn store(&self) -> &GroupStore {
305        &self.core().group_store
306    }
307
308    /// The group's manifest; [`Error::GroupNotFound`] without one.
309    pub(crate) async fn manifest(&self) -> Result<Manifest> {
310        self.store()
311            .read_manifest(&self.id)
312            .await?
313            .ok_or_else(|| Error::GroupNotFound(self.id.clone()))
314    }
315
316    /// Every member record of the group, in key order.
317    pub(crate) async fn members(&self) -> Result<Vec<MemberState>> {
318        self.store().members(&self.id).await
319    }
320
321    /// Submit `members` as the group's members. The first submission
322    /// writes the group's manifest; a later one with a different member
323    /// set is rejected with [`Error::GroupMismatch`], and one with the
324    /// same set submits every member whose last recorded termination is
325    /// not a success, so a group is re-submitted after a crash or to run
326    /// its failed members again. Two members with one key are rejected
327    /// with [`Error::DuplicateMemberKey`]. `options` applies to every
328    /// member. A submission that fails returns after the members
329    /// submitted so far.
330    pub async fn submit(&self, members: Vec<GroupMember>, options: &RunOptions) -> Result<()> {
331        let mut seen = std::collections::HashSet::new();
332        for member in &members {
333            if !seen.insert(member.key.as_str()) {
334                return Err(Error::DuplicateMemberKey(member.key.clone()));
335            }
336        }
337        let manifest = Manifest {
338            group_id: self.id.clone(),
339            members,
340        };
341        match self.store().read_manifest(&self.id).await? {
342            Some(existing) if existing.members != manifest.members => {
343                return Err(Error::GroupMismatch(self.id.clone()));
344            }
345            Some(_) => {}
346            None => self.store().write_manifest(&manifest).await?,
347        }
348        self.submit_members(manifest.members, options).await
349    }
350
351    /// Submit the members of the group's manifest whose last recorded
352    /// termination is not a success: a member still active continues,
353    /// and the rest run. Returns [`Error::GroupNotFound`] for a group
354    /// never submitted. `options` applies as in [`submit`](Self::submit).
355    pub async fn resume(&self, options: &RunOptions) -> Result<()> {
356        let manifest = self.manifest().await?;
357        self.submit_members(manifest.members, options).await
358    }
359
360    async fn submit_members(&self, members: Vec<GroupMember>, options: &RunOptions) -> Result<()> {
361        let succeeded: std::collections::HashSet<String> = self
362            .members()
363            .await?
364            .into_iter()
365            .filter(|member| member.status() == Some(TerminalStatus::Succeeded))
366            .map(|member| member.key)
367            .collect();
368        let mut submissions = stream::iter(
369            members
370                .into_iter()
371                .filter(|member| !succeeded.contains(&member.key)),
372        )
373        .map(|member| async move {
374            let membership = Membership {
375                group_id: self.id.clone(),
376                key: member.key,
377            };
378            self.runtime
379                .submit_member(
380                    &membership,
381                    RunSpec {
382                        run_id: Some(membership.run_id()),
383                        input: member.input,
384                        options: options.clone(),
385                        effects: SettlementEffects::default(),
386                    },
387                )
388                .await
389                .map(|_| ())
390        })
391        .buffer_unordered(SUBMIT_CONCURRENCY);
392        while submissions.try_next().await?.is_some() {}
393        Ok(())
394    }
395
396    /// The members of the manifest as each one terminates, in
397    /// completion order; a member already terminated is yielded at
398    /// once. Every member must have been submitted. Once every member
399    /// has been yielded, the group's terminal marker is written, from
400    /// which the group retention sweep counts the window; a failed
401    /// marker write is logged.
402    async fn terminations(&self) -> Result<impl Stream<Item = Result<MemberState>> + use<>> {
403        let manifest = self.manifest().await?;
404        let waits: FuturesUnordered<_> = manifest
405            .members
406            .into_iter()
407            .map(|member| {
408                let group = self.clone();
409                async move {
410                    let record = group.wait_member(&member.key).await?;
411                    Ok(MemberState {
412                        key: member.key,
413                        record,
414                    })
415                }
416            })
417            .collect();
418        let group = self.clone();
419        let marker = stream::once(async move {
420            if let Err(err) = group.mark_terminated().await {
421                warn!(group_id = %group.id, "group terminal marker write failed: {err}");
422            }
423        })
424        .filter_map(|()| async { None::<Result<MemberState>> });
425        Ok(waits.chain(marker))
426    }
427
428    /// The members' results as each one terminates, in completion
429    /// order; a member already terminated is yielded at once. Returns
430    /// [`Error::GroupNotFound`] for a group never submitted.
431    pub async fn results(&self) -> Result<impl Stream<Item = Result<MemberResult>> + use<>> {
432        let terminations = self.terminations().await?;
433        let group = self.clone();
434        Ok(terminations.then(move |member| {
435            let group = group.clone();
436            async move {
437                let member = member?;
438                let termination: RunTermination = member
439                    .record
440                    .terminated
441                    .ok_or_else(|| Error::InconsistentRunState(member.record.run_id.clone()))?
442                    .into();
443                let outcome = group
444                    .core()
445                    .view
446                    .run_result_of(&member.record.run_id, &termination)
447                    .await?;
448                Ok(MemberResult {
449                    key: member.key,
450                    run_id: member.record.run_id,
451                    termination,
452                    outcome: outcome.map(|result| result.outcome),
453                })
454            }
455        }))
456    }
457
458    /// The group's durable state. Returns [`Error::GroupNotFound`] for a
459    /// group never submitted.
460    pub async fn status(&self) -> Result<GroupStatus> {
461        let manifest = self.manifest().await?;
462        let mut status = GroupStatus {
463            group_id: self.id.clone(),
464            total: manifest.members.len(),
465            pending: 0,
466            succeeded: 0,
467            failed: 0,
468            cancelled: 0,
469        };
470        for member in self.members().await? {
471            match member.status() {
472                None => status.pending += 1,
473                Some(TerminalStatus::Succeeded) => status.succeeded += 1,
474                Some(TerminalStatus::Failed) => status.failed += 1,
475                Some(TerminalStatus::Cancelled) => status.cancelled += 1,
476            }
477        }
478        Ok(status)
479    }
480
481    /// Request cancellation of every active member, as
482    /// [`WorkflowRuntime::cancel`](crate::WorkflowRuntime::cancel) does for one run. Returns the number
483    /// of members whose request was recorded.
484    pub async fn cancel(&self) -> Result<usize> {
485        let mut cancellations = stream::iter(
486            self.members()
487                .await?
488                .into_iter()
489                .filter(|member| member.status().is_none()),
490        )
491        .map(|member| async move { self.runtime.cancel(&member.record.run_id).await })
492        .buffer_unordered(SUBMIT_CONCURRENCY);
493        let mut cancelled = 0;
494        while let Some(recorded) = cancellations.try_next().await? {
495            cancelled += usize::from(recorded);
496        }
497        Ok(cancelled)
498    }
499
500    /// Wait until the member `key` terminates and return its record.
501    /// Returns [`Error::MemberNotSubmitted`] for a member of the
502    /// manifest without a record.
503    async fn wait_member(&self, key: &str) -> Result<DurableMember> {
504        let member = |member: Option<DurableMember>| {
505            member.ok_or_else(|| Error::MemberNotSubmitted {
506                group_id: self.id.clone(),
507                key: key.to_string(),
508            })
509        };
510        let record = member(self.store().member(&self.id, key).await?)?;
511        if record.terminated.is_some() {
512            return Ok(record);
513        }
514        let run_id = member_run_id(&self.id, key);
515        self.core().wait_run(&run_id).await?;
516        // The pointer and the member record change in one transaction,
517        // so the record of a run without a pointer is terminated.
518        let record = member(self.store().member(&self.id, key).await?)?;
519        if record.terminated.is_none() {
520            return Err(Error::InconsistentRunState(run_id));
521        }
522        Ok(record)
523    }
524
525    /// Write the group's terminal marker; no marker is written without
526    /// [`WorkflowRuntimeBuilder::group_retention`](crate::WorkflowRuntimeBuilder::group_retention).
527    async fn mark_terminated(&self) -> Result<()> {
528        let core = self.core();
529        if let Some(sweep) = &core.group_sweep {
530            let effects = sweep.mark(SettlementEffects::default(), &self.id, core.clock.now_ms());
531            core.queue.commit_effects(effects).await?;
532        }
533        Ok(())
534    }
535
536    /// Remove the group's state: its manifest, member records and the
537    /// memo entries, run result records and terminal records of its
538    /// members. A later [`submit`](Self::submit) under the same id
539    /// starts from nothing.
540    pub async fn forget(&self) -> Result<()> {
541        self.store().forget(&self.id).await
542    }
543}
544
545/// The member record written with a member's submission.
546pub(crate) fn pending_member(run_id: &RunId) -> DurableMember {
547    DurableMember {
548        run_id: run_id.clone(),
549        terminated: None,
550    }
551}
552
553/// The member record written by the settlement that terminates a member.
554pub(crate) fn terminated_member(run_id: &RunId, termination: DurableTermination) -> DurableMember {
555    DurableMember {
556        run_id: run_id.clone(),
557        terminated: Some(termination),
558    }
559}
560
561#[cfg(test)]
562mod tests {
563    use std::time::Duration;
564
565    use super::*;
566    use crate::runner::{Step, StepError, StepOutcome, StepRunner};
567    use crate::runtime::{RunSpec, WorkflowRuntime};
568    use crate::terminal::NoopTerminalHook;
569    use crate::test_util::{open_queue, open_queue_at, rid};
570
571    /// Continues once, then succeeds with the step number.
572    struct TwoSteps;
573
574    impl StepRunner for TwoSteps {
575        async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
576            if step.step_number == 0 {
577                Ok(StepOutcome::continue_now(step.payload.clone()))
578            } else {
579                Ok(StepOutcome::Succeed {
580                    result: step.step_number.to_string().into_bytes(),
581                })
582            }
583        }
584    }
585
586    fn member(key: &str) -> GroupMember {
587        GroupMember {
588            key: key.to_string(),
589            input: key.as_bytes().to_vec(),
590        }
591    }
592
593    /// Fails permanently.
594    struct Rejecting;
595
596    impl StepRunner for Rejecting {
597        async fn run_step(&self, _: &Step) -> std::result::Result<StepOutcome, StepError> {
598            Err(StepError::permanent("rejected"))
599        }
600    }
601
602    #[tokio::test(start_paused = true)]
603    async fn a_result_record_of_an_earlier_termination_is_not_reported_for_a_re_run_member() {
604        let (queue, store, clock) = open_queue_at(10_000).await;
605        let runtime =
606            WorkflowRuntime::builder(queue.clone(), store, Rejecting, NoopTerminalHook).build();
607        let group = runtime.group(rid("g"));
608        group
609            .submit(vec![member("a")], &RunOptions::default())
610            .await
611            .unwrap();
612        let worker = runtime.spawn(std::future::pending::<()>());
613        let first: Vec<MemberResult> = tokio::time::timeout(Duration::from_secs(10), async {
614            group.results().await.unwrap().try_collect().await
615        })
616        .await
617        .expect("results finished in time")
618        .unwrap();
619        assert_eq!(first[0].termination.error.as_deref(), Some("rejected"));
620        assert_eq!(
621            first[0].termination.error_kind,
622            Some(crate::StepErrorKind::Permanent)
623        );
624        assert!(
625            first[0].outcome.is_some(),
626            "the worker recorded the failure"
627        );
628        worker.shutdown().await.unwrap();
629
630        // The member runs again and the queue dead-letters it outside
631        // the worker, so reconciliation terminates it and no new record
632        // is written.
633        clock.advance(Duration::from_secs(1));
634        group
635            .submit(vec![member("a")], &RunOptions::default())
636            .await
637            .unwrap();
638        let claim = queue
639            .claim("workflow-steps", Duration::from_secs(60))
640            .await
641            .unwrap()
642            .unwrap();
643        queue.dead_letter(&claim, "hung").await.unwrap();
644        assert_eq!(runtime.inner.core.reconcile_dead_steps().await.unwrap(), 1);
645
646        let second: Vec<MemberResult> = group.results().await.unwrap().try_collect().await.unwrap();
647        assert_eq!(second[0].termination.error.as_deref(), Some("hung"));
648        assert_eq!(second[0].termination.error_kind, None);
649        assert_eq!(second[0].termination.terminated_at_ms, 11_000);
650        assert!(
651            second[0].outcome.is_none(),
652            "the first termination's record does not belong to the second",
653        );
654    }
655
656    #[tokio::test(start_paused = true)]
657    async fn membership_holds_across_steps_and_terminations_wait_for_the_last_one() {
658        let (queue, store) = open_queue().await;
659        let runtime = WorkflowRuntime::builder(queue.clone(), store, TwoSteps, NoopTerminalHook)
660            .poll_interval(Duration::from_millis(10))
661            .build();
662        let group = runtime.group(rid("g"));
663        group
664            .submit(vec![member("a"), member("b")], &RunOptions::default())
665            .await
666            .unwrap();
667        let pending = group.members().await.unwrap();
668        assert_eq!(pending.len(), 2);
669        assert!(pending.iter().all(|m| m.status().is_none()));
670        assert_eq!(pending[0].record.run_id, member_run_id(&rid("g"), "a"));
671
672        assert!(matches!(
673            group.submit(vec![member("a"), member("a")], &RunOptions::default()).await,
674            Err(Error::DuplicateMemberKey(key)) if key == "a"
675        ));
676
677        let worker = runtime.spawn(std::future::pending::<()>());
678        let terminated: Vec<MemberResult> = tokio::time::timeout(Duration::from_secs(10), async {
679            group.results().await.unwrap().try_collect().await
680        })
681        .await
682        .expect("results finished in time")
683        .unwrap();
684        assert_eq!(terminated.len(), 2);
685        for m in &terminated {
686            assert_eq!(m.termination.status, TerminalStatus::Succeeded);
687            let outcome = m.outcome.as_ref().expect("the worker recorded the outcome");
688            assert_eq!(outcome.run_id, m.run_id);
689            assert_eq!(
690                outcome.final_step, 1,
691                "the member terminated at its second step"
692            );
693            assert_eq!(outcome.result.as_deref(), Some(b"1".as_slice()));
694        }
695        let status = group.status().await.unwrap();
696        assert_eq!((status.total, status.succeeded, status.pending), (2, 2, 0));
697
698        // A member already terminated is yielded again at once.
699        let again: Vec<MemberResult> = group.results().await.unwrap().try_collect().await.unwrap();
700        assert_eq!(again.len(), 2);
701        worker.shutdown().await.unwrap();
702    }
703
704    #[tokio::test(start_paused = true)]
705    async fn the_run_options_apply_to_every_member() {
706        let (queue, store) = open_queue().await;
707        let runtime =
708            WorkflowRuntime::builder(queue.clone(), store, TwoSteps, NoopTerminalHook).build();
709        let group = runtime.group(rid("g"));
710        let options = RunOptions {
711            priority: Some(3),
712            max_attempts_per_step: Some(5),
713            headers: HashMap::from([("tenant".to_string(), "acme".to_string())]),
714            ..Default::default()
715        };
716        group
717            .submit(vec![member("a"), member("b")], &options)
718            .await
719            .unwrap();
720
721        for _ in 0..2 {
722            let job = queue
723                .claim("workflow-steps", Duration::from_secs(30))
724                .await
725                .unwrap()
726                .expect("a member's step job");
727            assert_eq!(job.priority, 3);
728            assert_eq!(job.max_attempts, 5);
729            assert_eq!(job.headers.get("tenant").map(String::as_str), Some("acme"));
730        }
731    }
732
733    #[tokio::test(start_paused = true)]
734    async fn a_group_cancellation_records_the_member_cancelled() {
735        let (queue, store, _clock) = open_queue_at(10_000).await;
736        let runtime =
737            WorkflowRuntime::builder(queue.clone(), store, TwoSteps, NoopTerminalHook).build();
738        let group = runtime.group(rid("g"));
739        group
740            .submit(vec![member("a")], &RunOptions::default())
741            .await
742            .unwrap();
743        let run_id = member_run_id(&rid("g"), "a");
744        assert_eq!(group.cancel().await.unwrap(), 1);
745        assert_eq!(group.cancel().await.unwrap(), 0, "no member is active");
746
747        let results: Vec<MemberResult> =
748            group.results().await.unwrap().try_collect().await.unwrap();
749        assert_eq!(results[0].run_id, run_id);
750        assert_eq!(
751            results[0].termination,
752            RunTermination {
753                status: TerminalStatus::Cancelled,
754                error: None,
755                error_kind: None,
756                final_step: 0,
757                terminated_at_ms: 10_000,
758            }
759        );
760        assert!(
761            results[0].outcome.is_none(),
762            "a pending step is cancelled without a worker, so no result is recorded",
763        );
764
765        // The cancelled member is submitted again; a member is grouped by
766        // its key, so a plain submission of the same run id is not.
767        group
768            .submit(vec![member("a")], &RunOptions::default())
769            .await
770            .unwrap();
771        assert!(group.members().await.unwrap()[0].status().is_none());
772        let plain = runtime
773            .submit(RunSpec {
774                run_id: Some(run_id.clone()),
775                input: b"a".to_vec(),
776                ..Default::default()
777            })
778            .await
779            .unwrap();
780        assert!(!plain.newly_submitted);
781
782        group.forget().await.unwrap();
783        assert!(group.members().await.unwrap().is_empty());
784        assert!(matches!(
785            group.manifest().await,
786            Err(Error::GroupNotFound(_))
787        ));
788    }
789
790    #[tokio::test(start_paused = true)]
791    async fn the_group_sweep_removes_the_members_records_with_the_group() {
792        let (queue, store, clock) = open_queue_at(10_000).await;
793        let runtime =
794            WorkflowRuntime::builder(queue.clone(), store.clone(), TwoSteps, NoopTerminalHook)
795                .group_retention(Duration::from_secs(1))
796                .build();
797        let group = runtime.group(rid("g"));
798        group
799            .submit(vec![member("a")], &RunOptions::default())
800            .await
801            .unwrap();
802        let run_id = member_run_id(&rid("g"), "a");
803        assert_eq!(group.cancel().await.unwrap(), 1);
804        let memos = crate::memo::MemoStore::new(store, "workflow-steps-memo");
805        memos.new_run_memo(&run_id).put("k", b"v").await.unwrap();
806        let results: Vec<MemberResult> =
807            group.results().await.unwrap().try_collect().await.unwrap();
808        assert_eq!(results.len(), 1);
809        let sweep = runtime.inner.core.group_sweep.as_ref().unwrap();
810        let marker = sweep.marker_key(&rid("g"), 10_000);
811        assert!(
812            queue.view().kv_get(&marker).await.unwrap().is_some(),
813            "the marker is written when the last termination is observed"
814        );
815        let terminal_record = outcome_kv_key(&run_id);
816        assert!(
817            queue
818                .view()
819                .kv_get(&terminal_record)
820                .await
821                .unwrap()
822                .is_some()
823        );
824
825        clock.advance(Duration::from_millis(999));
826        assert_eq!(
827            runtime.inner.core.sweep_once().await.unwrap(),
828            0,
829            "the marker is not yet expired"
830        );
831        clock.advance(Duration::from_millis(1));
832        assert_eq!(runtime.inner.core.sweep_once().await.unwrap(), 1);
833        assert!(queue.view().kv_get(&marker).await.unwrap().is_none());
834        assert!(group.members().await.unwrap().is_empty());
835        assert!(matches!(
836            group.manifest().await,
837            Err(Error::GroupNotFound(_))
838        ));
839        assert!(
840            queue
841                .view()
842                .kv_get(&terminal_record)
843                .await
844                .unwrap()
845                .is_none(),
846            "the member's terminal record is removed with the group"
847        );
848        assert!(
849            memos
850                .new_run_memo(&run_id)
851                .get("k")
852                .await
853                .unwrap()
854                .is_none()
855        );
856        assert!(
857            runtime.status(&run_id).await.unwrap().is_none(),
858            "nothing of the member remains"
859        );
860    }
861}