Skip to main content

common_telemetry/logging/
file_retention.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! File-log retention based on the total size of managed log files.
16
17use std::collections::hash_map::Entry;
18use std::collections::{BTreeMap, HashMap, HashSet};
19use std::fs;
20use std::io::{self, Write};
21use std::path::PathBuf;
22use std::sync::Arc;
23use std::time::SystemTime;
24
25use common_base::readable_size::ReadableSize;
26use parking_lot::Mutex;
27use tracing_appender::rolling::{RollingFileAppender, Rotation};
28
29use crate::logging::LoggingOptions;
30
31/// A managed file-log kind.
32#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)]
33pub(crate) enum LogFileKind {
34    Default,
35    Error,
36    SlowQuery,
37}
38
39impl LogFileKind {
40    fn prefix(self) -> &'static str {
41        match self {
42            Self::Default => "greptimedb",
43            Self::Error => "greptimedb-err",
44            Self::SlowQuery => "greptimedb-slow-queries",
45        }
46    }
47
48    fn current_link_name(self) -> &'static str {
49        match self {
50            Self::Default => ".greptimedb.current",
51            Self::Error => ".greptimedb-err.current",
52            Self::SlowQuery => ".greptimedb-slow-queries.current",
53        }
54    }
55
56    fn kind_from_file_name(name: &str) -> Option<Self> {
57        [Self::SlowQuery, Self::Error, Self::Default]
58            .into_iter()
59            .find_map(|kind| {
60                name.strip_prefix(kind.prefix())
61                    .and_then(|suffix| suffix.strip_prefix('.'))
62                    .filter(|suffix| Self::is_hourly_suffix(suffix))
63                    .map(|_| kind)
64            })
65    }
66
67    fn is_hourly_suffix(suffix: &str) -> bool {
68        suffix.len() == 13
69            && suffix.bytes().enumerate().all(|(index, byte)| {
70                matches!(index, 4 | 7 | 10) && byte == b'-'
71                    || !matches!(index, 4 | 7 | 10) && byte.is_ascii_digit()
72            })
73    }
74}
75
76#[derive(Clone, Debug)]
77struct LogFile {
78    kind: LogFileKind,
79    path: PathBuf,
80    last_modified: SystemTime,
81    size: u64,
82}
83
84#[derive(Default)]
85struct FileIndex {
86    active: HashMap<LogFileKind, LogFile>,
87    closed: BTreeMap<(SystemTime, PathBuf), LogFile>,
88    total_size: u64,
89}
90
91impl FileIndex {
92    fn track(&mut self, kind: LogFileKind, path: PathBuf, size: u64, last_modified: SystemTime) {
93        let active = match self.active.entry(kind) {
94            Entry::Occupied(mut entry) if entry.get().path != path => {
95                let previous = entry.insert(LogFile {
96                    kind,
97                    path,
98                    last_modified,
99                    size: 0,
100                });
101                self.closed
102                    .insert((previous.last_modified, previous.path.clone()), previous);
103                entry.into_mut()
104            }
105            Entry::Occupied(entry) => entry.into_mut(),
106            Entry::Vacant(entry) => entry.insert(LogFile {
107                kind,
108                path,
109                last_modified,
110                size: 0,
111            }),
112        };
113
114        active.size = active.size.saturating_add(size);
115        active.last_modified = last_modified;
116        self.total_size = self.total_size.saturating_add(size);
117    }
118
119    fn count(&self, kind: LogFileKind) -> usize {
120        usize::from(self.active.contains_key(&kind))
121            + self
122                .closed
123                .values()
124                .filter(|file| file.kind == kind)
125                .count()
126    }
127
128    fn remove(&mut self, key: &(SystemTime, PathBuf)) {
129        if let Some(file) = self.closed.remove(key) {
130            self.total_size = self.total_size.saturating_sub(file.size);
131        }
132    }
133}
134
135#[derive(Default)]
136struct RetentionState {
137    kinds: HashSet<LogFileKind>,
138    files: FileIndex,
139    initialized: bool,
140    cleanup_error_reported: bool,
141}
142
143/// Retains managed file logs within configured directory-size and file-count limits.
144#[derive(Clone)]
145pub(crate) struct DirectoryRetention {
146    directory: Arc<PathBuf>,
147    max_size: u64,
148    max_log_files: usize,
149    state: Arc<Mutex<RetentionState>>,
150}
151
152impl DirectoryRetention {
153    /// Returns a retention manager when the configured limit is enabled.
154    pub(crate) fn new(
155        directory: impl Into<PathBuf>,
156        max_size: ReadableSize,
157        max_log_files: usize,
158    ) -> Option<Self> {
159        (max_size.as_bytes() > 0).then(|| Self {
160            directory: Arc::new(directory.into()),
161            max_size: max_size.as_bytes(),
162            max_log_files,
163            state: Arc::new(Mutex::new(RetentionState::default())),
164        })
165    }
166
167    /// Loads the initial file state after all enabled file-log kinds are registered.
168    pub(crate) fn initialize(&self) {
169        let mut state = self.state.lock();
170        if !self.ensure_initialized(&mut state) {
171            return;
172        }
173
174        self.prune_files(&mut state);
175        self.prune_size(&mut state, 0);
176    }
177
178    fn register(&self, kind: LogFileKind) {
179        let mut state = self.state.lock();
180        debug_assert!(!state.initialized);
181        state.kinds.insert(kind);
182    }
183
184    fn reclaim(&self, incoming_size: u64) {
185        let mut state = self.state.lock();
186        if !self.ensure_initialized(&mut state) {
187            return;
188        }
189
190        self.prune_size(&mut state, incoming_size);
191    }
192
193    fn track(&self, kind: LogFileKind, written: u64) {
194        let latest_file = self.current(kind);
195        let last_modified = SystemTime::now();
196        let mut state = self.state.lock();
197        if !self.ensure_initialized(&mut state) {
198            return;
199        }
200
201        let path = match latest_file {
202            Ok(latest_file) => latest_file,
203            Err(_) => {
204                self.report_error(&mut state, "resolving the latest log file");
205                self.reconcile(&mut state);
206                return;
207            }
208        };
209
210        state.files.track(kind, path, written, last_modified);
211        self.prune_files(&mut state);
212        self.prune_size(&mut state, 0);
213    }
214
215    fn scan(&self, kinds: &HashSet<LogFileKind>) -> io::Result<FileIndex> {
216        let active_paths = kinds
217            .iter()
218            .map(|kind| self.current(*kind).map(|path| (*kind, path)))
219            .collect::<io::Result<HashMap<_, _>>>()?;
220        let mut files = FileIndex::default();
221
222        for entry in fs::read_dir(self.directory.as_ref())? {
223            let entry = entry?;
224            let file_name = entry.file_name();
225            let Some(file_name) = file_name.to_str() else {
226                continue;
227            };
228            let Some(kind) = LogFileKind::kind_from_file_name(file_name) else {
229                continue;
230            };
231
232            let path = entry.path();
233            let metadata = fs::symlink_metadata(&path)?;
234            if !metadata.is_file() {
235                continue;
236            }
237
238            let file = LogFile {
239                kind,
240                path: path.clone(),
241                last_modified: metadata.modified()?,
242                size: metadata.len(),
243            };
244            files.total_size = files.total_size.saturating_add(file.size);
245            if active_paths.get(&kind) == Some(&path) {
246                files.active.insert(kind, file);
247            } else {
248                files.closed.insert((file.last_modified, path), file);
249            }
250        }
251
252        if files.active.len() != active_paths.len() {
253            return Err(io::Error::new(
254                io::ErrorKind::NotFound,
255                "latest log symlink target does not exist",
256            ));
257        }
258
259        Ok(files)
260    }
261
262    /// Rebuilds the in-memory index after an unexpected filesystem error.
263    fn reconcile(&self, state: &mut RetentionState) -> bool {
264        match self.scan(&state.kinds) {
265            Ok(files) => {
266                state.files = files;
267                state.initialized = true;
268                state.cleanup_error_reported = false;
269                true
270            }
271            Err(_) => {
272                self.report_error(state, "loading log directory");
273                false
274            }
275        }
276    }
277
278    fn ensure_initialized(&self, state: &mut RetentionState) -> bool {
279        state.initialized || self.reconcile(state)
280    }
281
282    fn current(&self, kind: LogFileKind) -> io::Result<PathBuf> {
283        let symlink = self.directory.join(kind.current_link_name());
284        let target = fs::read_link(&symlink)?;
285        let Some(file_name) = target.file_name().and_then(|name| name.to_str()) else {
286            return Err(io::Error::new(
287                io::ErrorKind::InvalidData,
288                format!("invalid latest log symlink {}", symlink.display()),
289            ));
290        };
291        let Some(candidate) = LogFileKind::kind_from_file_name(file_name) else {
292            return Err(io::Error::new(
293                io::ErrorKind::InvalidData,
294                format!("invalid latest log symlink {}", symlink.display()),
295            ));
296        };
297        if candidate != kind {
298            return Err(io::Error::new(
299                io::ErrorKind::InvalidData,
300                format!("invalid latest log symlink {}", symlink.display()),
301            ));
302        }
303
304        Ok(self.directory.join(file_name))
305    }
306
307    fn prune_files(&self, state: &mut RetentionState) {
308        if self.max_log_files == 0 {
309            return;
310        }
311
312        for kind in state.kinds.clone() {
313            while state.files.count(kind) > self.max_log_files {
314                let Some((key, entry)) = state
315                    .files
316                    .closed
317                    .iter()
318                    .find(|(_, file)| file.kind == kind)
319                    .map(|(key, entry)| (key.clone(), entry.clone()))
320                else {
321                    return;
322                };
323
324                if !self.delete(state, key, entry) {
325                    return;
326                }
327            }
328        }
329    }
330
331    fn prune_size(&self, state: &mut RetentionState, incoming_size: u64) {
332        while state.files.total_size.saturating_add(incoming_size) > self.max_size {
333            let Some((key, entry)) = state
334                .files
335                .closed
336                .first_key_value()
337                .map(|(key, entry)| (key.clone(), entry.clone()))
338            else {
339                // The limit is intentionally soft: active files and a single
340                // oversized log record are never removed or truncated.
341                self.report_error(
342                    state,
343                    "log directory exceeds the configured limit but only active files remain",
344                );
345                return;
346            };
347
348            if !self.delete(state, key, entry) {
349                return;
350            }
351        }
352    }
353
354    fn delete(
355        &self,
356        state: &mut RetentionState,
357        key: (SystemTime, PathBuf),
358        file: LogFile,
359    ) -> bool {
360        match fs::remove_file(&file.path) {
361            Ok(()) => {
362                state.files.remove(&key);
363                state.cleanup_error_reported = false;
364                true
365            }
366            Err(error) => {
367                self.report_error(state, &format!("removing {}: {error}", file.path.display()));
368                if error.kind() == io::ErrorKind::NotFound {
369                    self.reconcile(state);
370                }
371                false
372            }
373        }
374    }
375
376    #[allow(clippy::print_stderr)]
377    fn report_error(&self, state: &mut RetentionState, message: &str) {
378        if !state.cleanup_error_reported {
379            // Do not use tracing here: this writer is itself on the tracing path.
380            eprintln!(
381                "Failed to retain log directory {}: {message}",
382                self.directory.display()
383            );
384            state.cleanup_error_reported = true;
385        }
386    }
387}
388
389/// A [`RollingFileAppender`] with shared directory-size retention.
390pub(crate) struct RetentionAppender {
391    inner: RollingFileAppender,
392    retention: DirectoryRetention,
393    kind: LogFileKind,
394}
395
396impl RetentionAppender {
397    fn new(inner: RollingFileAppender, retention: DirectoryRetention, kind: LogFileKind) -> Self {
398        Self {
399            inner,
400            retention,
401            kind,
402        }
403    }
404}
405
406impl Write for RetentionAppender {
407    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
408        self.retention.reclaim(buf.len() as u64);
409        let written = self.inner.write(buf)?;
410        self.retention.track(self.kind, written as u64);
411        Ok(written)
412    }
413
414    fn flush(&mut self) -> io::Result<()> {
415        self.inner.flush()
416    }
417}
418
419/// Builds the existing hourly appender, optionally wrapped with retention.
420pub(crate) fn build_file_appender(
421    opts: &LoggingOptions,
422    kind: LogFileKind,
423    retention: Option<&DirectoryRetention>,
424) -> Box<dyn Write + Send> {
425    // Directory retention owns size and count pruning so its index remains authoritative.
426    let upstream_max_log_files = if retention.is_some() {
427        0
428    } else {
429        opts.max_log_files
430    };
431    let mut builder = RollingFileAppender::builder()
432        .rotation(Rotation::HOURLY)
433        .filename_prefix(kind.prefix())
434        .max_log_files(upstream_max_log_files);
435    if retention.is_some() {
436        builder = builder.latest_symlink(kind.current_link_name());
437    }
438
439    let appender = builder.build(&opts.dir).unwrap_or_else(|error| {
440        panic!(
441            "initializing rolling file appender at {} failed: {}",
442            opts.dir, error
443        )
444    });
445
446    if let Some(retention) = retention {
447        retention.register(kind);
448        Box::new(RetentionAppender::new(appender, retention.clone(), kind))
449    } else {
450        Box::new(appender)
451    }
452}
453
454#[cfg(test)]
455mod tests {
456    #[cfg(unix)]
457    use std::fs;
458    use std::fs::File;
459    use std::io::Write;
460    #[cfg(unix)]
461    use std::os::unix::fs::symlink;
462    use std::path::{Path, PathBuf};
463    use std::time::SystemTime;
464
465    use common_base::readable_size::ReadableSize;
466    use tempfile::TempDir;
467
468    use super::FileIndex;
469    #[cfg(unix)]
470    use crate::logging::LoggingOptions;
471    #[cfg(unix)]
472    use crate::logging::file_retention::build_file_appender;
473    use crate::logging::file_retention::{DirectoryRetention, LogFileKind};
474
475    fn write_file(path: &Path, contents: &[u8]) {
476        let mut file = File::create(path).unwrap();
477        file.write_all(contents).unwrap();
478    }
479
480    #[cfg(unix)]
481    fn register_default_kind(retention: &DirectoryRetention, directory: &Path, active: &Path) {
482        retention.register(LogFileKind::Default);
483        symlink(active, directory.join(".greptimedb.current")).unwrap();
484    }
485
486    #[test]
487    fn test_disabled_retention() {
488        assert!(DirectoryRetention::new("/tmp", ReadableSize::default(), 0).is_none());
489    }
490
491    #[test]
492    fn test_recognizes_managed_file_name() {
493        assert!(LogFileKind::kind_from_file_name("greptimedb.2026-01-01-00").is_some());
494        assert!(LogFileKind::kind_from_file_name("greptimedb-err.2026-01-01-00").is_some());
495        assert!(
496            LogFileKind::kind_from_file_name("greptimedb-slow-queries.2026-01-01-00").is_some()
497        );
498        assert!(LogFileKind::kind_from_file_name("greptimedb.2026-01-01-00.1").is_none());
499    }
500
501    #[test]
502    fn test_file_index_ignores_missing_closed_file() {
503        let mut files = FileIndex::default();
504
505        files.remove(&(SystemTime::UNIX_EPOCH, PathBuf::from("missing")));
506
507        assert_eq!(files.total_size, 0);
508    }
509
510    #[cfg(unix)]
511    #[test]
512    fn test_file_appender_creates_current_link_before_first_write() {
513        let directory = TempDir::new().unwrap();
514        let opts = LoggingOptions {
515            dir: directory.path().display().to_string(),
516            ..Default::default()
517        };
518        let retention = DirectoryRetention::new(directory.path(), ReadableSize::mb(1), 0).unwrap();
519
520        let _appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
521
522        let symlink = directory.path().join(".greptimedb.current");
523        assert!(symlink.is_symlink());
524        assert!(fs::read_link(symlink).unwrap().exists());
525    }
526
527    #[cfg(unix)]
528    #[test]
529    fn test_file_appender_prunes_closed_files_by_size() {
530        let directory = TempDir::new().unwrap();
531        let oldest = directory.path().join("greptimedb.2026-01-01-00");
532        let old = directory.path().join("greptimedb.2026-01-01-01");
533        write_file(&oldest, &vec![b'a'; 2 * 1024]);
534        write_file(&old, &vec![b'b'; 2 * 1024]);
535
536        let opts = LoggingOptions {
537            dir: directory.path().display().to_string(),
538            ..Default::default()
539        };
540        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1024), 0).unwrap();
541        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
542
543        retention.initialize();
544        appender.write_all(b"current").unwrap();
545
546        assert!(!oldest.exists());
547        assert!(!old.exists());
548        assert!(
549            fs::read_link(directory.path().join(".greptimedb.current"))
550                .unwrap()
551                .exists()
552        );
553    }
554
555    #[cfg(unix)]
556    #[test]
557    fn test_file_appender_prunes_oldest_closed_file_by_size() {
558        let directory = TempDir::new().unwrap();
559        let oldest = directory.path().join("greptimedb.2026-01-01-00");
560        let old = directory.path().join("greptimedb.2026-01-01-01");
561        write_file(&oldest, &vec![b'a'; 2 * 1024]);
562        write_file(&old, &vec![b'b'; 512]);
563
564        let opts = LoggingOptions {
565            dir: directory.path().display().to_string(),
566            ..Default::default()
567        };
568        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1024), 0).unwrap();
569        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
570
571        retention.initialize();
572        appender.write_all(b"current").unwrap();
573
574        assert!(!oldest.exists());
575        assert!(old.exists());
576    }
577
578    #[cfg(unix)]
579    #[test]
580    fn test_file_appender_prunes_closed_files_by_count() {
581        let directory = TempDir::new().unwrap();
582        let oldest = directory.path().join("greptimedb.2026-01-01-00");
583        let old = directory.path().join("greptimedb.2026-01-01-01");
584        write_file(&oldest, b"oldest");
585        write_file(&old, b"old");
586
587        let opts = LoggingOptions {
588            dir: directory.path().display().to_string(),
589            ..Default::default()
590        };
591        let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
592        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
593
594        retention.initialize();
595        appender.write_all(b"current").unwrap();
596
597        assert!(!oldest.exists());
598        assert!(old.exists());
599        assert!(
600            fs::read_link(directory.path().join(".greptimedb.current"))
601                .unwrap()
602                .exists()
603        );
604    }
605
606    #[cfg(unix)]
607    #[test]
608    fn test_file_appender_keeps_oversized_active_file() {
609        let directory = TempDir::new().unwrap();
610        let opts = LoggingOptions {
611            dir: directory.path().display().to_string(),
612            ..Default::default()
613        };
614        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
615        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
616
617        retention.initialize();
618        appender.write_all(b"current").unwrap();
619
620        let active = fs::metadata(directory.path().join(".greptimedb.current")).unwrap();
621        assert!(active.len() > 1);
622    }
623
624    #[cfg(unix)]
625    #[test]
626    fn test_retention_removes_closed_files() {
627        let directory = TempDir::new().unwrap();
628        let old = directory.path().join("greptimedb.2026-01-01-00");
629        let active = directory.path().join("greptimedb.2026-01-01-01");
630        write_file(&old, b"old");
631        write_file(&active, b"new");
632
633        let retention = DirectoryRetention::new(directory.path(), ReadableSize(4), 0).unwrap();
634        register_default_kind(&retention, directory.path(), &active);
635        retention.initialize();
636
637        assert!(!old.exists());
638        assert!(active.exists());
639    }
640
641    #[cfg(unix)]
642    #[test]
643    fn test_retention_keeps_active_file() {
644        let directory = TempDir::new().unwrap();
645        let active = directory.path().join("greptimedb.2026-01-01-00");
646        write_file(&active, b"active");
647
648        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
649        register_default_kind(&retention, directory.path(), &active);
650        retention.initialize();
651
652        assert!(active.exists());
653    }
654
655    #[cfg(unix)]
656    #[test]
657    fn test_retention_tracks_rotated_file_in_memory() {
658        let directory = TempDir::new().unwrap();
659        let old = directory.path().join("greptimedb.2026-01-01-00");
660        let active = directory.path().join("greptimedb.2026-01-01-01");
661        write_file(&old, b"old");
662        write_file(&active, b"");
663
664        let retention = DirectoryRetention::new(directory.path(), ReadableSize(3), 0).unwrap();
665        register_default_kind(&retention, directory.path(), &old);
666        retention.initialize();
667        fs::remove_file(directory.path().join(".greptimedb.current")).unwrap();
668        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
669        retention.track(LogFileKind::Default, 1);
670
671        assert!(!old.exists());
672        assert!(active.exists());
673    }
674
675    #[cfg(unix)]
676    #[test]
677    fn test_retention_retries_initialization() {
678        let directory = TempDir::new().unwrap();
679        let old = directory.path().join("greptimedb.2026-01-01-00");
680        let active = directory.path().join("greptimedb.2026-01-01-01");
681        let retention = DirectoryRetention::new(directory.path(), ReadableSize(5), 0).unwrap();
682
683        retention.register(LogFileKind::Default);
684        retention.initialize();
685        assert!(!retention.state.lock().initialized);
686
687        write_file(&old, b"old");
688        write_file(&active, b"new");
689        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
690        retention.reclaim(1);
691
692        assert!(!old.exists());
693        assert!(active.exists());
694        assert!(retention.state.lock().initialized);
695    }
696
697    #[cfg(unix)]
698    #[test]
699    fn test_retention_track_retries_initialization() {
700        let directory = TempDir::new().unwrap();
701        let active = directory.path().join("greptimedb.2026-01-01-00");
702        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
703
704        retention.register(LogFileKind::Default);
705        retention.initialize();
706        assert!(!retention.state.lock().initialized);
707
708        write_file(&active, b"");
709        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
710        retention.track(LogFileKind::Default, 0);
711
712        let state = retention.state.lock();
713        assert!(state.initialized);
714        assert_eq!(state.files.active[&LogFileKind::Default].path, active);
715    }
716
717    #[cfg(unix)]
718    #[test]
719    fn test_retention_reconciles_after_remove_failure() {
720        let directory = TempDir::new().unwrap();
721        let old = directory.path().join("greptimedb.2026-01-01-00");
722        let active = directory.path().join("greptimedb.2026-01-01-01");
723        write_file(&old, b"old");
724        write_file(&active, b"new");
725
726        let retention = DirectoryRetention::new(directory.path(), ReadableSize(6), 0).unwrap();
727        register_default_kind(&retention, directory.path(), &active);
728        retention.initialize();
729        fs::remove_file(&old).unwrap();
730
731        retention.reclaim(1);
732
733        let state = retention.state.lock();
734        assert!(state.files.closed.is_empty());
735        assert_eq!(state.files.total_size, 3);
736    }
737
738    #[cfg(unix)]
739    #[test]
740    fn test_retention_enforces_max_log_files() {
741        let directory = TempDir::new().unwrap();
742        let oldest = directory.path().join("greptimedb.2026-01-01-00");
743        let old = directory.path().join("greptimedb.2026-01-01-01");
744        let active = directory.path().join("greptimedb.2026-01-01-02");
745        write_file(&oldest, b"oldest");
746        write_file(&old, b"old");
747        write_file(&active, b"active");
748
749        let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
750        register_default_kind(&retention, directory.path(), &active);
751        retention.initialize();
752
753        assert!(!oldest.exists());
754        assert!(old.exists());
755        assert!(active.exists());
756    }
757
758    #[cfg(unix)]
759    #[test]
760    fn test_retention_enforces_max_log_files_after_rotation() {
761        let directory = TempDir::new().unwrap();
762        let oldest = directory.path().join("greptimedb.2026-01-01-00");
763        let old = directory.path().join("greptimedb.2026-01-01-01");
764        let active = directory.path().join("greptimedb.2026-01-01-02");
765        write_file(&oldest, b"oldest");
766        write_file(&old, b"old");
767
768        let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
769        register_default_kind(&retention, directory.path(), &old);
770        retention.initialize();
771        write_file(&active, b"");
772        fs::remove_file(directory.path().join(".greptimedb.current")).unwrap();
773        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
774        retention.track(LogFileKind::Default, 1);
775
776        assert!(!oldest.exists());
777        assert!(old.exists());
778        assert!(active.exists());
779    }
780
781    #[test]
782    fn test_unmanaged_files_are_ignored() {
783        let directory = TempDir::new().unwrap();
784        let unmanaged = directory.path().join("keep-me");
785        write_file(&unmanaged, b"unmanaged");
786        write_file(
787            &directory.path().join("greptimedb.2026-01-01-00"),
788            b"managed",
789        );
790
791        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
792        retention.initialize();
793
794        assert!(unmanaged.exists());
795    }
796
797    #[cfg(unix)]
798    #[test]
799    fn test_managed_file_symlink_is_ignored() {
800        let directory = TempDir::new().unwrap();
801        let target = directory.path().join("target");
802        write_file(&target, b"target");
803
804        let link = directory.path().join("greptimedb.2026-01-01-00");
805        symlink(&target, &link).unwrap();
806
807        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
808        retention.initialize();
809
810        assert!(link.exists());
811    }
812}