1use std::num::NonZeroU64;
4use std::pin::Pin;
5use std::sync::Arc;
6use std::task::{Context, Poll};
7
8use aion_core::{Event, WorkflowFilter, WorkflowId};
9use aion_proto::{
10 FilteredSubscription, FirehoseSubscription, PerWorkflowSubscription, ProtoWorkflowId,
11 SubscriptionRequest, subscription_request,
12};
13use futures::Stream;
14use futures::future::BoxFuture;
15use futures::stream::BoxStream;
16
17use crate::error::ClientError;
18use crate::transport::{SubscriptionAttempt, WorkflowTransport};
19
20pub type EventStream = Pin<Box<dyn Stream<Item = Result<Event, ClientError>> + Send>>;
22
23#[derive(Clone, Debug, PartialEq, Eq)]
25pub enum SubscribeTarget {
26 Workflow {
28 workflow_id: WorkflowId,
30 },
31 Filtered {
33 filter: WorkflowFilter,
35 },
36 Firehose,
38}
39
40impl SubscribeTarget {
41 pub(crate) fn request(&self, namespace: &str) -> SubscriptionRequest {
42 match self {
43 Self::Workflow { workflow_id } => SubscriptionRequest {
44 subscription: Some(subscription_request::Subscription::PerWorkflow(
45 PerWorkflowSubscription {
46 namespace: namespace.to_owned(),
47 workflow_id: Some(ProtoWorkflowId::from(workflow_id.clone())),
48 resume_from_seq: None,
49 },
50 )),
51 },
52 Self::Filtered { filter } => SubscriptionRequest {
53 subscription: Some(subscription_request::Subscription::Filtered(
54 FilteredSubscription {
55 namespace: namespace.to_owned(),
56 workflow_type: filter.workflow_type.clone(),
57 status: filter
58 .status
59 .map(|status| aion_proto::ProtoWorkflowStatus::from(status) as i32),
60 namespace_selector: None,
61 },
62 )),
63 },
64 Self::Firehose => SubscriptionRequest {
65 subscription: Some(subscription_request::Subscription::Firehose(
66 FirehoseSubscription {
67 namespace: namespace.to_owned(),
68 },
69 )),
70 },
71 }
72 }
73}
74
75pub struct ResumingEventStream {
94 transport: Arc<dyn WorkflowTransport>,
95 namespace: String,
96 target: SubscribeTarget,
97 last_seq: Option<u64>,
98 delivered_any: bool,
99 current: Option<BoxStream<'static, Result<Event, ClientError>>>,
100 pending_subscribe: Option<BoxFuture<'static, Result<SubscriptionAttempt, ClientError>>>,
101 terminal_error: Option<ClientError>,
102 finished: bool,
103}
104
105impl ResumingEventStream {
106 #[must_use]
108 pub fn new(
109 transport: Arc<dyn WorkflowTransport>,
110 namespace: impl Into<String>,
111 target: SubscribeTarget,
112 ) -> Self {
113 Self {
114 transport,
115 namespace: namespace.into(),
116 target,
117 last_seq: None,
118 delivered_any: false,
119 current: None,
120 pending_subscribe: None,
121 terminal_error: None,
122 finished: false,
123 }
124 }
125
126 #[must_use]
134 pub fn from_sequence(
135 transport: Arc<dyn WorkflowTransport>,
136 namespace: impl Into<String>,
137 workflow_id: WorkflowId,
138 resume_from: NonZeroU64,
139 ) -> Self {
140 let mut stream = Self::new(
141 transport,
142 namespace,
143 SubscribeTarget::Workflow { workflow_id },
144 );
145 stream.last_seq = Some(resume_from.get() - 1);
149 stream
150 }
151
152 fn is_per_workflow(&self) -> bool {
153 matches!(self.target, SubscribeTarget::Workflow { .. })
154 }
155
156 fn start_subscribe(&mut self) {
157 let transport = Arc::clone(&self.transport);
158 let request = self.target.request(&self.namespace);
159 let resume_from_sequence = if self.is_per_workflow() {
162 self.last_seq.map(|seq| seq.saturating_add(1))
163 } else {
164 None
165 };
166 self.pending_subscribe = Some(Box::pin(async move {
167 transport.subscribe(request, resume_from_sequence).await
168 }));
169 }
170}
171
172impl Stream for ResumingEventStream {
173 type Item = Result<Event, ClientError>;
174
175 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
176 let this = self.get_mut();
177 loop {
178 if this.finished {
179 return Poll::Ready(None);
180 }
181
182 if let Some(error) = this.terminal_error.take() {
183 this.finished = true;
184 return Poll::Ready(Some(Err(error)));
185 }
186
187 if this.current.is_none() && this.pending_subscribe.is_none() {
188 this.start_subscribe();
189 }
190
191 if let Some(pending) = this.pending_subscribe.as_mut() {
192 match pending.as_mut().poll(cx) {
193 Poll::Pending => return Poll::Pending,
194 Poll::Ready(Ok(attempt)) => {
195 this.pending_subscribe = None;
196 this.current = Some(attempt.events);
197 }
198 Poll::Ready(Err(error)) => {
199 this.pending_subscribe = None;
206 if is_retryable(&error) && (this.is_per_workflow() || !this.delivered_any) {
207 continue;
208 }
209 this.finished = true;
210 return Poll::Ready(Some(Err(error)));
211 }
212 }
213 }
214
215 let Some(current) = this.current.as_mut() else {
216 continue;
217 };
218 match current.as_mut().poll_next(cx) {
219 Poll::Pending => return Poll::Pending,
220 Poll::Ready(Some(Ok(event))) => {
221 if this.is_per_workflow() {
222 if this.last_seq.is_some_and(|seq| event.seq() <= seq) {
225 continue;
226 }
227 this.last_seq = Some(event.seq());
228 }
229 this.delivered_any = true;
230 return Poll::Ready(Some(Ok(event)));
231 }
232 Poll::Ready(Some(Err(error))) => {
233 this.current = None;
234 if is_retryable(&error) {
235 if this.is_per_workflow() {
236 continue;
237 }
238 if !this.delivered_any {
239 continue;
242 }
243 }
247 this.terminal_error = Some(error);
248 }
249 Poll::Ready(None) => {
250 this.current = None;
251 this.finished = true;
252 return Poll::Ready(None);
253 }
254 }
255 }
256 }
257}
258
259#[must_use]
261pub fn event_stream(
262 transport: Arc<dyn WorkflowTransport>,
263 namespace: impl Into<String>,
264 target: SubscribeTarget,
265) -> EventStream {
266 Box::pin(ResumingEventStream::new(transport, namespace, target))
267}
268
269#[must_use]
271pub fn event_stream_from(
272 transport: Arc<dyn WorkflowTransport>,
273 namespace: impl Into<String>,
274 workflow_id: WorkflowId,
275 resume_from: NonZeroU64,
276) -> EventStream {
277 Box::pin(ResumingEventStream::from_sequence(
278 transport,
279 namespace,
280 workflow_id,
281 resume_from,
282 ))
283}
284
285fn is_retryable(error: &ClientError) -> bool {
286 matches!(error, ClientError::Unavailable { .. })
287}
288
289#[cfg(test)]
290mod tests {
291 use std::collections::VecDeque;
292 use std::sync::Arc;
293
294 use aion_core::{ContentType, Event, EventEnvelope, Payload, WorkflowId};
295 use aion_proto::{
296 ProtoCancelResponse, ProtoDescribeWorkflowResponse, ProtoListWorkflowsResponse,
297 ProtoQueryResponse, ProtoSignalResponse, ProtoStartWorkflowResponse,
298 };
299 use async_trait::async_trait;
300 use chrono::Utc;
301 use futures::StreamExt;
302 use futures::stream;
303 use tokio::sync::Mutex;
304
305 use super::{ResumingEventStream, SubscribeTarget};
306 use crate::error::ClientError;
307 use crate::transport::{SubscriptionAttempt, WorkflowTransport};
308
309 #[derive(Default)]
310 struct SubscribeStub {
311 attach_failures: Mutex<VecDeque<ClientError>>,
314 attempts: Mutex<VecDeque<SubscriptionAttempt>>,
315 resume_points: Mutex<Vec<Option<u64>>>,
316 }
317
318 #[async_trait]
319 impl WorkflowTransport for SubscribeStub {
320 async fn start_workflow(
321 &self,
322 _: aion_proto::ProtoStartWorkflowRequest,
323 ) -> Result<ProtoStartWorkflowResponse, ClientError> {
324 Err(ClientError::unavailable("stub transport"))
325 }
326
327 async fn signal(
328 &self,
329 _: aion_proto::ProtoSignalRequest,
330 ) -> Result<ProtoSignalResponse, ClientError> {
331 Err(ClientError::unavailable("stub transport"))
332 }
333
334 async fn query(
335 &self,
336 _: aion_proto::ProtoQueryRequest,
337 ) -> Result<ProtoQueryResponse, ClientError> {
338 Err(ClientError::unavailable("stub transport"))
339 }
340
341 async fn cancel(
342 &self,
343 _: aion_proto::ProtoCancelRequest,
344 ) -> Result<ProtoCancelResponse, ClientError> {
345 Err(ClientError::unavailable("stub transport"))
346 }
347
348 async fn retire_workloop(
349 &self,
350 _: aion_proto::ProtoRetireWorkloopRequest,
351 ) -> Result<aion_proto::ProtoRetireWorkloopResponse, ClientError> {
352 Err(ClientError::unavailable("stub transport"))
353 }
354
355 async fn reopen(
356 &self,
357 _: aion_proto::ProtoReopenRequest,
358 ) -> Result<aion_proto::ProtoReopenResponse, ClientError> {
359 Err(ClientError::unavailable("stub transport"))
360 }
361
362 async fn pause(
363 &self,
364 _: aion_proto::ProtoPauseRequest,
365 ) -> Result<aion_proto::ProtoPauseResponse, ClientError> {
366 Err(ClientError::unavailable("stub transport"))
367 }
368
369 async fn resume(
370 &self,
371 _: aion_proto::ProtoResumeRequest,
372 ) -> Result<aion_proto::ProtoResumeResponse, ClientError> {
373 Err(ClientError::unavailable("stub transport"))
374 }
375
376 async fn list_workflows(
377 &self,
378 _: aion_proto::ProtoListWorkflowsRequest,
379 ) -> Result<ProtoListWorkflowsResponse, ClientError> {
380 Err(ClientError::unavailable("stub transport"))
381 }
382
383 async fn describe_workflow(
384 &self,
385 _: aion_proto::ProtoDescribeWorkflowRequest,
386 ) -> Result<ProtoDescribeWorkflowResponse, ClientError> {
387 Err(ClientError::unavailable("stub transport"))
388 }
389
390 async fn read_history(
391 &self,
392 _: aion_proto::ProtoReadHistoryRequest,
393 ) -> Result<aion_proto::ProtoReadHistoryResponse, ClientError> {
394 Err(ClientError::unavailable("stub transport"))
395 }
396
397 async fn subscribe(
398 &self,
399 _: aion_proto::SubscriptionRequest,
400 resume_from_sequence: Option<u64>,
401 ) -> Result<SubscriptionAttempt, ClientError> {
402 self.resume_points.lock().await.push(resume_from_sequence);
403 if let Some(failure) = self.attach_failures.lock().await.pop_front() {
404 return Err(failure);
405 }
406 self.attempts
407 .lock()
408 .await
409 .pop_front()
410 .ok_or_else(|| ClientError::server("missing subscribe attempt"))
411 }
412 }
413
414 fn event(seq: u64, workflow_id: &WorkflowId) -> Event {
415 Event::WorkflowStarted {
416 envelope: EventEnvelope {
417 seq,
418 recorded_at: Utc::now(),
419 workflow_id: workflow_id.clone(),
420 },
421 workflow_type: String::from("checkout"),
422 input: Payload::new(ContentType::Json, Vec::new()),
423 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
424 parent_run_id: None,
425 parent_workflow_id: None,
426 package_version: aion_core::PackageVersion::new("a".repeat(64)),
427 }
428 }
429
430 #[tokio::test]
431 async fn resumes_after_transient_disconnect_without_gaps_or_duplicates() {
432 let workflow_id = WorkflowId::new_v4();
433 let stub = Arc::new(SubscribeStub::default());
434 stub.attempts
435 .lock()
436 .await
437 .push_back(SubscriptionAttempt::new(
438 stream::iter(vec![
439 Ok(event(1, &workflow_id)),
440 Ok(event(2, &workflow_id)),
441 Err(ClientError::unavailable("transient disconnect")),
442 ])
443 .boxed(),
444 ));
445 stub.attempts
446 .lock()
447 .await
448 .push_back(SubscriptionAttempt::new(
449 stream::iter(vec![
450 Ok(event(2, &workflow_id)),
451 Ok(event(3, &workflow_id)),
452 Ok(event(4, &workflow_id)),
453 ])
454 .boxed(),
455 ));
456 let mut events = ResumingEventStream::new(
457 stub.clone(),
458 "tenant-a",
459 SubscribeTarget::Workflow {
460 workflow_id: workflow_id.clone(),
461 },
462 );
463
464 let mut seqs = Vec::new();
465 while let Some(item) = events.next().await {
466 let event = item
467 .map_err(|e| format!("unexpected stream error: {e}"))
468 .ok();
469 if let Some(event) = event {
470 seqs.push(event.seq());
471 }
472 }
473
474 assert_eq!(seqs, vec![1, 2, 3, 4]);
475 assert_eq!(*stub.resume_points.lock().await, vec![None, Some(3)]);
476 }
477
478 #[tokio::test]
479 async fn terminal_failure_is_yielded_before_end() {
480 let workflow_id = WorkflowId::new_v4();
481 let stub = Arc::new(SubscribeStub::default());
482 stub.attempts
483 .lock()
484 .await
485 .push_back(SubscriptionAttempt::new(
486 stream::iter(vec![Err(ClientError::unauthenticated("bad token"))]).boxed(),
487 ));
488 let mut events =
489 ResumingEventStream::new(stub, "tenant-a", SubscribeTarget::Workflow { workflow_id });
490
491 assert_eq!(
492 events.next().await,
493 Some(Err(ClientError::unauthenticated("bad token")))
494 );
495 assert_eq!(events.next().await, None);
496 }
497
498 #[tokio::test]
499 async fn namespace_denied_is_terminal_and_never_retried() {
500 let workflow_id = WorkflowId::new_v4();
501 let stub = Arc::new(SubscribeStub::default());
502 let denied =
503 ClientError::namespace_denied("namespace tenant-b is not granted to this caller");
504 stub.attempts
505 .lock()
506 .await
507 .push_back(SubscriptionAttempt::new(
508 stream::iter(vec![Err(denied.clone())]).boxed(),
509 ));
510 let mut events = ResumingEventStream::new(
511 stub.clone(),
512 "tenant-b",
513 SubscribeTarget::Workflow { workflow_id },
514 );
515
516 assert_eq!(events.next().await, Some(Err(denied)));
517 assert_eq!(events.next().await, None);
518 assert_eq!(stub.resume_points.lock().await.len(), 1);
519 }
520
521 #[tokio::test]
522 async fn from_sequence_passes_the_cursor_on_the_initial_attach() {
523 let workflow_id = WorkflowId::new_v4();
524 let stub = Arc::new(SubscribeStub::default());
525 stub.attempts
526 .lock()
527 .await
528 .push_back(SubscriptionAttempt::new(
529 stream::iter(vec![Ok(event(1, &workflow_id)), Ok(event(2, &workflow_id))]).boxed(),
530 ));
531 let Some(resume_from) = std::num::NonZeroU64::new(1) else {
532 unreachable!("1 is non-zero");
533 };
534 let mut events = super::ResumingEventStream::from_sequence(
535 stub.clone(),
536 "tenant-a",
537 workflow_id,
538 resume_from,
539 );
540
541 let mut seqs = Vec::new();
542 while let Some(item) = events.next().await {
543 if let Ok(event) = item {
544 seqs.push(event.seq());
545 }
546 }
547
548 assert_eq!(seqs, vec![1, 2]);
549 assert_eq!(
550 *stub.resume_points.lock().await,
551 vec![Some(1)],
552 "the initial attach must carry the explicit cursor"
553 );
554 }
555
556 #[tokio::test]
557 async fn live_only_streams_reconnect_only_before_any_delivery() {
558 let workflow_id = WorkflowId::new_v4();
561 let stub = Arc::new(SubscribeStub::default());
562 stub.attempts
563 .lock()
564 .await
565 .push_back(SubscriptionAttempt::new(
566 stream::iter(vec![Err(ClientError::unavailable("transient disconnect"))]).boxed(),
567 ));
568 stub.attempts
569 .lock()
570 .await
571 .push_back(SubscriptionAttempt::new(
572 stream::iter(vec![Ok(event(1, &workflow_id))]).boxed(),
573 ));
574 let mut events = ResumingEventStream::new(
575 stub.clone(),
576 "tenant-a",
577 SubscribeTarget::Filtered {
578 filter: aion_core::WorkflowFilter::default(),
579 },
580 );
581
582 let mut seqs = Vec::new();
583 while let Some(item) = events.next().await {
584 if let Ok(event) = item {
585 seqs.push(event.seq());
586 }
587 }
588
589 assert_eq!(seqs, vec![1]);
590 assert_eq!(
591 *stub.resume_points.lock().await,
592 vec![None, None],
593 "live-only streams never carry a resume cursor"
594 );
595 }
596
597 #[tokio::test]
598 async fn live_only_disconnect_after_delivery_is_honest_unavailable() {
599 for target in [
603 SubscribeTarget::Filtered {
604 filter: aion_core::WorkflowFilter::default(),
605 },
606 SubscribeTarget::Firehose,
607 ] {
608 let workflow_id = WorkflowId::new_v4();
609 let stub = Arc::new(SubscribeStub::default());
610 stub.attempts
611 .lock()
612 .await
613 .push_back(SubscriptionAttempt::new(
614 stream::iter(vec![
615 Ok(event(1, &workflow_id)),
616 Err(ClientError::unavailable("transient disconnect")),
617 ])
618 .boxed(),
619 ));
620 let mut events = ResumingEventStream::new(stub.clone(), "tenant-a", target);
621
622 let first = events.next().await;
623 assert!(matches!(first, Some(Ok(_))), "got {first:?}");
624 assert_eq!(
625 events.next().await,
626 Some(Err(ClientError::unavailable("transient disconnect")))
627 );
628 assert_eq!(events.next().await, None);
629 assert_eq!(
630 stub.resume_points.lock().await.len(),
631 1,
632 "no reattach may follow a post-delivery live-only disconnect"
633 );
634 }
635 }
636
637 #[tokio::test]
638 async fn live_only_streams_do_not_dedupe_sequence_numbers_across_workflows() {
639 let first_workflow = WorkflowId::new_v4();
642 let second_workflow = WorkflowId::new_v4();
643 let stub = Arc::new(SubscribeStub::default());
644 stub.attempts
645 .lock()
646 .await
647 .push_back(SubscriptionAttempt::new(
648 stream::iter(vec![
649 Ok(event(1, &first_workflow)),
650 Ok(event(1, &second_workflow)),
651 ])
652 .boxed(),
653 ));
654 let mut events = ResumingEventStream::new(stub, "tenant-a", SubscribeTarget::Firehose);
655
656 let mut delivered = Vec::new();
657 while let Some(item) = events.next().await {
658 if let Ok(event) = item {
659 delivered.push(event.envelope().workflow_id.clone());
660 }
661 }
662
663 assert_eq!(delivered, vec![first_workflow, second_workflow]);
664 }
665
666 #[tokio::test]
667 async fn not_found_is_terminal_and_never_retried() {
668 let workflow_id = WorkflowId::new_v4();
672 let stub = Arc::new(SubscribeStub::default());
673 stub.attempts
674 .lock()
675 .await
676 .push_back(SubscriptionAttempt::new(
677 stream::iter(vec![Err(ClientError::not_found("workflow was not found"))]).boxed(),
678 ));
679 let mut events = ResumingEventStream::new(
680 stub.clone(),
681 "tenant-a",
682 SubscribeTarget::Workflow { workflow_id },
683 );
684
685 assert_eq!(
686 events.next().await,
687 Some(Err(ClientError::not_found("workflow was not found")))
688 );
689 assert_eq!(events.next().await, None);
690 assert_eq!(stub.resume_points.lock().await.len(), 1);
691 }
692
693 #[tokio::test]
697 async fn unavailable_attach_failure_is_retried_until_attach_succeeds() -> Result<(), ClientError>
698 {
699 let workflow_id = WorkflowId::new_v4();
700 let stub = Arc::new(SubscribeStub::default());
701 stub.attach_failures
702 .lock()
703 .await
704 .push_back(ClientError::unavailable("connection refused"));
705 stub.attach_failures
706 .lock()
707 .await
708 .push_back(ClientError::unavailable("connection refused"));
709 stub.attempts
710 .lock()
711 .await
712 .push_back(SubscriptionAttempt::new(
713 stream::iter(vec![Ok(event(1, &workflow_id)), Ok(event(2, &workflow_id))]).boxed(),
714 ));
715 let mut events = ResumingEventStream::new(
716 stub.clone(),
717 "tenant-a",
718 SubscribeTarget::Workflow { workflow_id },
719 );
720
721 let mut seqs = Vec::new();
722 while let Some(item) = events.next().await {
723 seqs.push(item?.seq());
726 }
727
728 assert_eq!(seqs, vec![1, 2]);
729 assert_eq!(
730 *stub.resume_points.lock().await,
731 vec![None, None, None],
732 "every retried initial attach is still a live tail (no cursor)"
733 );
734 Ok(())
735 }
736
737 #[tokio::test]
740 async fn unavailable_reconnect_failure_retries_with_the_same_cursor() -> Result<(), ClientError>
741 {
742 let workflow_id = WorkflowId::new_v4();
743 let stub = Arc::new(SubscribeStub::default());
744 stub.attempts
745 .lock()
746 .await
747 .push_back(SubscriptionAttempt::new(
748 stream::iter(vec![
749 Ok(event(1, &workflow_id)),
750 Err(ClientError::unavailable("transient disconnect")),
751 ])
752 .boxed(),
753 ));
754 let mut events = ResumingEventStream::new(
755 stub.clone(),
756 "tenant-a",
757 SubscribeTarget::Workflow {
758 workflow_id: workflow_id.clone(),
759 },
760 );
761 let first = events.next().await;
762 assert!(matches!(first, Some(Ok(_))), "got {first:?}");
763 stub.attach_failures
765 .lock()
766 .await
767 .push_back(ClientError::unavailable("connection refused"));
768 stub.attempts
769 .lock()
770 .await
771 .push_back(SubscriptionAttempt::new(
772 stream::iter(vec![Ok(event(2, &workflow_id))]).boxed(),
773 ));
774
775 let mut seqs = vec![1];
776 while let Some(item) = events.next().await {
777 seqs.push(item?.seq());
780 }
781
782 assert_eq!(seqs, vec![1, 2]);
783 assert_eq!(
784 *stub.resume_points.lock().await,
785 vec![None, Some(2), Some(2)],
786 "the failed reconnect and the successful retry carry the same cursor"
787 );
788 Ok(())
789 }
790
791 #[tokio::test]
794 async fn non_retryable_attach_failure_is_terminal() {
795 let workflow_id = WorkflowId::new_v4();
796 let stub = Arc::new(SubscribeStub::default());
797 stub.attach_failures
798 .lock()
799 .await
800 .push_back(ClientError::unauthenticated("bad token"));
801 let mut events = ResumingEventStream::new(
802 stub.clone(),
803 "tenant-a",
804 SubscribeTarget::Workflow { workflow_id },
805 );
806
807 assert_eq!(
808 events.next().await,
809 Some(Err(ClientError::unauthenticated("bad token")))
810 );
811 assert_eq!(events.next().await, None);
812 assert_eq!(stub.resume_points.lock().await.len(), 1);
813 }
814
815 #[tokio::test]
818 async fn live_only_unavailable_attach_failure_is_retried_before_any_delivery() {
819 let workflow_id = WorkflowId::new_v4();
820 let stub = Arc::new(SubscribeStub::default());
821 stub.attach_failures
822 .lock()
823 .await
824 .push_back(ClientError::unavailable("connection refused"));
825 stub.attempts
826 .lock()
827 .await
828 .push_back(SubscriptionAttempt::new(
829 stream::iter(vec![Ok(event(1, &workflow_id))]).boxed(),
830 ));
831 let mut events =
832 ResumingEventStream::new(stub.clone(), "tenant-a", SubscribeTarget::Firehose);
833
834 let mut seqs = Vec::new();
835 while let Some(item) = events.next().await {
836 if let Ok(event) = item {
837 seqs.push(event.seq());
838 }
839 }
840
841 assert_eq!(seqs, vec![1]);
842 assert_eq!(*stub.resume_points.lock().await, vec![None, None]);
843 }
844}