Skip to main content

aion_client/
stream.rs

1//! Event subscription `Stream` and resumption.
2
3use 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
20/// Boxed event stream returned by subscribe operations.
21pub type EventStream = Pin<Box<dyn Stream<Item = Result<Event, ClientError>> + Send>>;
22
23/// Builder for the AW-owned subscription variants supported by the SDK.
24#[derive(Clone, Debug, PartialEq, Eq)]
25pub enum SubscribeTarget {
26    /// Subscribe to events for one workflow.
27    Workflow {
28        /// Workflow whose events are requested.
29        workflow_id: WorkflowId,
30    },
31    /// Subscribe to workflow metadata selected events.
32    Filtered {
33        /// Workflow metadata filter used for the subscription.
34        filter: WorkflowFilter,
35    },
36    /// Subscribe to all visible events in the client's namespace.
37    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
75/// Reconnecting subscription stream.
76///
77/// Resumption is per-workflow only: per-workflow `seq` is the only ordering
78/// that exists, so only [`SubscribeTarget::Workflow`] streams track a cursor
79/// (`resume_from_seq = last delivered + 1`) and deduplicate by sequence
80/// number. Filtered and firehose streams are live-only by design: a
81/// transient disconnect after at least one delivered event ends the stream
82/// with an honest [`ClientError::Unavailable`] instead of silently
83/// reattaching a gapped stream; reconnect-live-only is allowed only while
84/// nothing has been delivered yet.
85///
86/// Connect-failure contract (cross-SDK): a failed subscription attach is
87/// classified exactly like a mid-stream drop. [`ClientError::Unavailable`]
88/// (transport-level connect failure, DNS/TLS/socket failure, abnormal close)
89/// is retryable and the stream re-attaches — on the initial attach as well as
90/// after delivered events — until the caller drops the stream; every other
91/// taxonomy error (`Unauthenticated`, `NamespaceDenied`, `NotFound`,
92/// `InvalidArgument`, `Server`, ...) is terminal immediately.
93pub 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    /// Creates a subscription stream for `target`.
107    #[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    /// Creates a per-workflow subscription stream that attaches with an
127    /// explicit starting cursor.
128    ///
129    /// `resume_from` is the first per-workflow sequence number wanted
130    /// (`resume_from_seq` on the wire); `1` replays the full recorded
131    /// history before splicing into the live stream. The type makes the
132    /// invalid cursor `0` unrepresentable.
133    #[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        // The cursor sent on (re)attach is always `last_seq + 1`, so seeding
146        // `last_seq = resume_from - 1` makes the first attach request exactly
147        // `resume_from` and drops anything older on the dedupe path.
148        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        // Only per-workflow streams carry a resume cursor; filtered and
160        // firehose reattach live-only (and only before any delivery).
161        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                        // Cross-SDK connect-failure contract: an attach
200                        // failure is classified exactly like a mid-stream
201                        // drop. `Unavailable` is retryable — per-workflow
202                        // streams reconnect with their cursor, live-only
203                        // streams reconnect only while nothing has been
204                        // delivered. Every other taxonomy error is terminal.
205                        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                        // Sequence-number dedupe is coherent only within one
223                        // workflow's history.
224                        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                            // Nothing delivered yet: a live-only reattach
240                            // cannot gap, so reconnect.
241                            continue;
242                        }
243                        // Filtered/firehose streams have no resume cursor; a
244                        // reattach after delivered events would silently gap.
245                        // Surface an honest terminal Unavailable instead.
246                    }
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/// Boxes a resuming event stream behind the public return type.
260#[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/// Boxes a per-workflow stream attaching with an explicit starting cursor.
270#[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 consumed before any queued attempt: each entry is
312        /// one subscribe call that fails before a stream exists.
313        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        // A filtered stream that drops before delivering anything may
559        // reattach live-only — nothing can gap yet — and never with a cursor.
560        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        // Filtered/firehose streams have no resume cursor: a transient drop
600        // after >= 1 delivered event must surface Unavailable, never a silent
601        // gapped reattach.
602        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        // Per-workflow seq is the only ordering that exists; two workflows
640        // legitimately share sequence numbers on a filtered/firehose stream.
641        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        // A workflow-level visibility miss surfaces as NotFound (the server's
669        // anti-existence-leak contract); like every non-Unavailable error it
670        // must end the stream instead of reconnecting forever.
671        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    /// Connect-failure contract: an `Unavailable` initial attach failure is
694    /// retryable exactly like a mid-stream drop — the stream re-attaches and
695    /// delivers, never surfacing the transient error as terminal.
696    #[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            // A transient attach failure must not surface; `?` fails the test
724            // with the offending error if it does.
725            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    /// A mid-stream drop followed by an `Unavailable` reconnect failure keeps
738    /// retrying with the SAME cursor until the reconnect succeeds.
739    #[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        // The reconnect attempt fails transiently, then succeeds.
764        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            // A transient reconnect failure must not surface; `?` fails the
778            // test with the offending error if it does.
779            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    /// Non-`Unavailable` attach failures are terminal immediately: an
792    /// `Unauthenticated` connect rejection must never be retried.
793    #[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    /// Live-only streams also retry `Unavailable` attach failures while
816    /// nothing has been delivered (a live-only reattach cannot gap yet).
817    #[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}