1use 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#[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#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum ClusterSubscriptionFrameMetadata {
64 Data {
66 stream_generation: StreamGeneration,
68 partition: OutputPartitionId,
70 partition_sequence: PartitionSequence,
72 committed_epoch: u64,
74 },
75 Progress {
77 stream_generation: StreamGeneration,
79 epoch: u64,
81 checkpoint_id: u64,
83 },
84}
85
86#[derive(Debug, Clone)]
88pub struct SubscriptionEnvelope {
89 pub frame: PortalFrame,
91 pub cluster: Option<ClusterSubscriptionFrameMetadata>,
93 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#[derive(Debug, Clone)]
109pub enum PortalFrame {
110 Batch {
112 batch: RecordBatch,
114 sequence: u64,
117 #[doc(hidden)]
119 lease: SubscriptionFrameLease,
120 },
121 Barrier {
123 sequence: u64,
125 epoch: u64,
127 checkpoint_id: u64,
129 through_sequence: u64,
132 },
133 Lagged(u64),
136 Error {
138 message: String,
140 },
141}
142
143#[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 #[must_use]
204 pub fn schema(&self) -> SchemaRef {
205 Arc::clone(&self.schema)
206 }
207
208 pub async fn next_frame(&mut self) -> Option<PortalFrame> {
210 self.next_envelope().await.map(|envelope| envelope.frame)
211 }
212
213 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 pub fn try_next_frame(&mut self) -> Option<PortalFrame> {
229 self.try_next_envelope().map(|envelope| envelope.frame)
230 }
231
232 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 pub fn close(&mut self) {
412 self.closed = true;
413 self.reader = None;
414 }
415
416 #[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;