kumo_jsonl/
tailer.rs

1use crate::batch::LogBatch;
2use crate::checkpoint::CheckpointData;
3use crate::decompress::{FileDecompressor, NextLine, DEFAULT_MAX_LINE_SIZE};
4use camino::Utf8PathBuf;
5use filenamegen::Glob;
6use futures::Stream;
7use notify::event::{CreateKind, ModifyKind};
8use notify::{Event, EventKind, Watcher};
9use serde::{Deserialize, Serialize};
10use std::collections::HashSet;
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::Arc;
13use std::time::Duration;
14use tokio::sync::Notify;
15use tracing::warn;
16
17// ---------------------------------------------------------------------------
18// Default helpers
19// ---------------------------------------------------------------------------
20
21fn default_pattern() -> String {
22    "*".to_string()
23}
24
25fn default_max_batch_size() -> usize {
26    100
27}
28
29fn default_max_batch_latency() -> Duration {
30    Duration::from_secs(1)
31}
32
33fn default_max_line_size() -> usize {
34    DEFAULT_MAX_LINE_SIZE
35}
36
37// ---------------------------------------------------------------------------
38// ConsumerConfig
39// ---------------------------------------------------------------------------
40
41/// Per-consumer batching, checkpoint, and filter configuration.
42pub struct ConsumerConfig {
43    /// A name that identifies this consumer.  Returned by
44    /// [`LogBatch::consumer_name`].
45    pub name: String,
46    /// Maximum number of records per batch.
47    pub max_batch_size: usize,
48    /// Maximum time to wait for a partial batch to fill before yielding it.
49    pub max_batch_latency: Duration,
50    /// If set, enables checkpoint persistence with this name.
51    /// The checkpoint file will be stored as `.<name>` in the log directory.
52    pub checkpoint_name: Option<String>,
53    /// Optional filter applied to each record.  If the filter returns
54    /// `Ok(false)` the record is not added to this consumer's batch.
55    pub filter: Option<Box<dyn Fn(&serde_json::Value) -> anyhow::Result<bool> + Send>>,
56}
57
58impl ConsumerConfig {
59    pub fn new(name: impl Into<String>) -> Self {
60        Self {
61            name: name.into(),
62            max_batch_size: default_max_batch_size(),
63            max_batch_latency: default_max_batch_latency(),
64            checkpoint_name: None,
65            filter: None,
66        }
67    }
68
69    pub fn max_batch_size(mut self, size: usize) -> Self {
70        self.max_batch_size = size;
71        self
72    }
73
74    pub fn max_batch_latency(mut self, latency: Duration) -> Self {
75        self.max_batch_latency = latency;
76        self
77    }
78
79    pub fn checkpoint_name(mut self, name: impl Into<String>) -> Self {
80        self.checkpoint_name = Some(name.into());
81        self
82    }
83
84    pub fn filter<F>(mut self, f: F) -> Self
85    where
86        F: Fn(&serde_json::Value) -> anyhow::Result<bool> + Send + 'static,
87    {
88        self.filter = Some(Box::new(f));
89        self
90    }
91}
92
93// ---------------------------------------------------------------------------
94// MultiConsumerTailerConfig
95// ---------------------------------------------------------------------------
96
97/// Configuration for a tailer that fans out records to multiple consumers.
98pub struct MultiConsumerTailerConfig {
99    /// The directory containing zstd-compressed JSONL log files.
100    pub directory: Utf8PathBuf,
101    /// Glob pattern for matching log filenames.
102    pub pattern: String,
103    /// If set, use a polling-based filesystem watcher.
104    pub poll_watcher: Option<Duration>,
105    /// If true, ignore checkpoints and start from the most recent segment.
106    pub tail: bool,
107    /// Largest a decompressed record may be before the segment is treated as
108    /// corrupt. Bounds the transient memory used to decompress an unusually
109    /// large record.
110    pub max_line_size: usize,
111    /// The set of consumers that receive records.
112    pub consumers: Vec<ConsumerConfig>,
113}
114
115impl MultiConsumerTailerConfig {
116    pub fn new(directory: Utf8PathBuf, consumers: Vec<ConsumerConfig>) -> Self {
117        Self {
118            directory,
119            pattern: default_pattern(),
120            poll_watcher: None,
121            tail: false,
122            max_line_size: default_max_line_size(),
123            consumers,
124        }
125    }
126
127    pub fn pattern(mut self, pattern: impl Into<String>) -> Self {
128        self.pattern = pattern.into();
129        self
130    }
131
132    pub fn poll_watcher(mut self, interval: Duration) -> Self {
133        self.poll_watcher = Some(interval);
134        self
135    }
136
137    pub fn tail(mut self, enable: bool) -> Self {
138        self.tail = enable;
139        self
140    }
141
142    pub fn max_line_size(mut self, size: usize) -> Self {
143        self.max_line_size = size;
144        self
145    }
146
147    /// Build the multi-consumer tailer.
148    pub async fn build(self) -> anyhow::Result<MultiConsumerTailer> {
149        // Collect checkpoint paths without borrowing consumers across
150        // an await (consumers contains non-Sync filter closures).
151        let cp_paths: Vec<Option<Utf8PathBuf>> = self
152            .consumers
153            .iter()
154            .map(|c| {
155                c.checkpoint_name
156                    .as_ref()
157                    .map(|name| self.directory.join(format!(".{name}")))
158            })
159            .collect();
160
161        // Now load checkpoints (async) without borrowing consumers.
162        let mut consumer_checkpoints: Vec<Option<CheckpointData>> =
163            Vec::with_capacity(cp_paths.len());
164        let mut earliest_checkpoint: Option<CheckpointData> = None;
165
166        for cp_path in &cp_paths {
167            let cp = if self.tail {
168                resolve_tail_checkpoint(&self.directory, &self.pattern)?
169            } else if let Some(cp_path) = cp_path {
170                CheckpointData::load(cp_path).await?
171            } else {
172                None
173            };
174
175            match (&earliest_checkpoint, &cp) {
176                (None, Some(cp)) => {
177                    earliest_checkpoint = Some(cp.clone());
178                }
179                (Some(existing), Some(cp)) => {
180                    if cp.file < existing.file
181                        || (cp.file == existing.file && cp.line < existing.line)
182                    {
183                        earliest_checkpoint = Some(cp.clone());
184                    }
185                }
186                _ => {}
187            }
188
189            consumer_checkpoints.push(cp);
190        }
191
192        let closed = Arc::new(AtomicBool::new(false));
193        let close_notify = Arc::new(Notify::new());
194
195        let fs_notify = Arc::new(Notify::new());
196        let fs_notify_tx = fs_notify.clone();
197        let event_handler = move |res: Result<Event, _>| match res {
198            Ok(event) => match event.kind {
199                EventKind::Create(CreateKind::File) | EventKind::Modify(ModifyKind::Data(_)) => {
200                    fs_notify_tx.notify_one();
201                }
202                _ => {}
203            },
204            Err(_) => {}
205        };
206        let mut watcher: Box<dyn Watcher + Send> = if let Some(interval) = self.poll_watcher {
207            Box::new(notify::PollWatcher::new(
208                event_handler,
209                notify::Config::default().with_poll_interval(interval),
210            )?)
211        } else {
212            Box::new(notify::recommended_watcher(event_handler)?)
213        };
214        watcher.watch(
215            &self.directory.clone().into_std_path_buf(),
216            notify::RecursiveMode::NonRecursive,
217        )?;
218
219        let shared = Arc::new(TailerShared {
220            closed,
221            close_notify,
222        });
223
224        let stream = make_multi_stream(
225            self.directory,
226            self.pattern,
227            self.max_line_size,
228            self.consumers,
229            earliest_checkpoint,
230            consumer_checkpoints,
231            cp_paths,
232            fs_notify,
233            shared.clone(),
234        );
235
236        Ok(MultiConsumerTailer {
237            close_handle: CloseHandle { shared },
238            _watcher: watcher,
239            stream: Box::pin(stream),
240        })
241    }
242}
243
244// ---------------------------------------------------------------------------
245// Shared internals
246// ---------------------------------------------------------------------------
247
248struct TailerShared {
249    closed: Arc<AtomicBool>,
250    close_notify: Arc<Notify>,
251}
252
253/// A `Send + Sync` handle that can close a tailer from any context.
254#[derive(Clone)]
255pub struct CloseHandle {
256    shared: Arc<TailerShared>,
257}
258
259impl CloseHandle {
260    /// Signal the stream to terminate.
261    pub fn close(&self) {
262        self.shared.closed.store(true, Ordering::SeqCst);
263        self.shared.close_notify.notify_waiters();
264    }
265}
266
267// ---------------------------------------------------------------------------
268// MultiConsumerTailer
269// ---------------------------------------------------------------------------
270
271/// An async Stream that yields vectors of [`LogBatch`], one per consumer
272/// whose batch is ready.
273pub struct MultiConsumerTailer {
274    close_handle: CloseHandle,
275    _watcher: Box<dyn Watcher + Send>,
276    stream: std::pin::Pin<Box<dyn Stream<Item = anyhow::Result<Vec<LogBatch>>> + Send>>,
277}
278
279impl MultiConsumerTailer {
280    pub fn close_handle(&self) -> CloseHandle {
281        self.close_handle.clone()
282    }
283
284    pub fn close(&self) {
285        self.close_handle.close();
286    }
287}
288
289impl Stream for MultiConsumerTailer {
290    type Item = anyhow::Result<Vec<LogBatch>>;
291
292    fn poll_next(
293        mut self: std::pin::Pin<&mut Self>,
294        cx: &mut std::task::Context<'_>,
295    ) -> std::task::Poll<Option<Self::Item>> {
296        if self.close_handle.shared.closed.load(Ordering::SeqCst) {
297            return std::task::Poll::Ready(None);
298        }
299        self.stream.as_mut().poll_next(cx)
300    }
301}
302
303// ---------------------------------------------------------------------------
304// Single-consumer LogTailerConfig / LogTailer (delegates to multi-consumer)
305// ---------------------------------------------------------------------------
306
307/// Configuration for constructing a single-consumer [`LogTailer`].
308#[derive(Deserialize, Serialize)]
309pub struct LogTailerConfig {
310    pub directory: Utf8PathBuf,
311    #[serde(default = "default_pattern")]
312    pub pattern: String,
313    #[serde(default = "default_max_batch_size")]
314    pub max_batch_size: usize,
315    #[serde(default = "default_max_batch_latency", with = "duration_serde")]
316    pub max_batch_latency: Duration,
317    #[serde(default, skip_serializing_if = "Option::is_none")]
318    pub checkpoint_name: Option<String>,
319    #[serde(
320        default,
321        with = "duration_serde",
322        skip_serializing_if = "Option::is_none"
323    )]
324    pub poll_watcher: Option<Duration>,
325    #[serde(default)]
326    pub tail: bool,
327    #[serde(default = "default_max_line_size")]
328    pub max_line_size: usize,
329}
330
331impl LogTailerConfig {
332    pub fn new(directory: Utf8PathBuf) -> Self {
333        Self {
334            directory,
335            pattern: default_pattern(),
336            max_batch_size: default_max_batch_size(),
337            max_batch_latency: default_max_batch_latency(),
338            checkpoint_name: None,
339            poll_watcher: None,
340            tail: false,
341            max_line_size: default_max_line_size(),
342        }
343    }
344
345    pub fn pattern(mut self, pattern: impl Into<String>) -> Self {
346        self.pattern = pattern.into();
347        self
348    }
349
350    pub fn max_batch_size(mut self, size: usize) -> Self {
351        self.max_batch_size = size;
352        self
353    }
354
355    pub fn max_line_size(mut self, size: usize) -> Self {
356        self.max_line_size = size;
357        self
358    }
359
360    pub fn max_batch_latency(mut self, latency: Duration) -> Self {
361        self.max_batch_latency = latency;
362        self
363    }
364
365    pub fn checkpoint_name(mut self, name: impl Into<String>) -> Self {
366        self.checkpoint_name = Some(name.into());
367        self
368    }
369
370    pub fn poll_watcher(mut self, interval: Duration) -> Self {
371        self.poll_watcher = Some(interval);
372        self
373    }
374
375    pub fn tail(mut self, enable: bool) -> Self {
376        self.tail = enable;
377        self
378    }
379
380    /// Build a single-consumer tailer.
381    pub async fn build(self) -> anyhow::Result<LogTailer> {
382        self.build_with_filter(None::<fn(&serde_json::Value) -> anyhow::Result<bool>>)
383            .await
384    }
385
386    /// Build a single-consumer tailer with an optional record filter.
387    pub async fn build_with_filter<F>(self, filter: Option<F>) -> anyhow::Result<LogTailer>
388    where
389        F: Fn(&serde_json::Value) -> anyhow::Result<bool> + Send + 'static,
390    {
391        let mut consumer = ConsumerConfig::new("default")
392            .max_batch_size(self.max_batch_size)
393            .max_batch_latency(self.max_batch_latency);
394        if let Some(name) = self.checkpoint_name.clone() {
395            consumer = consumer.checkpoint_name(name);
396        }
397        if let Some(f) = filter {
398            consumer = consumer.filter(f);
399        }
400
401        let multi_config = MultiConsumerTailerConfig {
402            directory: self.directory,
403            pattern: self.pattern,
404            poll_watcher: self.poll_watcher,
405            tail: self.tail,
406            max_line_size: self.max_line_size,
407            consumers: vec![consumer],
408        };
409
410        let multi = multi_config.build().await?;
411
412        Ok(LogTailer { inner: multi })
413    }
414}
415
416/// A single-consumer async Stream that yields one [`LogBatch`] at a time.
417///
418/// This is a convenience wrapper around [`MultiConsumerTailer`] with
419/// exactly one consumer.
420pub struct LogTailer {
421    inner: MultiConsumerTailer,
422}
423
424impl LogTailer {
425    pub fn close_handle(&self) -> CloseHandle {
426        self.inner.close_handle()
427    }
428
429    pub fn close(&self) {
430        self.inner.close();
431    }
432}
433
434impl Stream for LogTailer {
435    type Item = anyhow::Result<LogBatch>;
436
437    fn poll_next(
438        mut self: std::pin::Pin<&mut Self>,
439        cx: &mut std::task::Context<'_>,
440    ) -> std::task::Poll<Option<Self::Item>> {
441        use std::task::Poll;
442        // The inner multi-consumer stream yields Vec<LogBatch> with exactly
443        // one element.  Unwrap it.
444        match std::pin::Pin::new(&mut self.inner).poll_next(cx) {
445            Poll::Ready(Some(Ok(mut batches))) => Poll::Ready(Some(Ok(batches
446                .pop()
447                .expect("single consumer yields one batch")))),
448            Poll::Ready(Some(Err(e))) => Poll::Ready(Some(Err(e))),
449            Poll::Ready(None) => Poll::Ready(None),
450            Poll::Pending => Poll::Pending,
451        }
452    }
453}
454
455// ---------------------------------------------------------------------------
456// Shared utilities
457// ---------------------------------------------------------------------------
458
459fn resolve_tail_checkpoint(
460    directory: &Utf8PathBuf,
461    pattern: &str,
462) -> anyhow::Result<Option<CheckpointData>> {
463    let glob = Glob::new(pattern)?;
464    let mut files = vec![];
465    for path in glob.walk(directory) {
466        let path = directory.join(Utf8PathBuf::try_from(path).map_err(|e| anyhow::anyhow!("{e}"))?);
467        if path.is_file() {
468            files.push(path);
469        }
470    }
471    files.sort();
472    Ok(files.last().map(|f| CheckpointData {
473        file: f.to_string(),
474        line: 0,
475    }))
476}
477
478/// Build the file plan: sorted list of matching files in the directory.
479fn build_plan(
480    directory: &Utf8PathBuf,
481    pattern: &str,
482    checkpoint_paths: &[Option<Utf8PathBuf>],
483    last_processed: &Option<Utf8PathBuf>,
484    checkpoint: &Option<CheckpointData>,
485    bad_files: &HashSet<Utf8PathBuf>,
486) -> anyhow::Result<Vec<Utf8PathBuf>> {
487    let glob = Glob::new(pattern)?;
488    let mut result = vec![];
489    for path in glob.walk(directory) {
490        let path = directory.join(Utf8PathBuf::try_from(path).map_err(|e| anyhow::anyhow!("{e}"))?);
491        // Skip checkpoint files
492        if checkpoint_paths.iter().any(|cp| cp.as_ref() == Some(&path)) {
493            continue;
494        }
495        // Skip files we've previously determined to be unreadable
496        // (foreign content or unrecoverable corruption).
497        if bad_files.contains(&path) {
498            continue;
499        }
500        if path.is_file() {
501            result.push(path);
502        }
503    }
504    result.sort();
505
506    if let Some(last) = last_processed {
507        result.retain(|item| item > last);
508    } else if let Some(cp) = checkpoint {
509        let cp_file = &cp.file;
510        result.retain(|item| item.as_str() >= cp_file.as_str());
511    }
512
513    Ok(result)
514}
515
516fn is_file_done(path: &Utf8PathBuf) -> bool {
517    path.metadata()
518        .map(|m| m.permissions().readonly())
519        .unwrap_or(false)
520}
521
522// ---------------------------------------------------------------------------
523// Per-consumer state used during stream construction
524// ---------------------------------------------------------------------------
525
526// ---------------------------------------------------------------------------
527// Multi-consumer stream
528// ---------------------------------------------------------------------------
529
530fn make_multi_stream(
531    directory: Utf8PathBuf,
532    pattern: String,
533    max_line_size: usize,
534    consumers: Vec<ConsumerConfig>,
535    earliest_checkpoint: Option<CheckpointData>,
536    mut consumer_checkpoints: Vec<Option<CheckpointData>>,
537    cp_paths: Vec<Option<Utf8PathBuf>>,
538    fs_notify: Arc<Notify>,
539    shared: Arc<TailerShared>,
540) -> impl Stream<Item = anyhow::Result<Vec<LogBatch>>> + Send {
541    let num_consumers = consumers.len();
542
543    // Extract per-consumer config into parallel vecs
544    let consumer_names: Vec<String> = consumers.iter().map(|c| c.name.clone()).collect();
545    let max_batch_sizes: Vec<usize> = consumers.iter().map(|c| c.max_batch_size).collect();
546    let max_batch_latencies: Vec<Duration> =
547        consumers.iter().map(|c| c.max_batch_latency).collect();
548    let filters: Vec<Option<Box<dyn Fn(&serde_json::Value) -> anyhow::Result<bool> + Send>>> =
549        consumers.into_iter().map(|c| c.filter).collect();
550
551    async_stream::try_stream! {
552        let mut last_processed: Option<Utf8PathBuf> = None;
553        let mut global_checkpoint = earliest_checkpoint;
554        let mut skip_lines: usize;
555        let retry_delay = Duration::from_millis(200);
556        // Files that produced unrecoverable read errors (foreign content,
557        // corruption, etc).  Kept in-memory for the lifetime of the tailer
558        // so we don't re-attempt them on every directory rescan.  Unlike
559        // `last_processed`, this does not interact with file ordering, so
560        // a bad file whose name sorts after future legitimate segments
561        // does not hide them.
562        let mut bad_files: HashSet<Utf8PathBuf> = HashSet::new();
563
564        // Per-consumer skip lines (for the first file only, when
565        // resuming from checkpoint).  The global skip_lines is the
566        // minimum across all consumers for that file, and individual
567        // consumers that are further ahead will have their records
568        // filtered out by index comparison.
569        let mut consumer_skip: Vec<usize> = vec![0; num_consumers];
570
571        'outer: loop {
572            if shared.closed.load(Ordering::SeqCst) {
573                break;
574            }
575
576            let plan = build_plan(
577                &directory,
578                &pattern,
579                &cp_paths,
580                &last_processed,
581                &global_checkpoint,
582                &bad_files,
583            )?;
584
585            if plan.is_empty() {
586                tokio::select! {
587                    _ = shared.close_notify.notified() => break,
588                    _ = fs_notify.notified() => continue,
589                }
590            }
591
592            // Determine skip_lines from the earliest checkpoint
593            if let Some(cp) = &global_checkpoint {
594                if plan.first().map(|p| p.as_str()) == Some(cp.file.as_str()) {
595                    skip_lines = cp.line;
596                } else {
597                    skip_lines = 0;
598                }
599            } else {
600                skip_lines = 0;
601            }
602
603            // Determine per-consumer skip lines
604            for i in 0..num_consumers {
605                if let Some(cp) = &consumer_checkpoints[i] {
606                    if plan.first().map(|p| p.as_str()) == Some(cp.file.as_str()) {
607                        consumer_skip[i] = cp.line;
608                    } else {
609                        consumer_skip[i] = 0;
610                    }
611                } else {
612                    consumer_skip[i] = 0;
613                }
614            }
615            global_checkpoint.take();
616            for cp in consumer_checkpoints.iter_mut() {
617                cp.take();
618            }
619
620            let mut plan_index = 0;
621            let mut decomp: Option<FileDecompressor> = None;
622            let mut current_path: Option<&Utf8PathBuf> = None;
623            let mut last_lines_consumed: usize = 0;
624            // Track the global line number for the current file so we
625            // can apply per-consumer skip logic.
626            let mut global_line_in_file: usize = skip_lines;
627
628            if let Some(path) = plan.get(plan_index) {
629                let path_std = path.as_std_path().to_owned();
630                decomp = Some(FileDecompressor::open_with_max_line_size(
631                    &path_std,
632                    max_line_size,
633                )?);
634                current_path = Some(path);
635            }
636
637            // Per-consumer batches and deadlines persist across
638            // fill/yield cycles.  A consumer's batch is "ready" when
639            // it is full or its deadline has expired.  Only ready
640            // batches are yielded; others keep accumulating.
641            let mut batches: Vec<LogBatch> = (0..num_consumers)
642                .map(|i| LogBatch::with_consumer_name(consumer_names[i].clone()))
643                .collect();
644            let mut deadlines: Vec<Option<tokio::time::Instant>> = vec![None; num_consumers];
645
646            while decomp.is_some() {
647                if shared.closed.load(Ordering::SeqCst) {
648                    break 'outer;
649                }
650
651                // Fill batches until at least one is ready
652                'fill: loop {
653                    if shared.closed.load(Ordering::SeqCst) {
654                        break 'outer;
655                    }
656
657                    // Check if any consumer already has a ready batch
658                    let now = tokio::time::Instant::now();
659                    let any_ready = (0..num_consumers).any(|i| {
660                        !batches[i].is_empty()
661                            && (batches[i].len() >= max_batch_sizes[i]
662                                || deadlines[i].map_or(false, |d| now >= d))
663                    });
664                    if any_ready {
665                        break 'fill;
666                    }
667
668                    let d = decomp.as_mut().expect("checked above");
669                    let path = current_path.expect("set with decomp");
670
671                    // When set after processing the current decompressor
672                    // call, advance to the next file in the plan.
673                    // `Some(true)`  -> mark the file as bad (foreign or
674                    //                  unrecoverable corruption); do not
675                    //                  update `last_processed`.
676                    // `Some(false)` -> file completed normally; update
677                    //                  `last_processed`.
678                    let mut advance_file: Option<bool> = None;
679
680                    match d.next_line(skip_lines) {
681                        Ok(NextLine::Line(line)) => {
682                            match serde_json::from_str::<serde_json::Value>(&line.text) {
683                                Ok(value) => {
684                                    for i in 0..num_consumers {
685                                        if global_line_in_file < consumer_skip[i] {
686                                            continue;
687                                        }
688                                        if let Some(ref f) = filters[i] {
689                                            if !f(&value)? {
690                                                continue;
691                                            }
692                                        }
693                                        batches[i].push_value(
694                                            value.clone(),
695                                            path,
696                                            line.byte_offset,
697                                        );
698                                        // Start the deadline timer on first record
699                                        if deadlines[i].is_none() {
700                                            deadlines[i] = Some(
701                                                tokio::time::Instant::now() + max_batch_latencies[i],
702                                            );
703                                        }
704                                    }
705                                    global_line_in_file += 1;
706                                }
707                                Err(err) => {
708                                    warn!(
709                                        "Failed to parse a line from {path} (byte offset {}) \
710                                         as json: {err}. Skipping remainder of this file; \
711                                         move it aside if it is not a kumo-jsonl segment.",
712                                        line.byte_offset
713                                    );
714                                    advance_file = Some(true);
715                                }
716                            }
717                        }
718                        Ok(NextLine::Skipped { byte_offset, bytes }) => {
719                            // A record exceeded max_line_size and was
720                            // discarded, but the stream is intact, so the rest
721                            // of the segment is still readable. Count it as a
722                            // consumed line (the decompressor already advanced
723                            // its checkpoint) and keep going.
724                            warn!(
725                                "Skipping a record from {path} at byte offset {byte_offset} \
726                                 ({bytes} bytes) that exceeds max_line_size; the rest of the \
727                                 segment is still processed."
728                            );
729                            global_line_in_file += 1;
730                        }
731                        Ok(NextLine::None) => {
732                            // EOF on current file
733                            if is_file_done(path) {
734                                if d.is_discarding_oversized_record() {
735                                    // The final record was still being
736                                    // discarded for exceeding max_line_size
737                                    // when the segment ended, without ever
738                                    // reaching its terminating newline. Drop
739                                    // it and treat the file as complete.
740                                    warn!(
741                                        "segment {path} ended while discarding a trailing record \
742                                         that exceeds max_line_size; the record is dropped and \
743                                         the segment is treated as complete."
744                                    );
745                                } else if d.has_partial_data() {
746                                    // The writer was killed before it
747                                    // could finish the zstd stream and the
748                                    // segment was later marked done by the
749                                    // producer's startup sweep over
750                                    // abandoned segments.  Discard the
751                                    // partial trailing line and treat the
752                                    // file as complete.
753                                    warn!(
754                                        "unexpected EOF for {} with partial line data remaining; \
755                                         the writer exited without flushing the zstd stream. \
756                                         Discarding partial trailing data and advancing to next \
757                                         segment.",
758                                        path
759                                    );
760                                }
761                                advance_file = Some(false);
762                            } else {
763                                // File not done; find the earliest deadline
764                                // among non-empty batches to bound the wait.
765                                let earliest_deadline = (0..num_consumers)
766                                    .filter(|&i| !batches[i].is_empty())
767                                    .filter_map(|i| deadlines[i])
768                                    .min();
769
770                                if let Some(deadline) = earliest_deadline {
771                                    let remaining = deadline.saturating_duration_since(
772                                        tokio::time::Instant::now(),
773                                    );
774                                    if remaining.is_zero() {
775                                        break 'fill;
776                                    }
777                                    tokio::select! {
778                                        _ = shared.close_notify.notified() => break 'outer,
779                                        _ = tokio::time::sleep(remaining.min(retry_delay)) => {},
780                                        _ = fs_notify.notified() => {},
781                                    }
782                                } else {
783                                    // All batches empty, file not done — wait
784                                    tokio::select! {
785                                        _ = shared.close_notify.notified() => break 'outer,
786                                        _ = tokio::time::sleep(retry_delay) => {},
787                                        _ = fs_notify.notified() => {},
788                                    }
789                                }
790                                d.reset_eof();
791                            }
792                        }
793                        Err(e) => {
794                            // Any decompression error means this file is
795                            // not consumable -- either it is foreign
796                            // content that was dropped into the directory,
797                            // or its bytes are corrupt in a way that
798                            // waiting cannot resolve (an incomplete-but-
799                            // well-formed stream surfaces as Ok(None), not
800                            // Err).  Log and skip.
801                            warn!(
802                                "Error decompressing {path}: {e:#}. Skipping this file; \
803                                 move it aside if it is not a kumo-jsonl segment."
804                            );
805                            advance_file = Some(true);
806                        }
807                    }
808
809                    if let Some(mark_bad) = advance_file {
810                        if mark_bad {
811                            bad_files.insert(path.clone());
812                        } else {
813                            last_lines_consumed = d.lines_consumed;
814                            last_processed = Some(path.clone());
815                        }
816                        skip_lines = 0;
817                        global_line_in_file = 0;
818                        for cs in consumer_skip.iter_mut() {
819                            *cs = 0;
820                        }
821                        plan_index += 1;
822                        if let Some(next_path) = plan.get(plan_index) {
823                            let path_std = next_path.as_std_path().to_owned();
824                            decomp = Some(FileDecompressor::open_with_max_line_size(
825                                &path_std,
826                                max_line_size,
827                            )?);
828                            current_path = Some(next_path);
829                            continue 'fill;
830                        } else {
831                            decomp = None;
832                            current_path = None;
833                            break 'fill;
834                        }
835                    }
836                }
837
838                // Determine which batches are ready to yield
839                let now = tokio::time::Instant::now();
840                let mut ready: Vec<LogBatch> = Vec::new();
841                for i in 0..num_consumers {
842                    let is_ready = !batches[i].is_empty()
843                        && (batches[i].len() >= max_batch_sizes[i]
844                            || deadlines[i].map_or(false, |d| now >= d)
845                            || decomp.is_none()); // end of plan: flush all
846
847                    if !is_ready {
848                        continue;
849                    }
850
851                    // Swap out the ready batch, replace with a fresh one
852                    let mut batch = std::mem::replace(
853                        &mut batches[i],
854                        LogBatch::with_consumer_name(consumer_names[i].clone()),
855                    );
856                    deadlines[i] = None;
857
858                    // Set the commit callback
859                    if let Some(ref cp_path) = cp_paths[i] {
860                        let (cp_file, cp_line) = if let Some(d) = &decomp {
861                            let path = current_path.expect("set with decomp");
862                            (path.clone(), d.lines_consumed)
863                        } else if let Some(last) = &last_processed {
864                            (last.clone(), last_lines_consumed)
865                        } else {
866                            unreachable!("non-empty batch without a source");
867                        };
868                        let cp_path = cp_path.clone();
869                        batch.set_commit_fn(Box::new(move || {
870                            CheckpointData::save_atomic(&cp_path, &cp_file, cp_line)
871                        }));
872                    }
873                    ready.push(batch);
874                }
875
876                if !ready.is_empty() {
877                    skip_lines = 0;
878                    yield ready;
879                }
880            }
881        }
882    }
883}