Skip to main content

laminar_db/subscription/portal/
mod.rs

1//! Per-subscriber cursor over the shared subscription log.
2
3use std::sync::Arc;
4
5use arrow_array::RecordBatch;
6use arrow_schema::SchemaRef;
7use datafusion::physical_expr::PhysicalExpr;
8use futures::FutureExt;
9use laminar_core::checkpoint::{OutputPartitionId, PartitionSequence, StreamGeneration};
10
11#[cfg(feature = "cluster")]
12use super::cluster::{ClusterReaderFrame, ClusterReaderRead, ClusterSubscriptionReader};
13use super::registry::{ChargedUpdate, MvUpdate, SubscriptionRead, SubscriptionReader};
14
15#[derive(Debug)]
16enum PortalReader {
17    Local(SubscriptionReader),
18    #[cfg(feature = "cluster")]
19    Cluster(ClusterSubscriptionReader),
20}
21
22impl PortalReader {
23    async fn next(&mut self) -> PortalRead {
24        match self {
25            Self::Local(reader) => PortalRead::Local(reader.next().await),
26            #[cfg(feature = "cluster")]
27            Self::Cluster(reader) => PortalRead::Cluster(reader.next().await),
28        }
29    }
30
31    fn try_next(&mut self) -> Option<PortalRead> {
32        match self {
33            Self::Local(reader) => reader.next().now_or_never().map(PortalRead::Local),
34            #[cfg(feature = "cluster")]
35            Self::Cluster(reader) => reader.try_next().map(PortalRead::Cluster),
36        }
37    }
38}
39
40enum PortalRead {
41    Local(SubscriptionRead),
42    #[cfg(feature = "cluster")]
43    Cluster(ClusterReaderRead),
44}
45
46/// Keeps the process-wide subscription charge alive with an emitted batch.
47#[doc(hidden)]
48#[derive(Clone)]
49pub struct SubscriptionFrameLease {
50    _local_owner: Option<ChargedUpdate>,
51    #[cfg(feature = "cluster")]
52    _cluster_owner: Option<Arc<tokio::sync::OwnedSemaphorePermit>>,
53}
54
55impl std::fmt::Debug for SubscriptionFrameLease {
56    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
57        formatter.write_str("SubscriptionFrameLease")
58    }
59}
60
61/// Optional durable identity carried alongside a compatibility [`PortalFrame`].
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum ClusterSubscriptionFrameMetadata {
64    /// Partition-local identity of one committed data frame.
65    Data {
66        /// Durable stream incarnation.
67        stream_generation: StreamGeneration,
68        /// Stable vnode output partition.
69        partition: OutputPartitionId,
70        /// Monotonic sequence within this generation and partition.
71        partition_sequence: PartitionSequence,
72        /// Whole-cluster checkpoint that first exposed this frame.
73        committed_epoch: u64,
74    },
75    /// Whole-cluster committed progress after every partition interval was delivered.
76    Progress {
77        /// Durable stream incarnation.
78        stream_generation: StreamGeneration,
79        /// Committed checkpoint epoch.
80        epoch: u64,
81        /// Committed checkpoint identifier.
82        checkpoint_id: u64,
83    },
84}
85
86/// Backward-compatible frame plus optional cluster delivery metadata.
87#[derive(Debug, Clone)]
88pub struct SubscriptionEnvelope {
89    /// Existing local/standalone-compatible frame.
90    pub frame: PortalFrame,
91    /// Present only for committed cluster frames.
92    pub cluster: Option<ClusterSubscriptionFrameMetadata>,
93    /// Stable terminal error code when `frame` is [`PortalFrame::Error`].
94    pub error_code: Option<&'static str>,
95}
96
97impl SubscriptionEnvelope {
98    fn local(frame: PortalFrame) -> Self {
99        Self {
100            frame,
101            cluster: None,
102            error_code: None,
103        }
104    }
105}
106
107/// One frame emitted toward the wire.
108#[derive(Debug, Clone)]
109pub enum PortalFrame {
110    /// Rows produced in a cycle.
111    Batch {
112        /// Arrow rows in the shared-log entry.
113        batch: RecordBatch,
114        /// Portal-local delivery sequence. In cluster mode this is neither durable nor global;
115        /// use [`ClusterSubscriptionFrameMetadata::Data`] for partition identity.
116        sequence: u64,
117        /// Internal process-memory ownership token.
118        #[doc(hidden)]
119        lease: SubscriptionFrameLease,
120    },
121    /// Progress frontier for a durably committed checkpoint.
122    Barrier {
123        /// Portal-local delivery sequence; not a cluster-wide ordering position.
124        sequence: u64,
125        /// Engine checkpoint epoch.
126        epoch: u64,
127        /// Engine checkpoint id.
128        checkpoint_id: u64,
129        /// Local-log cut, or the gateway-local delivery cut in cluster mode. The durable cluster
130        /// position is the checkpoint's partition-frontier vector, addressed by `epoch`.
131        through_sequence: u64,
132    },
133    /// Consumer fell behind by exactly `skipped` shared-log entries. This is
134    /// terminal because continuing would hide rows or checkpoint markers.
135    Lagged(u64),
136    /// The subscription cannot continue without returning an invalid result.
137    Error {
138        /// Human-readable failure detail.
139        message: String,
140    },
141}
142
143/// One `SUBSCRIBE` consumer.
144#[derive(Debug)]
145pub struct SubscriptionPortal {
146    name: String,
147    schema: SchemaRef,
148    reader: Option<PortalReader>,
149    closed: bool,
150    filter: Option<Arc<dyn PhysicalExpr>>,
151}
152
153impl SubscriptionPortal {
154    pub(crate) fn open(
155        name: impl Into<String>,
156        schema: SchemaRef,
157        reader: SubscriptionReader,
158    ) -> Self {
159        Self::new_inner(name, schema, reader, None)
160    }
161
162    pub(crate) fn open_with_filter(
163        name: impl Into<String>,
164        schema: SchemaRef,
165        reader: SubscriptionReader,
166        filter: Arc<dyn PhysicalExpr>,
167    ) -> Self {
168        Self::new_inner(name, schema, reader, Some(filter))
169    }
170
171    fn new_inner(
172        name: impl Into<String>,
173        schema: SchemaRef,
174        reader: SubscriptionReader,
175        filter: Option<Arc<dyn PhysicalExpr>>,
176    ) -> Self {
177        Self {
178            name: name.into(),
179            schema,
180            reader: Some(PortalReader::Local(reader)),
181            closed: false,
182            filter,
183        }
184    }
185
186    #[cfg(feature = "cluster")]
187    pub(crate) fn open_cluster(
188        name: impl Into<String>,
189        schema: SchemaRef,
190        reader: ClusterSubscriptionReader,
191        filter: Option<Arc<dyn PhysicalExpr>>,
192    ) -> Self {
193        Self {
194            name: name.into(),
195            schema,
196            reader: Some(PortalReader::Cluster(reader)),
197            closed: false,
198            filter,
199        }
200    }
201
202    /// Schema of the subscribed object.
203    #[must_use]
204    pub fn schema(&self) -> SchemaRef {
205        Arc::clone(&self.schema)
206    }
207
208    /// Next frame, or `None` after a terminal frame or explicit close.
209    pub async fn next_frame(&mut self) -> Option<PortalFrame> {
210        self.next_envelope().await.map(|envelope| envelope.frame)
211    }
212
213    /// Next frame with optional partition/checkpoint metadata for cluster-aware clients.
214    pub async fn next_envelope(&mut self) -> Option<SubscriptionEnvelope> {
215        if self.closed {
216            return None;
217        }
218
219        loop {
220            let read = self.reader.as_mut()?.next().await;
221            if let Some(frame) = self.process_read(read) {
222                return Some(frame);
223            }
224        }
225    }
226
227    /// Return the next immediately available frame without waiting.
228    pub fn try_next_frame(&mut self) -> Option<PortalFrame> {
229        self.try_next_envelope().map(|envelope| envelope.frame)
230    }
231
232    /// Return the next immediately available cluster-aware envelope without waiting.
233    pub fn try_next_envelope(&mut self) -> Option<SubscriptionEnvelope> {
234        if self.closed {
235            return None;
236        }
237
238        loop {
239            let read = self.reader.as_mut()?.try_next()?;
240            if let Some(frame) = self.process_read(read) {
241                return Some(frame);
242            }
243        }
244    }
245
246    fn process_read(&mut self, read: PortalRead) -> Option<SubscriptionEnvelope> {
247        match read {
248            PortalRead::Local(read) => self.process_local_read(read),
249            #[cfg(feature = "cluster")]
250            PortalRead::Cluster(read) => self.process_cluster_read(read),
251        }
252    }
253
254    fn process_local_read(&mut self, read: SubscriptionRead) -> Option<SubscriptionEnvelope> {
255        let frame = match read {
256            SubscriptionRead::Update { sequence, update } => translate(sequence, update),
257            SubscriptionRead::Lagged(skipped) => {
258                tracing::warn!(
259                    subscription = %self.name,
260                    skipped,
261                    "subscription cursor was evicted; closing"
262                );
263                self.close();
264                return Some(SubscriptionEnvelope::local(PortalFrame::Lagged(skipped)));
265            }
266            SubscriptionRead::Terminal(message) => {
267                tracing::warn!(
268                    subscription = %self.name,
269                    %message,
270                    "subscription log terminated; closing"
271                );
272                self.close();
273                return Some(SubscriptionEnvelope::local(PortalFrame::Error { message }));
274            }
275        };
276
277        self.process_envelope(SubscriptionEnvelope::local(frame))
278    }
279
280    #[cfg(feature = "cluster")]
281    fn process_cluster_read(&mut self, read: ClusterReaderRead) -> Option<SubscriptionEnvelope> {
282        let envelope = match read {
283            ClusterReaderRead::Frame(ClusterReaderFrame::Batch {
284                batch,
285                delivery_sequence,
286                stream_generation,
287                partition,
288                partition_sequence,
289                committed_epoch,
290                permit,
291            }) => SubscriptionEnvelope {
292                frame: PortalFrame::Batch {
293                    batch,
294                    sequence: delivery_sequence,
295                    lease: SubscriptionFrameLease {
296                        _local_owner: None,
297                        _cluster_owner: Some(permit),
298                    },
299                },
300                cluster: Some(ClusterSubscriptionFrameMetadata::Data {
301                    stream_generation,
302                    partition,
303                    partition_sequence,
304                    committed_epoch,
305                }),
306                error_code: None,
307            },
308            ClusterReaderRead::Frame(ClusterReaderFrame::Progress {
309                delivery_sequence,
310                through_sequence,
311                stream_generation,
312                epoch,
313                checkpoint_id,
314            }) => SubscriptionEnvelope {
315                frame: PortalFrame::Barrier {
316                    sequence: delivery_sequence,
317                    epoch,
318                    checkpoint_id,
319                    through_sequence,
320                },
321                cluster: Some(ClusterSubscriptionFrameMetadata::Progress {
322                    stream_generation,
323                    epoch,
324                    checkpoint_id,
325                }),
326                error_code: None,
327            },
328            ClusterReaderRead::Terminal(error) => {
329                let code = error.code();
330                tracing::warn!(
331                    subscription = %self.name,
332                    code,
333                    error = %error,
334                    "committed cluster subscription terminated"
335                );
336                self.close();
337                return Some(SubscriptionEnvelope {
338                    frame: PortalFrame::Error {
339                        message: format!("[{code}] {error}"),
340                    },
341                    cluster: None,
342                    error_code: Some(code),
343                });
344            }
345        };
346        self.process_envelope(envelope)
347    }
348
349    fn process_envelope(&mut self, envelope: SubscriptionEnvelope) -> Option<SubscriptionEnvelope> {
350        let SubscriptionEnvelope {
351            frame,
352            cluster,
353            error_code,
354        } = envelope;
355
356        let PortalFrame::Batch {
357            batch,
358            sequence,
359            lease,
360        } = frame
361        else {
362            if matches!(&frame, PortalFrame::Error { .. }) {
363                self.close();
364            }
365            return Some(SubscriptionEnvelope {
366                frame,
367                cluster,
368                error_code,
369            });
370        };
371        let Some(filter) = self.filter.as_ref() else {
372            return Some(SubscriptionEnvelope {
373                frame: PortalFrame::Batch {
374                    batch,
375                    sequence,
376                    lease,
377                },
378                cluster,
379                error_code,
380            });
381        };
382        match crate::filter_compile::apply(&batch, filter.as_ref()) {
383            Ok(Some(filtered)) => Some(SubscriptionEnvelope {
384                frame: PortalFrame::Batch {
385                    batch: filtered,
386                    sequence,
387                    lease,
388                },
389                cluster,
390                error_code,
391            }),
392            Ok(None) => None,
393            Err(error) => {
394                tracing::warn!(
395                    subscription = %self.name,
396                    %error,
397                    "subscription filter failed; closing"
398                );
399                let message = error.to_string();
400                self.close();
401                Some(SubscriptionEnvelope {
402                    frame: PortalFrame::Error { message },
403                    cluster: None,
404                    error_code: None,
405                })
406            }
407        }
408    }
409
410    /// Stop reading and release the subscriber registration. Idempotent.
411    pub fn close(&mut self) {
412        self.closed = true;
413        self.reader = None;
414    }
415
416    /// True after `close()` has been called.
417    #[must_use]
418    pub fn is_closed(&self) -> bool {
419        self.closed
420    }
421}
422
423impl Drop for SubscriptionPortal {
424    fn drop(&mut self) {
425        self.close();
426    }
427}
428
429fn translate(sequence: u64, update: ChargedUpdate) -> PortalFrame {
430    match update.as_ref() {
431        MvUpdate::Batch(batch) => {
432            let batch = batch.clone();
433            PortalFrame::Batch {
434                batch,
435                sequence,
436                lease: SubscriptionFrameLease {
437                    _local_owner: Some(update),
438                    #[cfg(feature = "cluster")]
439                    _cluster_owner: None,
440                },
441            }
442        }
443        MvUpdate::Barrier {
444            epoch,
445            checkpoint_id,
446            through_sequence,
447        } => PortalFrame::Barrier {
448            sequence,
449            epoch: *epoch,
450            checkpoint_id: *checkpoint_id,
451            through_sequence: *through_sequence,
452        },
453        MvUpdate::Error(message) => PortalFrame::Error {
454            message: message.clone(),
455        },
456    }
457}
458
459#[cfg(test)]
460mod tests;