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    use std::fs::{self, File};
457    use std::io::Write;
458    #[cfg(unix)]
459    use std::os::unix::fs::symlink;
460    use std::path::{Path, PathBuf};
461    use std::time::SystemTime;
462
463    use common_base::readable_size::ReadableSize;
464    use tempfile::TempDir;
465
466    use super::FileIndex;
467    use crate::logging::LoggingOptions;
468    use crate::logging::file_retention::{DirectoryRetention, LogFileKind, build_file_appender};
469
470    fn write_file(path: &Path, contents: &[u8]) {
471        let mut file = File::create(path).unwrap();
472        file.write_all(contents).unwrap();
473    }
474
475    #[cfg(unix)]
476    fn register_default_kind(retention: &DirectoryRetention, directory: &Path, active: &Path) {
477        retention.register(LogFileKind::Default);
478        symlink(active, directory.join(".greptimedb.current")).unwrap();
479    }
480
481    #[test]
482    fn test_disabled_retention() {
483        assert!(DirectoryRetention::new("/tmp", ReadableSize::default(), 0).is_none());
484    }
485
486    #[test]
487    fn test_recognizes_managed_file_name() {
488        assert!(LogFileKind::kind_from_file_name("greptimedb.2026-01-01-00").is_some());
489        assert!(LogFileKind::kind_from_file_name("greptimedb-err.2026-01-01-00").is_some());
490        assert!(
491            LogFileKind::kind_from_file_name("greptimedb-slow-queries.2026-01-01-00").is_some()
492        );
493        assert!(LogFileKind::kind_from_file_name("greptimedb.2026-01-01-00.1").is_none());
494    }
495
496    #[test]
497    fn test_file_index_ignores_missing_closed_file() {
498        let mut files = FileIndex::default();
499
500        files.remove(&(SystemTime::UNIX_EPOCH, PathBuf::from("missing")));
501
502        assert_eq!(files.total_size, 0);
503    }
504
505    #[cfg(unix)]
506    #[test]
507    fn test_file_appender_creates_current_link_before_first_write() {
508        let directory = TempDir::new().unwrap();
509        let opts = LoggingOptions {
510            dir: directory.path().display().to_string(),
511            ..Default::default()
512        };
513        let retention = DirectoryRetention::new(directory.path(), ReadableSize::mb(1), 0).unwrap();
514
515        let _appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
516
517        let symlink = directory.path().join(".greptimedb.current");
518        assert!(symlink.is_symlink());
519        assert!(fs::read_link(symlink).unwrap().exists());
520    }
521
522    #[cfg(unix)]
523    #[test]
524    fn test_file_appender_prunes_closed_files_by_size() {
525        let directory = TempDir::new().unwrap();
526        let oldest = directory.path().join("greptimedb.2026-01-01-00");
527        let old = directory.path().join("greptimedb.2026-01-01-01");
528        write_file(&oldest, &vec![b'a'; 2 * 1024]);
529        write_file(&old, &vec![b'b'; 2 * 1024]);
530
531        let opts = LoggingOptions {
532            dir: directory.path().display().to_string(),
533            ..Default::default()
534        };
535        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1024), 0).unwrap();
536        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
537
538        retention.initialize();
539        appender.write_all(b"current").unwrap();
540
541        assert!(!oldest.exists());
542        assert!(!old.exists());
543        assert!(
544            fs::read_link(directory.path().join(".greptimedb.current"))
545                .unwrap()
546                .exists()
547        );
548    }
549
550    #[cfg(unix)]
551    #[test]
552    fn test_file_appender_prunes_oldest_closed_file_by_size() {
553        let directory = TempDir::new().unwrap();
554        let oldest = directory.path().join("greptimedb.2026-01-01-00");
555        let old = directory.path().join("greptimedb.2026-01-01-01");
556        write_file(&oldest, &vec![b'a'; 2 * 1024]);
557        write_file(&old, &vec![b'b'; 512]);
558
559        let opts = LoggingOptions {
560            dir: directory.path().display().to_string(),
561            ..Default::default()
562        };
563        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1024), 0).unwrap();
564        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
565
566        retention.initialize();
567        appender.write_all(b"current").unwrap();
568
569        assert!(!oldest.exists());
570        assert!(old.exists());
571    }
572
573    #[cfg(unix)]
574    #[test]
575    fn test_file_appender_prunes_closed_files_by_count() {
576        let directory = TempDir::new().unwrap();
577        let oldest = directory.path().join("greptimedb.2026-01-01-00");
578        let old = directory.path().join("greptimedb.2026-01-01-01");
579        write_file(&oldest, b"oldest");
580        write_file(&old, b"old");
581
582        let opts = LoggingOptions {
583            dir: directory.path().display().to_string(),
584            ..Default::default()
585        };
586        let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
587        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
588
589        retention.initialize();
590        appender.write_all(b"current").unwrap();
591
592        assert!(!oldest.exists());
593        assert!(old.exists());
594        assert!(
595            fs::read_link(directory.path().join(".greptimedb.current"))
596                .unwrap()
597                .exists()
598        );
599    }
600
601    #[cfg(unix)]
602    #[test]
603    fn test_file_appender_keeps_oversized_active_file() {
604        let directory = TempDir::new().unwrap();
605        let opts = LoggingOptions {
606            dir: directory.path().display().to_string(),
607            ..Default::default()
608        };
609        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
610        let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
611
612        retention.initialize();
613        appender.write_all(b"current").unwrap();
614
615        let active = fs::metadata(directory.path().join(".greptimedb.current")).unwrap();
616        assert!(active.len() > 1);
617    }
618
619    #[cfg(unix)]
620    #[test]
621    fn test_retention_removes_closed_files() {
622        let directory = TempDir::new().unwrap();
623        let old = directory.path().join("greptimedb.2026-01-01-00");
624        let active = directory.path().join("greptimedb.2026-01-01-01");
625        write_file(&old, b"old");
626        write_file(&active, b"new");
627
628        let retention = DirectoryRetention::new(directory.path(), ReadableSize(4), 0).unwrap();
629        register_default_kind(&retention, directory.path(), &active);
630        retention.initialize();
631
632        assert!(!old.exists());
633        assert!(active.exists());
634    }
635
636    #[cfg(unix)]
637    #[test]
638    fn test_retention_keeps_active_file() {
639        let directory = TempDir::new().unwrap();
640        let active = directory.path().join("greptimedb.2026-01-01-00");
641        write_file(&active, b"active");
642
643        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
644        register_default_kind(&retention, directory.path(), &active);
645        retention.initialize();
646
647        assert!(active.exists());
648    }
649
650    #[cfg(unix)]
651    #[test]
652    fn test_retention_tracks_rotated_file_in_memory() {
653        let directory = TempDir::new().unwrap();
654        let old = directory.path().join("greptimedb.2026-01-01-00");
655        let active = directory.path().join("greptimedb.2026-01-01-01");
656        write_file(&old, b"old");
657        write_file(&active, b"");
658
659        let retention = DirectoryRetention::new(directory.path(), ReadableSize(3), 0).unwrap();
660        register_default_kind(&retention, directory.path(), &old);
661        retention.initialize();
662        fs::remove_file(directory.path().join(".greptimedb.current")).unwrap();
663        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
664        retention.track(LogFileKind::Default, 1);
665
666        assert!(!old.exists());
667        assert!(active.exists());
668    }
669
670    #[cfg(unix)]
671    #[test]
672    fn test_retention_retries_initialization() {
673        let directory = TempDir::new().unwrap();
674        let old = directory.path().join("greptimedb.2026-01-01-00");
675        let active = directory.path().join("greptimedb.2026-01-01-01");
676        let retention = DirectoryRetention::new(directory.path(), ReadableSize(5), 0).unwrap();
677
678        retention.register(LogFileKind::Default);
679        retention.initialize();
680        assert!(!retention.state.lock().initialized);
681
682        write_file(&old, b"old");
683        write_file(&active, b"new");
684        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
685        retention.reclaim(1);
686
687        assert!(!old.exists());
688        assert!(active.exists());
689        assert!(retention.state.lock().initialized);
690    }
691
692    #[cfg(unix)]
693    #[test]
694    fn test_retention_track_retries_initialization() {
695        let directory = TempDir::new().unwrap();
696        let active = directory.path().join("greptimedb.2026-01-01-00");
697        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
698
699        retention.register(LogFileKind::Default);
700        retention.initialize();
701        assert!(!retention.state.lock().initialized);
702
703        write_file(&active, b"");
704        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
705        retention.track(LogFileKind::Default, 0);
706
707        let state = retention.state.lock();
708        assert!(state.initialized);
709        assert_eq!(state.files.active[&LogFileKind::Default].path, active);
710    }
711
712    #[cfg(unix)]
713    #[test]
714    fn test_retention_reconciles_after_remove_failure() {
715        let directory = TempDir::new().unwrap();
716        let old = directory.path().join("greptimedb.2026-01-01-00");
717        let active = directory.path().join("greptimedb.2026-01-01-01");
718        write_file(&old, b"old");
719        write_file(&active, b"new");
720
721        let retention = DirectoryRetention::new(directory.path(), ReadableSize(6), 0).unwrap();
722        register_default_kind(&retention, directory.path(), &active);
723        retention.initialize();
724        fs::remove_file(&old).unwrap();
725
726        retention.reclaim(1);
727
728        let state = retention.state.lock();
729        assert!(state.files.closed.is_empty());
730        assert_eq!(state.files.total_size, 3);
731    }
732
733    #[cfg(unix)]
734    #[test]
735    fn test_retention_enforces_max_log_files() {
736        let directory = TempDir::new().unwrap();
737        let oldest = directory.path().join("greptimedb.2026-01-01-00");
738        let old = directory.path().join("greptimedb.2026-01-01-01");
739        let active = directory.path().join("greptimedb.2026-01-01-02");
740        write_file(&oldest, b"oldest");
741        write_file(&old, b"old");
742        write_file(&active, b"active");
743
744        let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
745        register_default_kind(&retention, directory.path(), &active);
746        retention.initialize();
747
748        assert!(!oldest.exists());
749        assert!(old.exists());
750        assert!(active.exists());
751    }
752
753    #[cfg(unix)]
754    #[test]
755    fn test_retention_enforces_max_log_files_after_rotation() {
756        let directory = TempDir::new().unwrap();
757        let oldest = directory.path().join("greptimedb.2026-01-01-00");
758        let old = directory.path().join("greptimedb.2026-01-01-01");
759        let active = directory.path().join("greptimedb.2026-01-01-02");
760        write_file(&oldest, b"oldest");
761        write_file(&old, b"old");
762
763        let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
764        register_default_kind(&retention, directory.path(), &old);
765        retention.initialize();
766        write_file(&active, b"");
767        fs::remove_file(directory.path().join(".greptimedb.current")).unwrap();
768        symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
769        retention.track(LogFileKind::Default, 1);
770
771        assert!(!oldest.exists());
772        assert!(old.exists());
773        assert!(active.exists());
774    }
775
776    #[test]
777    fn test_unmanaged_files_are_ignored() {
778        let directory = TempDir::new().unwrap();
779        let unmanaged = directory.path().join("keep-me");
780        write_file(&unmanaged, b"unmanaged");
781        write_file(
782            &directory.path().join("greptimedb.2026-01-01-00"),
783            b"managed",
784        );
785
786        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
787        retention.initialize();
788
789        assert!(unmanaged.exists());
790    }
791
792    #[cfg(unix)]
793    #[test]
794    fn test_managed_file_symlink_is_ignored() {
795        let directory = TempDir::new().unwrap();
796        let target = directory.path().join("target");
797        write_file(&target, b"target");
798
799        let link = directory.path().join("greptimedb.2026-01-01-00");
800        symlink(&target, &link).unwrap();
801
802        let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
803        retention.initialize();
804
805        assert!(link.exists());
806    }
807}