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
17fn 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
37pub struct ConsumerConfig {
43 pub name: String,
46 pub max_batch_size: usize,
48 pub max_batch_latency: Duration,
50 pub checkpoint_name: Option<String>,
53 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
93pub struct MultiConsumerTailerConfig {
99 pub directory: Utf8PathBuf,
101 pub pattern: String,
103 pub poll_watcher: Option<Duration>,
105 pub tail: bool,
107 pub max_line_size: usize,
111 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 pub async fn build(self) -> anyhow::Result<MultiConsumerTailer> {
149 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 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
244struct TailerShared {
249 closed: Arc<AtomicBool>,
250 close_notify: Arc<Notify>,
251}
252
253#[derive(Clone)]
255pub struct CloseHandle {
256 shared: Arc<TailerShared>,
257}
258
259impl CloseHandle {
260 pub fn close(&self) {
262 self.shared.closed.store(true, Ordering::SeqCst);
263 self.shared.close_notify.notify_waiters();
264 }
265}
266
267pub 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#[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 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 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
416pub 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 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
455fn 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
478fn 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 if checkpoint_paths.iter().any(|cp| cp.as_ref() == Some(&path)) {
493 continue;
494 }
495 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
522fn 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 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 let mut bad_files: HashSet<Utf8PathBuf> = HashSet::new();
563
564 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 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 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 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 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: loop {
653 if shared.closed.load(Ordering::SeqCst) {
654 break 'outer;
655 }
656
657 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 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 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 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 if is_file_done(path) {
734 if d.is_discarding_oversized_record() {
735 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 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 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 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 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 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()); if !is_ready {
848 continue;
849 }
850
851 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 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}