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 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
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, &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 = 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 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 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#[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 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 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 pub(crate) async fn members(&self) -> Result<Vec<MemberState>> {
316 self.store().members(&self.id).await
317 }
318
319 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 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 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 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 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 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 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 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 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 pub async fn forget(&self) -> Result<()> {
538 self.store().forget(&self.id).await
539 }
540}
541
542pub(crate) fn pending_member(run_id: &RunId) -> DurableMember {
544 DurableMember {
545 run_id: run_id.clone(),
546 terminated: None,
547 }
548}
549
550pub(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 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 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 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 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 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}