Repository navigation
Expand file tree
/
Copy pathsubscription.rs
More file actions
388 lines (348 loc) · 13.1 KB
/
Copy pathsubscription.rs
File metadata and controls
388 lines (348 loc) · 13.1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
//! Subscription handles for observing session and lifecycle events.
//!
//! Returned by [`Session::subscribe`](crate::session::Session::subscribe) and
//! [`Client::subscribe_lifecycle`](crate::Client::subscribe_lifecycle).
//!
//! Each subscription is an opt-in **observer** of events that are also
//! delivered to the per-event handlers installed on the session config
//! (see [`crate::handler`]). Subscribers receive a clone of every event but
//! cannot influence permission decisions, tool results, or any other event
//! whose handler return value affects the runtime.
//!
//! # Async iteration
//!
//! The subscription types implement [`tokio_stream::Stream`], so consumers
//! can use adapter combinators from [`tokio_stream::StreamExt`] or
//! `futures::StreamExt` (filtering, mapping, batching, racing with
//! `tokio::select!`, etc.) without learning the SDK's internal channel
//! choice. A simple `while let Ok(event) = sub.recv().await { ... }` loop
//! also works for callers who don't need the [`Stream`](tokio_stream::Stream)
//! surface.
//!
//! # Resume bootstrap and lag policy
//!
//! The first subscription on a resumed session may begin with a lossless,
//! ordered bootstrap prefix. Once its owner catches up, delivery switches
//! atomically to the bounded live stream. See
//! [`Session::subscribe`](crate::session::Session::subscribe) for the
//! unbounded retention, eager ownership, cleanup, and router limits.
//!
//! Each live subscriber maintains its own finite queue. If a consumer cannot
//! keep up, the oldest live events are dropped and the next call yields
//! [`Lagged`](crate::subscription::Lagged) reporting how many events were skipped.
//! Slow subscribers do not block the producer.
use std::collections::VecDeque;
use std::fmt;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll};
use parking_lot::Mutex;
use tokio::sync::broadcast::{Receiver, Sender, WeakSender};
use tokio_stream::wrappers::BroadcastStream;
use tokio_stream::wrappers::errors::BroadcastStreamRecvError;
use tokio_stream::{Stream, StreamExt as _};
use crate::types::{SessionEvent, SessionLifecycleEvent};
use crate::{Custom, Repr};
/// The subscription fell behind the producer.
///
/// Reports the number of events that were dropped from this subscriber's
/// queue because the consumer didn't keep up. The subscription continues
/// after this error, starting from the next live event — callers who care
/// about lag should match on it and decide whether to resync, re-fetch, or
/// log and continue.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Lagged(pub(crate) u64);
impl Lagged {
/// Number of events skipped before this consumer could read them.
pub fn skipped(&self) -> u64 {
self.0
}
}
impl fmt::Display for Lagged {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "subscription lagged behind by {} events", self.0)
}
}
impl std::error::Error for Lagged {}
/// Error kind for subscription receive operations.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum RecvErrorKind {
/// The producer is gone — the session has shut down or the client has
/// stopped. No further events will be delivered.
Closed,
/// The subscriber fell behind. See [`Lagged`].
Lagged(Lagged),
}
impl fmt::Display for RecvErrorKind {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
RecvErrorKind::Closed => write!(f, "subscription closed"),
RecvErrorKind::Lagged(l) => write!(f, "{l}"),
}
}
}
/// Error returned by [`crate::subscription::EventSubscription::recv`] and
/// [`crate::subscription::LifecycleSubscription::recv`].
#[derive(Debug)]
pub struct RecvError {
repr: Repr<RecvErrorKind>,
}
impl RecvError {
/// The [`RecvErrorKind`] of this error.
pub fn kind(&self) -> &RecvErrorKind {
match &self.repr {
Repr::Simple(k) | Repr::SimpleMessage(k, ..) | Repr::Custom(Custom { kind: k, .. }) => {
k
}
}
}
}
impl fmt::Display for RecvError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match &self.repr {
Repr::Simple(k) => write!(f, "{k}"),
Repr::SimpleMessage(_, m) => write!(f, "{m}"),
Repr::Custom(Custom { error, .. }) => write!(f, "{error}"),
}
}
}
impl std::error::Error for RecvError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match &self.repr {
Repr::Custom(Custom { error, .. }) => Some(&**error),
_ => None,
}
}
}
impl From<RecvErrorKind> for RecvError {
fn from(kind: RecvErrorKind) -> Self {
Self {
repr: Repr::Simple(kind),
}
}
}
impl From<Lagged> for RecvError {
fn from(lagged: Lagged) -> Self {
Self::from(RecvErrorKind::Lagged(lagged))
}
}
enum ResumeBootstrapState {
Unclaimed(VecDeque<SessionEvent>),
Claimed(VecDeque<SessionEvent>),
Disabled,
}
/// Publication, ownership, and the empty-queue handoff share one lock so
/// the bootstrap owner observes an exact prefix without gaps or duplicates.
pub(crate) struct ResumeBootstrap {
state: Mutex<ResumeBootstrapState>,
live: WeakSender<SessionEvent>,
}
/// Releases only unclaimed events when the publishing task exits or unwinds.
pub(crate) struct ResumeBootstrapCleanup(Arc<ResumeBootstrap>);
impl Drop for ResumeBootstrapCleanup {
fn drop(&mut self) {
self.0.release_unclaimed();
}
}
impl ResumeBootstrap {
pub(crate) fn new(event_tx: &Sender<SessionEvent>) -> Arc<Self> {
Arc::new(Self {
state: Mutex::new(ResumeBootstrapState::Unclaimed(VecDeque::new())),
live: event_tx.downgrade(),
})
}
pub(crate) fn cleanup_guard(self: &Arc<Self>) -> ResumeBootstrapCleanup {
ResumeBootstrapCleanup(self.clone())
}
pub(crate) fn publish(&self, event_tx: &Sender<SessionEvent>, event: SessionEvent) {
let mut state = self.state.lock();
match &mut *state {
ResumeBootstrapState::Unclaimed(events) | ResumeBootstrapState::Claimed(events) => {
events.push_back(event.clone());
}
ResumeBootstrapState::Disabled => {}
}
// Other observers stay live while the bootstrap owner catches up.
let _ = event_tx.send(event);
}
pub(crate) fn subscribe(
self: &Arc<Self>,
event_tx: &Sender<SessionEvent>,
) -> EventSubscription {
let mut state = self.state.lock();
match &mut *state {
ResumeBootstrapState::Unclaimed(events) => {
let events = std::mem::take(events);
*state = ResumeBootstrapState::Claimed(events);
EventSubscription {
inner: None,
bootstrap: Some(self.clone()),
}
}
ResumeBootstrapState::Claimed(_) | ResumeBootstrapState::Disabled => {
EventSubscription::new(event_tx.subscribe())
}
}
}
fn pop(&self, live: &mut Option<BroadcastStream<SessionEvent>>) -> Option<SessionEvent> {
let mut state = self.state.lock();
let ResumeBootstrapState::Claimed(events) = &mut *state else {
return None;
};
if let Some(event) = events.pop_front() {
return Some(event);
}
// Subscribe under the publication lock, never replaying the broadcast
// copy of an event already delivered from the bootstrap queue.
*live = self
.live
.upgrade()
.map(|sender| BroadcastStream::new(sender.subscribe()));
*state = ResumeBootstrapState::Disabled;
None
}
pub(crate) fn release_unclaimed(&self) {
let mut state = self.state.lock();
if matches!(*state, ResumeBootstrapState::Unclaimed(_)) {
*state = ResumeBootstrapState::Disabled;
}
}
fn abandon(&self) {
let mut state = self.state.lock();
if matches!(*state, ResumeBootstrapState::Claimed(_)) {
*state = ResumeBootstrapState::Disabled;
}
}
}
/// Subscription to runtime events for a single
/// [`Session`](crate::session::Session).
///
/// Created by [`Session::subscribe`](crate::session::Session::subscribe).
/// Implements [`Stream`] yielding `Result<SessionEvent, Lagged>`.
/// Drop the value to unsubscribe; there is no separate cancel handle.
/// A resume bootstrap is claimed at construction, not on the first poll.
/// Dropping its owner discards any unread bootstrap events.
#[must_use = "dropping the subscription unsubscribes and discards any owned resume bootstrap backlog"]
pub struct EventSubscription {
inner: Option<BroadcastStream<SessionEvent>>,
bootstrap: Option<Arc<ResumeBootstrap>>,
}
impl EventSubscription {
pub(crate) fn new(rx: Receiver<SessionEvent>) -> Self {
Self {
inner: Some(BroadcastStream::new(rx)),
bootstrap: None,
}
}
fn next_bootstrap_event(&mut self) -> Option<SessionEvent> {
let event = self
.bootstrap
.as_ref()
.and_then(|bootstrap| bootstrap.pop(&mut self.inner));
if event.is_none() {
self.bootstrap = None;
}
event
}
/// Receive the next event.
///
/// Returns:
///
/// - `Ok(event)` for the next delivered event.
/// - [`RecvErrorKind::Lagged`] if live delivery fell behind; call again
/// to continue from the next available live event.
/// - [`RecvErrorKind::Closed`] once the producer is gone and any retained
/// events have been drained.
///
/// # Cancel safety
///
/// **Cancel-safe.** Bootstrap events are removed before the future's
/// first suspension point. Once live delivery begins, this wraps a
/// cancel-safe `tokio::sync::broadcast::Receiver` via `BroadcastStream`.
pub async fn recv(&mut self) -> Result<SessionEvent, RecvError> {
match self.next().await {
Some(Ok(event)) => Ok(event),
Some(Err(lagged)) => Err(lagged.into()),
None => Err(RecvErrorKind::Closed.into()),
}
}
}
impl Stream for EventSubscription {
type Item = Result<SessionEvent, Lagged>;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
if let Some(event) = self.next_bootstrap_event() {
return Poll::Ready(Some(Ok(event)));
}
let Some(inner) = self.inner.as_mut() else {
return Poll::Ready(None);
};
match Pin::new(inner).poll_next(cx) {
Poll::Ready(Some(Ok(event))) => Poll::Ready(Some(Ok(event))),
Poll::Ready(Some(Err(BroadcastStreamRecvError::Lagged(n)))) => {
Poll::Ready(Some(Err(Lagged(n))))
}
Poll::Ready(None) => Poll::Ready(None),
Poll::Pending => Poll::Pending,
}
}
}
impl Drop for EventSubscription {
fn drop(&mut self) {
if let Some(bootstrap) = &self.bootstrap {
bootstrap.abandon();
}
}
}
/// Subscription to lifecycle events on a [`Client`](crate::Client).
///
/// Created by [`Client::subscribe_lifecycle`](crate::Client::subscribe_lifecycle).
/// Implements [`Stream`] yielding `Result<SessionLifecycleEvent, Lagged>`.
/// Drop the value to unsubscribe; there is no separate cancel handle.
#[must_use = "dropping the subscription unsubscribes"]
pub struct LifecycleSubscription {
inner: BroadcastStream<SessionLifecycleEvent>,
}
impl LifecycleSubscription {
pub(crate) fn new(rx: Receiver<SessionLifecycleEvent>) -> Self {
Self {
inner: BroadcastStream::new(rx),
}
}
/// Receive the next event.
///
/// Returns:
///
/// - `Ok(event)` for the next delivered event.
/// - [`RecvErrorKind::Lagged`] if the subscriber fell behind; call again
/// to continue from the next available event.
/// - [`RecvErrorKind::Closed`] once the producer is gone.
///
/// # Cancel safety
///
/// **Cancel-safe.** Wraps a `tokio::sync::broadcast::Receiver` via
/// `BroadcastStream`; both are cancel-safe by design. Dropping the future
/// before completion leaves buffered events available for the next call.
pub async fn recv(&mut self) -> Result<SessionLifecycleEvent, RecvError> {
match self.next().await {
Some(Ok(event)) => Ok(event),
Some(Err(lagged)) => Err(lagged.into()),
None => Err(RecvErrorKind::Closed.into()),
}
}
}
impl Stream for LifecycleSubscription {
type Item = Result<SessionLifecycleEvent, Lagged>;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
match Pin::new(&mut self.inner).poll_next(cx) {
Poll::Ready(Some(Ok(event))) => Poll::Ready(Some(Ok(event))),
Poll::Ready(Some(Err(BroadcastStreamRecvError::Lagged(n)))) => {
Poll::Ready(Some(Err(Lagged(n))))
}
Poll::Ready(None) => Poll::Ready(None),
Poll::Pending => Poll::Pending,
}
}
}
#[cfg(test)]
mod tests;