1use 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
29const SUBMIT_CONCURRENCY: usize = 32;
32const MEMBER_PAGE_SIZE: usize = 1000;
35
36pub(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#[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 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 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 pub(crate) fn kv_key(&self) -> Vec<u8> {
75 group_member_kv_key(&self.group_id, &self.key)
76 }
77
78 pub(crate) fn run_id(&self) -> RunId {
80 member_run_id(&self.group_id, &self.key)
81 }
82}
83
84#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
87pub struct GroupMember {
88 pub key: String,
90 pub input: Vec<u8>,
92}
93
94#[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#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct GroupStatus {
105 pub group_id: RunId,
107 pub total: usize,
109 pub pending: usize,
111 pub succeeded: usize,
113 pub failed: usize,
115 pub cancelled: usize,
117}
118
119#[derive(Debug, Clone)]
121pub struct MemberResult {
122 pub key: String,
124 pub run_id: RunId,
126 pub termination: RunTermination,
128 pub outcome: Option<RunOutcome>,
131}
132
133pub(crate) struct MemberState {
135 pub(crate) key: String,
136 pub(crate) record: DurableMember,
137}
138
139impl MemberState {
140 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#[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 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 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 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 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#[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 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 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 pub(crate) async fn members(&self) -> Result<Vec<MemberState>> {
318 self.store().members(&self.id).await
319 }
320
321 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 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 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 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 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 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 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 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 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 pub async fn forget(&self) -> Result<()> {
541 self.store().forget(&self.id).await
542 }
543}
544
545pub(crate) fn pending_member(run_id: &RunId) -> DurableMember {
547 DurableMember {
548 run_id: run_id.clone(),
549 terminated: None,
550 }
551}
552
553pub(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 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 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 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 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 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}