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    group_terminal_kv_key, 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, &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 = std::pin::pin!(self.queue.kv_entries(&prefix, MEMBER_PAGE_SIZE));
209        while let Some((kv_key, value)) = entries.try_next().await? {
210            let key = String::from_utf8_lossy(&kv_key[prefix.len()..]).into_owned();
211            if let Some(record) = durable::decode_or_absent(
212                &value,
213                "group member record",
214                &format_args!("{group_id}/{key}"),
215            ) {
216                members.push(MemberState { key, record });
217            }
218        }
219        Ok(members)
220    }
221
222    /// Remove the state of `group_id`: the memo entries and the terminal
223    /// record of every member in its manifest, its member records and
224    /// the manifest. A group without a manifest has its member records
225    /// removed and nothing else.
226    pub(crate) async fn forget(&self, group_id: &RunId) -> Result<()> {
227        let mut keys = Vec::new();
228        if let Some(manifest) = self.read_manifest(group_id).await? {
229            for member in &manifest.members {
230                let run_id = member_run_id(group_id, &member.key);
231                self.memo_store.clear_memos_for_run(&run_id).await?;
232                keys.push(outcome_kv_key(&run_id));
233                if keys.len() == MEMBER_PAGE_SIZE {
234                    self.delete_keys(std::mem::take(&mut keys)).await?;
235                }
236            }
237        }
238        let prefix = group_members_kv_prefix(group_id);
239        let mut entries = std::pin::pin!(self.queue.kv_entries(&prefix, MEMBER_PAGE_SIZE));
240        while let Some((key, _)) = entries.try_next().await? {
241            keys.push(key);
242            if keys.len() == MEMBER_PAGE_SIZE {
243                self.delete_keys(std::mem::take(&mut keys)).await?;
244            }
245        }
246        if !keys.is_empty() {
247            self.delete_keys(keys).await?;
248        }
249        self.objects.delete(&self.manifest_path(group_id)).await?;
250        Ok(())
251    }
252
253    /// Delete the KV entries under `keys` in one transaction.
254    async fn delete_keys(&self, keys: Vec<Vec<u8>>) -> Result<()> {
255        self.queue
256            .commit_effects(SettlementEffects::default().kv_deletes(keys))
257            .await?;
258        Ok(())
259    }
260}
261
262impl Clearable for GroupStore {
263    type Error = Error;
264
265    async fn clear(&self, group_id: &RunId) -> Result<Vec<Vec<u8>>> {
266        self.forget(group_id).await.map(|()| Vec::new())
267    }
268}
269
270/// A group of runs of one runtime, identified by a group id. Obtained
271/// from [`WorkflowRuntime::group`](crate::WorkflowRuntime::group) or [`WorkflowRuntime::new_group`](crate::WorkflowRuntime::new_group);
272/// cheap to clone.
273///
274/// The group's members are identified by key. A member's run id is
275/// derived from the group id and its key, so the same input submitted
276/// to two groups runs twice, and a second [`submit`](Self::submit) of
277/// the group runs again every member that did not succeed. The group's
278/// membership is a durable manifest, so [`results`](Self::results),
279/// [`status`](Self::status), [`cancel`](Self::cancel) and
280/// [`forget`](Self::forget) answer after a restart and from any runtime
281/// over the same queue.
282#[derive(Clone)]
283pub struct RunGroup {
284    runtime: Arc<RuntimeCore>,
285    id: RunId,
286}
287
288impl RunGroup {
289    pub(crate) fn new(runtime: Arc<RuntimeCore>, id: RunId) -> Self {
290        Self { runtime, id }
291    }
292
293    /// The group id.
294    pub fn id(&self) -> &RunId {
295        &self.id
296    }
297
298    fn core(&self) -> &RuntimeCore {
299        &self.runtime
300    }
301
302    fn store(&self) -> &GroupStore {
303        &self.core().group_store
304    }
305
306    /// The group's manifest; [`Error::GroupNotFound`] without one.
307    pub(crate) async fn manifest(&self) -> Result<Manifest> {
308        self.store()
309            .read_manifest(&self.id)
310            .await?
311            .ok_or_else(|| Error::GroupNotFound(self.id.clone()))
312    }
313
314    /// Every member record of the group, in key order.
315    pub(crate) async fn members(&self) -> Result<Vec<MemberState>> {
316        self.store().members(&self.id).await
317    }
318
319    /// Submit `members` as the group's members. The first submission
320    /// writes the group's manifest; a later one with a different member
321    /// set is rejected with [`Error::GroupMismatch`], and one with the
322    /// same set submits every member whose last recorded termination is
323    /// not a success, so a group is re-submitted after a crash or to run
324    /// its failed members again. Two members with one key are rejected
325    /// with [`Error::DuplicateMemberKey`]. `options` applies to every
326    /// member. A submission that fails returns after the members
327    /// submitted so far.
328    pub async fn submit(&self, members: Vec<GroupMember>, options: &RunOptions) -> Result<()> {
329        let mut seen = std::collections::HashSet::new();
330        for member in &members {
331            if !seen.insert(member.key.as_str()) {
332                return Err(Error::DuplicateMemberKey(member.key.clone()));
333            }
334        }
335        let manifest = Manifest {
336            group_id: self.id.clone(),
337            members,
338        };
339        match self.store().read_manifest(&self.id).await? {
340            Some(existing) if existing.members != manifest.members => {
341                return Err(Error::GroupMismatch(self.id.clone()));
342            }
343            Some(_) => {}
344            None => self.store().write_manifest(&manifest).await?,
345        }
346        self.submit_members(manifest.members, options).await
347    }
348
349    /// Submit the members of the group's manifest whose last recorded
350    /// termination is not a success: a member still active continues,
351    /// and the rest run. Returns [`Error::GroupNotFound`] for a group
352    /// never submitted. `options` applies as in [`submit`](Self::submit).
353    pub async fn resume(&self, options: &RunOptions) -> Result<()> {
354        let manifest = self.manifest().await?;
355        self.submit_members(manifest.members, options).await
356    }
357
358    async fn submit_members(&self, members: Vec<GroupMember>, options: &RunOptions) -> Result<()> {
359        let succeeded: std::collections::HashSet<String> = self
360            .members()
361            .await?
362            .into_iter()
363            .filter(|member| member.status() == Some(TerminalStatus::Succeeded))
364            .map(|member| member.key)
365            .collect();
366        let mut submissions = stream::iter(
367            members
368                .into_iter()
369                .filter(|member| !succeeded.contains(&member.key)),
370        )
371        .map(|member| async move {
372            let membership = Membership {
373                group_id: self.id.clone(),
374                key: member.key,
375            };
376            self.runtime
377                .submit_member(
378                    &membership,
379                    RunSpec {
380                        run_id: Some(membership.run_id()),
381                        input: member.input,
382                        options: options.clone(),
383                        kv_writes: HashMap::new(),
384                    },
385                )
386                .await
387                .map(|_| ())
388        })
389        .buffer_unordered(SUBMIT_CONCURRENCY);
390        while submissions.try_next().await?.is_some() {}
391        Ok(())
392    }
393
394    /// The members of the manifest as each one terminates, in
395    /// completion order; a member already terminated is yielded at
396    /// once. Every member must have been submitted. Once every member
397    /// has been yielded, the group's terminal marker is written, from
398    /// which the group retention sweep counts the window; a failed
399    /// marker write is logged.
400    async fn terminations(&self) -> Result<impl Stream<Item = Result<MemberState>> + use<>> {
401        let manifest = self.manifest().await?;
402        let waits: FuturesUnordered<_> = manifest
403            .members
404            .into_iter()
405            .map(|member| {
406                let group = self.clone();
407                async move {
408                    let record = group.wait_member(&member.key).await?;
409                    Ok(MemberState {
410                        key: member.key,
411                        record,
412                    })
413                }
414            })
415            .collect();
416        let group = self.clone();
417        let marker = stream::once(async move {
418            if let Err(err) = group.mark_terminated().await {
419                warn!(group_id = %group.id, "group terminal marker write failed: {err}");
420            }
421        })
422        .filter_map(|()| async { None::<Result<MemberState>> });
423        Ok(waits.chain(marker))
424    }
425
426    /// The members' results as each one terminates, in completion
427    /// order; a member already terminated is yielded at once. Returns
428    /// [`Error::GroupNotFound`] for a group never submitted.
429    pub async fn results(&self) -> Result<impl Stream<Item = Result<MemberResult>> + use<>> {
430        let terminations = self.terminations().await?;
431        let group = self.clone();
432        Ok(terminations.then(move |member| {
433            let group = group.clone();
434            async move {
435                let member = member?;
436                let termination: RunTermination = member
437                    .record
438                    .terminated
439                    .ok_or_else(|| Error::InconsistentRunState(member.record.run_id.clone()))?
440                    .into();
441                let outcome = group
442                    .core()
443                    .run_result_of(&member.record.run_id, &termination)
444                    .await?;
445                Ok(MemberResult {
446                    key: member.key,
447                    run_id: member.record.run_id,
448                    termination,
449                    outcome: outcome.map(|result| result.outcome),
450                })
451            }
452        }))
453    }
454
455    /// The group's durable state. Returns [`Error::GroupNotFound`] for a
456    /// group never submitted.
457    pub async fn status(&self) -> Result<GroupStatus> {
458        let manifest = self.manifest().await?;
459        let mut status = GroupStatus {
460            group_id: self.id.clone(),
461            total: manifest.members.len(),
462            pending: 0,
463            succeeded: 0,
464            failed: 0,
465            cancelled: 0,
466        };
467        for member in self.members().await? {
468            match member.status() {
469                None => status.pending += 1,
470                Some(TerminalStatus::Succeeded) => status.succeeded += 1,
471                Some(TerminalStatus::Failed) => status.failed += 1,
472                Some(TerminalStatus::Cancelled) => status.cancelled += 1,
473            }
474        }
475        Ok(status)
476    }
477
478    /// Request cancellation of every active member, as
479    /// [`WorkflowRuntime::cancel`](crate::WorkflowRuntime::cancel) does for one run. Returns the number
480    /// of members whose request was recorded.
481    pub async fn cancel(&self) -> Result<usize> {
482        let mut cancellations = stream::iter(
483            self.members()
484                .await?
485                .into_iter()
486                .filter(|member| member.status().is_none()),
487        )
488        .map(|member| async move { self.runtime.cancel(&member.record.run_id).await })
489        .buffer_unordered(SUBMIT_CONCURRENCY);
490        let mut cancelled = 0;
491        while let Some(recorded) = cancellations.try_next().await? {
492            cancelled += usize::from(recorded);
493        }
494        Ok(cancelled)
495    }
496
497    /// Wait until the member `key` terminates and return its record.
498    /// Returns [`Error::MemberNotSubmitted`] for a member of the
499    /// manifest without a record.
500    async fn wait_member(&self, key: &str) -> Result<DurableMember> {
501        let member = |member: Option<DurableMember>| {
502            member.ok_or_else(|| Error::MemberNotSubmitted {
503                group_id: self.id.clone(),
504                key: key.to_string(),
505            })
506        };
507        let record = member(self.store().member(&self.id, key).await?)?;
508        if record.terminated.is_some() {
509            return Ok(record);
510        }
511        let run_id = member_run_id(&self.id, key);
512        self.core().wait_run(&run_id).await?;
513        // The pointer and the member record change in one transaction,
514        // so the record of a run without a pointer is terminated.
515        let record = member(self.store().member(&self.id, key).await?)?;
516        if record.terminated.is_none() {
517            return Err(Error::InconsistentRunState(run_id));
518        }
519        Ok(record)
520    }
521
522    /// Write the group's terminal marker; no marker is written without
523    /// [`WorkflowRuntimeBuilder::group_retention`](crate::WorkflowRuntimeBuilder::group_retention).
524    async fn mark_terminated(&self) -> Result<()> {
525        let core = self.core();
526        if core.group_retention.is_some() {
527            let key = group_terminal_kv_key(&self.id, core.clock.now_ms());
528            core.queue.kv_put(&key, b"").await?;
529        }
530        Ok(())
531    }
532
533    /// Remove the group's state: its manifest, member records and the
534    /// memo entries, run result records and terminal records of its
535    /// members. A later [`submit`](Self::submit) under the same id
536    /// starts from nothing.
537    pub async fn forget(&self) -> Result<()> {
538        self.store().forget(&self.id).await
539    }
540}
541
542/// The member record written with a member's submission.
543pub(crate) fn pending_member(run_id: &RunId) -> DurableMember {
544    DurableMember {
545        run_id: run_id.clone(),
546        terminated: None,
547    }
548}
549
550/// The member record written by the settlement that terminates a member.
551pub(crate) fn terminated_member(run_id: &RunId, termination: DurableTermination) -> DurableMember {
552    DurableMember {
553        run_id: run_id.clone(),
554        terminated: Some(termination),
555    }
556}
557
558#[cfg(test)]
559mod tests {
560    use std::time::Duration;
561
562    use super::*;
563    use crate::runner::{Step, StepError, StepOutcome, StepRunner};
564    use crate::runtime::{RunSpec, WorkflowRuntime};
565    use crate::terminal::NoopTerminalHook;
566    use crate::test_util::{open_queue, open_queue_at, rid};
567
568    /// Continues once, then succeeds with the step number.
569    struct TwoSteps;
570
571    impl StepRunner for TwoSteps {
572        async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
573            if step.step_number == 0 {
574                Ok(StepOutcome::continue_now(step.payload.clone()))
575            } else {
576                Ok(StepOutcome::Succeed {
577                    result: step.step_number.to_string().into_bytes(),
578                })
579            }
580        }
581    }
582
583    fn member(key: &str) -> GroupMember {
584        GroupMember {
585            key: key.to_string(),
586            input: key.as_bytes().to_vec(),
587        }
588    }
589
590    /// Fails permanently.
591    struct Rejecting;
592
593    impl StepRunner for Rejecting {
594        async fn run_step(&self, _: &Step) -> std::result::Result<StepOutcome, StepError> {
595            Err(StepError::permanent("rejected"))
596        }
597    }
598
599    #[tokio::test(start_paused = true)]
600    async fn a_result_record_of_an_earlier_termination_is_not_reported_for_a_re_run_member() {
601        let (queue, store, clock) = open_queue_at(10_000).await;
602        let runtime =
603            WorkflowRuntime::builder(queue.clone(), store, Rejecting, NoopTerminalHook).build();
604        let group = runtime.group(rid("g"));
605        group
606            .submit(vec![member("a")], &RunOptions::default())
607            .await
608            .unwrap();
609        let worker = runtime.spawn(std::future::pending::<()>());
610        let first: Vec<MemberResult> = tokio::time::timeout(Duration::from_secs(10), async {
611            group.results().await.unwrap().try_collect().await
612        })
613        .await
614        .expect("results finished in time")
615        .unwrap();
616        assert_eq!(first[0].termination.error.as_deref(), Some("rejected"));
617        assert_eq!(
618            first[0].termination.error_kind,
619            Some(crate::StepErrorKind::Permanent)
620        );
621        assert!(
622            first[0].outcome.is_some(),
623            "the worker recorded the failure"
624        );
625        worker.shutdown().await.unwrap();
626
627        // The member runs again and the queue dead-letters it outside
628        // the worker, so reconciliation terminates it and no new record
629        // is written.
630        clock.advance(Duration::from_secs(1));
631        group
632            .submit(vec![member("a")], &RunOptions::default())
633            .await
634            .unwrap();
635        let claim = queue
636            .claim("workflow-steps", Duration::from_secs(60))
637            .await
638            .unwrap()
639            .unwrap();
640        queue.dead_letter(&claim, "hung").await.unwrap();
641        assert_eq!(runtime.inner.core.reconcile_dead_steps().await.unwrap(), 1);
642
643        let second: Vec<MemberResult> = group.results().await.unwrap().try_collect().await.unwrap();
644        assert_eq!(second[0].termination.error.as_deref(), Some("hung"));
645        assert_eq!(second[0].termination.error_kind, None);
646        assert_eq!(second[0].termination.terminated_at_ms, 11_000);
647        assert!(
648            second[0].outcome.is_none(),
649            "the first termination's record does not belong to the second",
650        );
651    }
652
653    #[tokio::test(start_paused = true)]
654    async fn membership_holds_across_steps_and_terminations_wait_for_the_last_one() {
655        let (queue, store) = open_queue().await;
656        let runtime = WorkflowRuntime::builder(queue.clone(), store, TwoSteps, NoopTerminalHook)
657            .poll_interval(Duration::from_millis(10))
658            .build();
659        let group = runtime.group(rid("g"));
660        group
661            .submit(vec![member("a"), member("b")], &RunOptions::default())
662            .await
663            .unwrap();
664        let pending = group.members().await.unwrap();
665        assert_eq!(pending.len(), 2);
666        assert!(pending.iter().all(|m| m.status().is_none()));
667        assert_eq!(pending[0].record.run_id, member_run_id(&rid("g"), "a"));
668
669        assert!(matches!(
670            group.submit(vec![member("a"), member("a")], &RunOptions::default()).await,
671            Err(Error::DuplicateMemberKey(key)) if key == "a"
672        ));
673
674        let worker = runtime.spawn(std::future::pending::<()>());
675        let terminated: Vec<MemberResult> = tokio::time::timeout(Duration::from_secs(10), async {
676            group.results().await.unwrap().try_collect().await
677        })
678        .await
679        .expect("results finished in time")
680        .unwrap();
681        assert_eq!(terminated.len(), 2);
682        for m in &terminated {
683            assert_eq!(m.termination.status, TerminalStatus::Succeeded);
684            let outcome = m.outcome.as_ref().expect("the worker recorded the outcome");
685            assert_eq!(outcome.run_id, m.run_id);
686            assert_eq!(
687                outcome.final_step, 1,
688                "the member terminated at its second step"
689            );
690            assert_eq!(outcome.result.as_deref(), Some(b"1".as_slice()));
691        }
692        let status = group.status().await.unwrap();
693        assert_eq!((status.total, status.succeeded, status.pending), (2, 2, 0));
694
695        // A member already terminated is yielded again at once.
696        let again: Vec<MemberResult> = group.results().await.unwrap().try_collect().await.unwrap();
697        assert_eq!(again.len(), 2);
698        worker.shutdown().await.unwrap();
699    }
700
701    #[tokio::test(start_paused = true)]
702    async fn the_run_options_apply_to_every_member() {
703        let (queue, store) = open_queue().await;
704        let runtime =
705            WorkflowRuntime::builder(queue.clone(), store, TwoSteps, NoopTerminalHook).build();
706        let group = runtime.group(rid("g"));
707        let options = RunOptions {
708            priority: Some(3),
709            max_attempts_per_step: Some(5),
710            headers: HashMap::from([("tenant".to_string(), "acme".to_string())]),
711            ..Default::default()
712        };
713        group
714            .submit(vec![member("a"), member("b")], &options)
715            .await
716            .unwrap();
717
718        for _ in 0..2 {
719            let job = queue
720                .claim("workflow-steps", Duration::from_secs(30))
721                .await
722                .unwrap()
723                .expect("a member's step job");
724            assert_eq!(job.priority, 3);
725            assert_eq!(job.max_attempts, 5);
726            assert_eq!(job.headers.get("tenant").map(String::as_str), Some("acme"));
727        }
728    }
729
730    #[tokio::test(start_paused = true)]
731    async fn a_group_cancellation_records_the_member_cancelled() {
732        let (queue, store, _clock) = open_queue_at(10_000).await;
733        let runtime =
734            WorkflowRuntime::builder(queue.clone(), store, TwoSteps, NoopTerminalHook).build();
735        let group = runtime.group(rid("g"));
736        group
737            .submit(vec![member("a")], &RunOptions::default())
738            .await
739            .unwrap();
740        let run_id = member_run_id(&rid("g"), "a");
741        assert_eq!(group.cancel().await.unwrap(), 1);
742        assert_eq!(group.cancel().await.unwrap(), 0, "no member is active");
743
744        let results: Vec<MemberResult> =
745            group.results().await.unwrap().try_collect().await.unwrap();
746        assert_eq!(results[0].run_id, run_id);
747        assert_eq!(
748            results[0].termination,
749            RunTermination {
750                status: TerminalStatus::Cancelled,
751                error: None,
752                error_kind: None,
753                final_step: 0,
754                terminated_at_ms: 10_000,
755            }
756        );
757        assert!(
758            results[0].outcome.is_none(),
759            "a pending step is cancelled without a worker, so no result is recorded",
760        );
761
762        // The cancelled member is submitted again; a member is grouped by
763        // its key, so a plain submission of the same run id is not.
764        group
765            .submit(vec![member("a")], &RunOptions::default())
766            .await
767            .unwrap();
768        assert!(group.members().await.unwrap()[0].status().is_none());
769        let plain = runtime
770            .submit(RunSpec {
771                run_id: Some(run_id.clone()),
772                input: b"a".to_vec(),
773                ..Default::default()
774            })
775            .await
776            .unwrap();
777        assert!(!plain.newly_submitted);
778
779        group.forget().await.unwrap();
780        assert!(group.members().await.unwrap().is_empty());
781        assert!(matches!(
782            group.manifest().await,
783            Err(Error::GroupNotFound(_))
784        ));
785    }
786
787    #[tokio::test(start_paused = true)]
788    async fn the_group_sweep_removes_the_members_records_with_the_group() {
789        let (queue, store, clock) = open_queue_at(10_000).await;
790        let runtime =
791            WorkflowRuntime::builder(queue.clone(), store.clone(), TwoSteps, NoopTerminalHook)
792                .group_retention(Duration::from_secs(1))
793                .build();
794        let group = runtime.group(rid("g"));
795        group
796            .submit(vec![member("a")], &RunOptions::default())
797            .await
798            .unwrap();
799        let run_id = member_run_id(&rid("g"), "a");
800        assert_eq!(group.cancel().await.unwrap(), 1);
801        let memos = crate::memo::MemoStore::new(store, "workflow-steps-memo");
802        memos.new_run_memo(&run_id).put("k", b"v").await.unwrap();
803        let results: Vec<MemberResult> =
804            group.results().await.unwrap().try_collect().await.unwrap();
805        assert_eq!(results.len(), 1);
806        let marker = group_terminal_kv_key(&rid("g"), 10_000);
807        assert!(
808            queue.kv_get(&marker).await.unwrap().is_some(),
809            "the marker is written when the last termination is observed"
810        );
811        let terminal_record = outcome_kv_key(&run_id);
812        assert!(queue.kv_get(&terminal_record).await.unwrap().is_some());
813
814        clock.advance(Duration::from_secs(1));
815        assert_eq!(
816            runtime.inner.core.sweep_once().await.unwrap(),
817            0,
818            "the marker is not yet expired"
819        );
820        clock.advance(Duration::from_millis(1));
821        assert_eq!(runtime.inner.core.sweep_once().await.unwrap(), 1);
822        assert!(queue.kv_get(&marker).await.unwrap().is_none());
823        assert!(group.members().await.unwrap().is_empty());
824        assert!(matches!(
825            group.manifest().await,
826            Err(Error::GroupNotFound(_))
827        ));
828        assert!(
829            queue.kv_get(&terminal_record).await.unwrap().is_none(),
830            "the member's terminal record is removed with the group"
831        );
832        assert!(
833            memos
834                .new_run_memo(&run_id)
835                .get("k")
836                .await
837                .unwrap()
838                .is_none()
839        );
840        assert!(
841            runtime.status(&run_id).await.unwrap().is_none(),
842            "nothing of the member remains"
843        );
844    }
845}