diff --git a/Cargo.lock b/Cargo.lock
index fb277f0b0b..00e3d85442 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -2918,6 +2918,7 @@ dependencies = [
"prometheus 0.14.0",
"serde",
"serde_json",
+ "tempfile",
"tokio",
"tracing",
"tracing-appender",
@@ -14597,6 +14598,12 @@ dependencies = [
"symbolic-common",
]
+[[package]]
+name = "symlink"
+version = "0.1.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "a7973cce6668464ea31f176d85b13c7ab3bba2cb3b77a2ed26abd7801688010a"
+
[[package]]
name = "syn"
version = "1.0.109"
@@ -15774,12 +15781,13 @@ dependencies = [
[[package]]
name = "tracing-appender"
-version = "0.2.3"
+version = "0.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "3566e8ce28cc0a3fe42519fc80e6b4c943cc4c8cef275620eb8dac2d3d4e06cf"
+checksum = "050686193eb999b4bb3bc2acfa891a13da00f79734704c4b8b4ef1a10b368a3c"
dependencies = [
"crossbeam-channel",
- "thiserror 1.0.69",
+ "symlink",
+ "thiserror 2.0.17",
"time",
"tracing-subscriber",
]
diff --git a/config/config.md b/config/config.md
index 173b99f71b..a8d12c3557 100644
--- a/config/config.md
+++ b/config/config.md
@@ -232,6 +232,7 @@
| `logging.enable_file_logging` | Bool | `true` | Whether to write logs to files in `dir`. |
| `logging.log_format` | String | `text` | The log format. Can be `text`/`json`. |
| `logging.max_log_files` | Integer | `720` | The maximum amount of log files. |
+| `logging.max_log_dir_size` | String | `0B` | The maximum total size of managed log files in `dir`.
Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable. |
| `logging.enable_per_region_metrics` | Bool | `false` | Whether to enable per-region metrics.
Default to false. |
| `logging.otlp_export_protocol` | String | `http` | The OTLP tracing export protocol. Can be `grpc`/`http`. |
| `logging.otlp_headers` | -- | -- | Additional OTLP headers, only valid when using OTLP http |
@@ -375,6 +376,7 @@
| `logging.enable_file_logging` | Bool | `true` | Whether to write logs to files in `dir`. |
| `logging.log_format` | String | `text` | The log format. Can be `text`/`json`. |
| `logging.max_log_files` | Integer | `720` | The maximum amount of log files. |
+| `logging.max_log_dir_size` | String | `0B` | The maximum total size of managed log files in `dir`.
Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable. |
| `logging.enable_per_region_metrics` | Bool | `false` | Whether to enable per-region metrics.
Default to false. |
| `logging.otlp_export_protocol` | String | `http` | The OTLP tracing export protocol. Can be `grpc`/`http`. |
| `logging.otlp_headers` | -- | -- | Additional OTLP headers, only valid when using OTLP http |
@@ -491,6 +493,7 @@
| `logging.enable_file_logging` | Bool | `true` | Whether to write logs to files in `dir`. |
| `logging.log_format` | String | `text` | The log format. Can be `text`/`json`. |
| `logging.max_log_files` | Integer | `720` | The maximum amount of log files. |
+| `logging.max_log_dir_size` | String | `0B` | The maximum total size of managed log files in `dir`.
Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable. |
| `logging.otlp_export_protocol` | String | `http` | The OTLP tracing export protocol. Can be `grpc`/`http`. |
| `logging.otlp_headers` | -- | -- | Additional OTLP headers, only valid when using OTLP http |
| `logging.tracing_sample_ratio` | -- | Unset | The percentage of tracing will be sampled and exported.
Valid range `[0, 1]`, 1 means all traces are sampled, 0 means all traces are not sampled, the default value is 1.
ratio > 1 are treated as 1. Fractions < 0 are treated as 0 |
@@ -680,6 +683,7 @@
| `logging.enable_file_logging` | Bool | `true` | Whether to write logs to files in `dir`. |
| `logging.log_format` | String | `text` | The log format. Can be `text`/`json`. |
| `logging.max_log_files` | Integer | `720` | The maximum amount of log files. |
+| `logging.max_log_dir_size` | String | `0B` | The maximum total size of managed log files in `dir`.
Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable. |
| `logging.otlp_export_protocol` | String | `http` | The OTLP tracing export protocol. Can be `grpc`/`http`. |
| `logging.otlp_headers` | -- | -- | Additional OTLP headers, only valid when using OTLP http |
| `logging.tracing_sample_ratio` | -- | Unset | The percentage of tracing will be sampled and exported.
Valid range `[0, 1]`, 1 means all traces are sampled, 0 means all traces are not sampled, the default value is 1.
ratio > 1 are treated as 1. Fractions < 0 are treated as 0 |
@@ -741,6 +745,7 @@
| `logging.enable_file_logging` | Bool | `true` | Whether to write logs to files in `dir`. |
| `logging.log_format` | String | `text` | The log format. Can be `text`/`json`. |
| `logging.max_log_files` | Integer | `720` | The maximum amount of log files. |
+| `logging.max_log_dir_size` | String | `0B` | The maximum total size of managed log files in `dir`.
Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable. |
| `logging.otlp_export_protocol` | String | `http` | The OTLP tracing export protocol. Can be `grpc`/`http`. |
| `logging.otlp_headers` | -- | -- | Additional OTLP headers, only valid when using OTLP http |
| `logging.tracing_sample_ratio` | -- | Unset | The percentage of tracing will be sampled and exported.
Valid range `[0, 1]`, 1 means all traces are sampled, 0 means all traces are not sampled, the default value is 1.
ratio > 1 are treated as 1. Fractions < 0 are treated as 0 |
diff --git a/config/datanode.example.toml b/config/datanode.example.toml
index decc69dbf5..d378d55cc2 100644
--- a/config/datanode.example.toml
+++ b/config/datanode.example.toml
@@ -771,6 +771,10 @@ log_format = "text"
## The maximum amount of log files.
max_log_files = 720
+## The maximum total size of managed log files in `dir`.
+## Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable.
+max_log_dir_size = "0B"
+
## The OTLP tracing export protocol. Can be `grpc`/`http`.
otlp_export_protocol = "http"
diff --git a/config/flownode.example.toml b/config/flownode.example.toml
index 66caf2430a..192c6af697 100644
--- a/config/flownode.example.toml
+++ b/config/flownode.example.toml
@@ -123,6 +123,10 @@ log_format = "text"
## The maximum amount of log files.
max_log_files = 720
+## The maximum total size of managed log files in `dir`.
+## Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable.
+max_log_dir_size = "0B"
+
## The OTLP tracing export protocol. Can be `grpc`/`http`.
otlp_export_protocol = "http"
diff --git a/config/frontend.example.toml b/config/frontend.example.toml
index e85773c52e..7c6862916e 100644
--- a/config/frontend.example.toml
+++ b/config/frontend.example.toml
@@ -375,6 +375,10 @@ log_format = "text"
## The maximum amount of log files.
max_log_files = 720
+## The maximum total size of managed log files in `dir`.
+## Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable.
+max_log_dir_size = "0B"
+
## Whether to enable per-region metrics.
## Default to false.
enable_per_region_metrics = false
diff --git a/config/metasrv.example.toml b/config/metasrv.example.toml
index f2dbe53ba5..514993a4ca 100644
--- a/config/metasrv.example.toml
+++ b/config/metasrv.example.toml
@@ -380,6 +380,10 @@ log_format = "text"
## The maximum amount of log files.
max_log_files = 720
+## The maximum total size of managed log files in `dir`.
+## Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable.
+max_log_dir_size = "0B"
+
## The OTLP tracing export protocol. Can be `grpc`/`http`.
otlp_export_protocol = "http"
diff --git a/config/standalone.example.toml b/config/standalone.example.toml
index 10d6876b9c..b05cf0177d 100644
--- a/config/standalone.example.toml
+++ b/config/standalone.example.toml
@@ -913,6 +913,10 @@ log_format = "text"
## The maximum amount of log files.
max_log_files = 720
+## The maximum total size of managed log files in `dir`.
+## Old closed log files are removed before writing when necessary. Active files may exceed it. Set to `0B` to disable.
+max_log_dir_size = "0B"
+
## Whether to enable per-region metrics.
## Default to false.
enable_per_region_metrics = false
diff --git a/src/common/telemetry/Cargo.toml b/src/common/telemetry/Cargo.toml
index ae25a43c5d..1e146a5a75 100644
--- a/src/common/telemetry/Cargo.toml
+++ b/src/common/telemetry/Cargo.toml
@@ -37,3 +37,6 @@ tracing-appender.workspace = true
tracing-log = "0.2"
tracing-opentelemetry.workspace = true
tracing-subscriber.workspace = true
+
+[dev-dependencies]
+tempfile.workspace = true
diff --git a/src/common/telemetry/src/logging.rs b/src/common/telemetry/src/logging.rs
index 2f25dfc55e..83ceb7aef9 100644
--- a/src/common/telemetry/src/logging.rs
+++ b/src/common/telemetry/src/logging.rs
@@ -13,13 +13,17 @@
// limitations under the License.
//! logging stuffs, inspired by databend
+mod file_retention;
+
use std::collections::HashMap;
use std::env;
use std::io::IsTerminal;
use std::sync::{Arc, Mutex, Once, RwLock};
use std::time::Duration;
+use common_base::readable_size::ReadableSize;
use common_base::serde::empty_string_as_default;
+use file_retention::{DirectoryRetention, LogFileKind, build_file_appender};
use once_cell::sync::{Lazy, OnceCell};
use opentelemetry::trace::TracerProvider;
use opentelemetry::{KeyValue, global};
@@ -31,7 +35,6 @@ use serde::{Deserialize, Serialize};
use tracing::callsite;
use tracing::metadata::LevelFilter;
use tracing_appender::non_blocking::WorkerGuard;
-use tracing_appender::rolling::{RollingFileAppender, Rotation};
use tracing_log::LogTracer;
use tracing_subscriber::filter::{FilterFn, Targets};
use tracing_subscriber::fmt::Layer;
@@ -253,6 +256,9 @@ pub struct LoggingOptions {
/// The maximum number of log files set by default.
pub max_log_files: usize,
+ /// The maximum total size of managed log files in `dir`. Zero disables size-based retention.
+ pub max_log_dir_size: ReadableSize,
+
/// Whether to append logs to stdout. Default is true.
pub append_stdout: bool,
@@ -365,6 +371,7 @@ impl Default for LoggingOptions {
enable_file_logging: true,
// Rotation hourly, 24 files per day, keeps info log files of 30 days
max_log_files: 720,
+ max_log_dir_size: ReadableSize::default(),
otlp_export_protocol: None,
otlp_headers: HashMap::new(),
enable_per_region_metrics: false,
@@ -461,19 +468,16 @@ pub fn init_global_logging(
let file_logging_enabled = opts.enable_file_logging && !opts.dir.is_empty();
+ let retention = file_logging_enabled
+ .then(|| {
+ DirectoryRetention::new(opts.dir.clone(), opts.max_log_dir_size, opts.max_log_files)
+ })
+ .flatten();
+
// Configure the file logging layer with rolling policy.
let file_logging_layer = if file_logging_enabled {
- let rolling_appender = RollingFileAppender::builder()
- .rotation(Rotation::HOURLY)
- .filename_prefix("greptimedb")
- .max_log_files(opts.max_log_files)
- .build(&opts.dir)
- .unwrap_or_else(|e| {
- panic!(
- "initializing rolling file appender at {} failed: {}",
- &opts.dir, e
- )
- });
+ let rolling_appender =
+ build_file_appender(opts, LogFileKind::Default, retention.as_ref());
let (writer, guard) = tracing_appender::non_blocking(rolling_appender);
guards.push(guard);
@@ -494,17 +498,8 @@ pub fn init_global_logging(
// Configure the error file logging layer with rolling policy.
let err_file_logging_layer = if file_logging_enabled {
- let rolling_appender = RollingFileAppender::builder()
- .rotation(Rotation::HOURLY)
- .filename_prefix("greptimedb-err")
- .max_log_files(opts.max_log_files)
- .build(&opts.dir)
- .unwrap_or_else(|e| {
- panic!(
- "initializing rolling file appender at {} failed: {}",
- &opts.dir, e
- )
- });
+ let rolling_appender =
+ build_file_appender(opts, LogFileKind::Error, retention.as_ref());
let (writer, guard) = tracing_appender::non_blocking(rolling_appender);
guards.push(guard);
@@ -530,7 +525,12 @@ pub fn init_global_logging(
None
};
- let slow_query_logging_layer = build_slow_query_logger(opts, slow_query_opts, &mut guards);
+ let slow_query_logging_layer =
+ build_slow_query_logger(opts, slow_query_opts, retention.as_ref(), &mut guards);
+
+ if let Some(retention) = &retention {
+ retention.initialize();
+ }
// resolve log level settings from:
// - options from command line or config files
@@ -709,6 +709,7 @@ fn build_otlp_exporter(opts: &LoggingOptions) -> SpanExporter {
fn build_slow_query_logger(
opts: &LoggingOptions,
slow_query_opts: Option<&SlowQueryOptions>,
+ retention: Option<&DirectoryRetention>,
guards: &mut Vec,
) -> Option + Send + Sync + 'static>>
where
@@ -723,17 +724,7 @@ where
&& slow_query_opts.enable
&& slow_query_opts.record_type == SlowQueriesRecordType::Log
{
- let rolling_appender = RollingFileAppender::builder()
- .rotation(Rotation::HOURLY)
- .filename_prefix("greptimedb-slow-queries")
- .max_log_files(opts.max_log_files)
- .build(&opts.dir)
- .unwrap_or_else(|e| {
- panic!(
- "initializing rolling file appender at {} failed: {}",
- &opts.dir, e
- )
- });
+ let rolling_appender = build_file_appender(opts, LogFileKind::SlowQuery, retention);
let (writer, guard) = tracing_appender::non_blocking(rolling_appender);
guards.push(guard);
@@ -795,6 +786,14 @@ mod tests {
assert!(!opts.enable_file_logging);
}
+ #[test]
+ fn test_logging_options_deserialization_max_log_dir_size() {
+ let json = r#"{"max_log_dir_size": "1MiB"}"#;
+ let opts: LoggingOptions = serde_json::from_str(json).unwrap();
+
+ assert_eq!(opts.max_log_dir_size, ReadableSize::mb(1));
+ }
+
#[test]
fn test_logging_options_deserialization_empty_log_format() {
let json = r#"{"log_format": ""}"#;
@@ -889,6 +888,10 @@ mod tests {
max_log_files.max_log_files += 1;
assert_ne!(base, max_log_files);
+ let mut max_log_dir_size = base.clone();
+ max_log_dir_size.max_log_dir_size = ReadableSize::mb(1);
+ assert_ne!(base, max_log_dir_size);
+
let mut otlp_export_protocol = base.clone();
otlp_export_protocol.otlp_export_protocol = Some(OtlpExportProtocol::Http);
assert_ne!(base, otlp_export_protocol);
diff --git a/src/common/telemetry/src/logging/file_retention.rs b/src/common/telemetry/src/logging/file_retention.rs
new file mode 100644
index 0000000000..49cc976db3
--- /dev/null
+++ b/src/common/telemetry/src/logging/file_retention.rs
@@ -0,0 +1,807 @@
+// Copyright 2023 Greptime Team
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//! File-log retention based on the total size of managed log files.
+
+use std::collections::hash_map::Entry;
+use std::collections::{BTreeMap, HashMap, HashSet};
+use std::fs;
+use std::io::{self, Write};
+use std::path::PathBuf;
+use std::sync::Arc;
+use std::time::SystemTime;
+
+use common_base::readable_size::ReadableSize;
+use parking_lot::Mutex;
+use tracing_appender::rolling::{RollingFileAppender, Rotation};
+
+use crate::logging::LoggingOptions;
+
+/// A managed file-log kind.
+#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)]
+pub(crate) enum LogFileKind {
+ Default,
+ Error,
+ SlowQuery,
+}
+
+impl LogFileKind {
+ fn prefix(self) -> &'static str {
+ match self {
+ Self::Default => "greptimedb",
+ Self::Error => "greptimedb-err",
+ Self::SlowQuery => "greptimedb-slow-queries",
+ }
+ }
+
+ fn current_link_name(self) -> &'static str {
+ match self {
+ Self::Default => ".greptimedb.current",
+ Self::Error => ".greptimedb-err.current",
+ Self::SlowQuery => ".greptimedb-slow-queries.current",
+ }
+ }
+
+ fn kind_from_file_name(name: &str) -> Option {
+ [Self::SlowQuery, Self::Error, Self::Default]
+ .into_iter()
+ .find_map(|kind| {
+ name.strip_prefix(kind.prefix())
+ .and_then(|suffix| suffix.strip_prefix('.'))
+ .filter(|suffix| Self::is_hourly_suffix(suffix))
+ .map(|_| kind)
+ })
+ }
+
+ fn is_hourly_suffix(suffix: &str) -> bool {
+ suffix.len() == 13
+ && suffix.bytes().enumerate().all(|(index, byte)| {
+ matches!(index, 4 | 7 | 10) && byte == b'-'
+ || !matches!(index, 4 | 7 | 10) && byte.is_ascii_digit()
+ })
+ }
+}
+
+#[derive(Clone, Debug)]
+struct LogFile {
+ kind: LogFileKind,
+ path: PathBuf,
+ last_modified: SystemTime,
+ size: u64,
+}
+
+#[derive(Default)]
+struct FileIndex {
+ active: HashMap,
+ closed: BTreeMap<(SystemTime, PathBuf), LogFile>,
+ total_size: u64,
+}
+
+impl FileIndex {
+ fn track(&mut self, kind: LogFileKind, path: PathBuf, size: u64, last_modified: SystemTime) {
+ let active = match self.active.entry(kind) {
+ Entry::Occupied(mut entry) if entry.get().path != path => {
+ let previous = entry.insert(LogFile {
+ kind,
+ path,
+ last_modified,
+ size: 0,
+ });
+ self.closed
+ .insert((previous.last_modified, previous.path.clone()), previous);
+ entry.into_mut()
+ }
+ Entry::Occupied(entry) => entry.into_mut(),
+ Entry::Vacant(entry) => entry.insert(LogFile {
+ kind,
+ path,
+ last_modified,
+ size: 0,
+ }),
+ };
+
+ active.size = active.size.saturating_add(size);
+ active.last_modified = last_modified;
+ self.total_size = self.total_size.saturating_add(size);
+ }
+
+ fn count(&self, kind: LogFileKind) -> usize {
+ usize::from(self.active.contains_key(&kind))
+ + self
+ .closed
+ .values()
+ .filter(|file| file.kind == kind)
+ .count()
+ }
+
+ fn remove(&mut self, key: &(SystemTime, PathBuf)) {
+ if let Some(file) = self.closed.remove(key) {
+ self.total_size = self.total_size.saturating_sub(file.size);
+ }
+ }
+}
+
+#[derive(Default)]
+struct RetentionState {
+ kinds: HashSet,
+ files: FileIndex,
+ initialized: bool,
+ cleanup_error_reported: bool,
+}
+
+/// Retains managed file logs within configured directory-size and file-count limits.
+#[derive(Clone)]
+pub(crate) struct DirectoryRetention {
+ directory: Arc,
+ max_size: u64,
+ max_log_files: usize,
+ state: Arc>,
+}
+
+impl DirectoryRetention {
+ /// Returns a retention manager when the configured limit is enabled.
+ pub(crate) fn new(
+ directory: impl Into,
+ max_size: ReadableSize,
+ max_log_files: usize,
+ ) -> Option {
+ (max_size.as_bytes() > 0).then(|| Self {
+ directory: Arc::new(directory.into()),
+ max_size: max_size.as_bytes(),
+ max_log_files,
+ state: Arc::new(Mutex::new(RetentionState::default())),
+ })
+ }
+
+ /// Loads the initial file state after all enabled file-log kinds are registered.
+ pub(crate) fn initialize(&self) {
+ let mut state = self.state.lock();
+ if !self.ensure_initialized(&mut state) {
+ return;
+ }
+
+ self.prune_files(&mut state);
+ self.prune_size(&mut state, 0);
+ }
+
+ fn register(&self, kind: LogFileKind) {
+ let mut state = self.state.lock();
+ debug_assert!(!state.initialized);
+ state.kinds.insert(kind);
+ }
+
+ fn reclaim(&self, incoming_size: u64) {
+ let mut state = self.state.lock();
+ if !self.ensure_initialized(&mut state) {
+ return;
+ }
+
+ self.prune_size(&mut state, incoming_size);
+ }
+
+ fn track(&self, kind: LogFileKind, written: u64) {
+ let latest_file = self.current(kind);
+ let last_modified = SystemTime::now();
+ let mut state = self.state.lock();
+ if !self.ensure_initialized(&mut state) {
+ return;
+ }
+
+ let path = match latest_file {
+ Ok(latest_file) => latest_file,
+ Err(_) => {
+ self.report_error(&mut state, "resolving the latest log file");
+ self.reconcile(&mut state);
+ return;
+ }
+ };
+
+ state.files.track(kind, path, written, last_modified);
+ self.prune_files(&mut state);
+ self.prune_size(&mut state, 0);
+ }
+
+ fn scan(&self, kinds: &HashSet) -> io::Result {
+ let active_paths = kinds
+ .iter()
+ .map(|kind| self.current(*kind).map(|path| (*kind, path)))
+ .collect::>>()?;
+ let mut files = FileIndex::default();
+
+ for entry in fs::read_dir(self.directory.as_ref())? {
+ let entry = entry?;
+ let file_name = entry.file_name();
+ let Some(file_name) = file_name.to_str() else {
+ continue;
+ };
+ let Some(kind) = LogFileKind::kind_from_file_name(file_name) else {
+ continue;
+ };
+
+ let path = entry.path();
+ let metadata = fs::symlink_metadata(&path)?;
+ if !metadata.is_file() {
+ continue;
+ }
+
+ let file = LogFile {
+ kind,
+ path: path.clone(),
+ last_modified: metadata.modified()?,
+ size: metadata.len(),
+ };
+ files.total_size = files.total_size.saturating_add(file.size);
+ if active_paths.get(&kind) == Some(&path) {
+ files.active.insert(kind, file);
+ } else {
+ files.closed.insert((file.last_modified, path), file);
+ }
+ }
+
+ if files.active.len() != active_paths.len() {
+ return Err(io::Error::new(
+ io::ErrorKind::NotFound,
+ "latest log symlink target does not exist",
+ ));
+ }
+
+ Ok(files)
+ }
+
+ /// Rebuilds the in-memory index after an unexpected filesystem error.
+ fn reconcile(&self, state: &mut RetentionState) -> bool {
+ match self.scan(&state.kinds) {
+ Ok(files) => {
+ state.files = files;
+ state.initialized = true;
+ state.cleanup_error_reported = false;
+ true
+ }
+ Err(_) => {
+ self.report_error(state, "loading log directory");
+ false
+ }
+ }
+ }
+
+ fn ensure_initialized(&self, state: &mut RetentionState) -> bool {
+ state.initialized || self.reconcile(state)
+ }
+
+ fn current(&self, kind: LogFileKind) -> io::Result {
+ let symlink = self.directory.join(kind.current_link_name());
+ let target = fs::read_link(&symlink)?;
+ let Some(file_name) = target.file_name().and_then(|name| name.to_str()) else {
+ return Err(io::Error::new(
+ io::ErrorKind::InvalidData,
+ format!("invalid latest log symlink {}", symlink.display()),
+ ));
+ };
+ let Some(candidate) = LogFileKind::kind_from_file_name(file_name) else {
+ return Err(io::Error::new(
+ io::ErrorKind::InvalidData,
+ format!("invalid latest log symlink {}", symlink.display()),
+ ));
+ };
+ if candidate != kind {
+ return Err(io::Error::new(
+ io::ErrorKind::InvalidData,
+ format!("invalid latest log symlink {}", symlink.display()),
+ ));
+ }
+
+ Ok(self.directory.join(file_name))
+ }
+
+ fn prune_files(&self, state: &mut RetentionState) {
+ if self.max_log_files == 0 {
+ return;
+ }
+
+ for kind in state.kinds.clone() {
+ while state.files.count(kind) > self.max_log_files {
+ let Some((key, entry)) = state
+ .files
+ .closed
+ .iter()
+ .find(|(_, file)| file.kind == kind)
+ .map(|(key, entry)| (key.clone(), entry.clone()))
+ else {
+ return;
+ };
+
+ if !self.delete(state, key, entry) {
+ return;
+ }
+ }
+ }
+ }
+
+ fn prune_size(&self, state: &mut RetentionState, incoming_size: u64) {
+ while state.files.total_size.saturating_add(incoming_size) > self.max_size {
+ let Some((key, entry)) = state
+ .files
+ .closed
+ .first_key_value()
+ .map(|(key, entry)| (key.clone(), entry.clone()))
+ else {
+ // The limit is intentionally soft: active files and a single
+ // oversized log record are never removed or truncated.
+ self.report_error(
+ state,
+ "log directory exceeds the configured limit but only active files remain",
+ );
+ return;
+ };
+
+ if !self.delete(state, key, entry) {
+ return;
+ }
+ }
+ }
+
+ fn delete(
+ &self,
+ state: &mut RetentionState,
+ key: (SystemTime, PathBuf),
+ file: LogFile,
+ ) -> bool {
+ match fs::remove_file(&file.path) {
+ Ok(()) => {
+ state.files.remove(&key);
+ state.cleanup_error_reported = false;
+ true
+ }
+ Err(error) => {
+ self.report_error(state, &format!("removing {}: {error}", file.path.display()));
+ if error.kind() == io::ErrorKind::NotFound {
+ self.reconcile(state);
+ }
+ false
+ }
+ }
+ }
+
+ #[allow(clippy::print_stderr)]
+ fn report_error(&self, state: &mut RetentionState, message: &str) {
+ if !state.cleanup_error_reported {
+ // Do not use tracing here: this writer is itself on the tracing path.
+ eprintln!(
+ "Failed to retain log directory {}: {message}",
+ self.directory.display()
+ );
+ state.cleanup_error_reported = true;
+ }
+ }
+}
+
+/// A [`RollingFileAppender`] with shared directory-size retention.
+pub(crate) struct RetentionAppender {
+ inner: RollingFileAppender,
+ retention: DirectoryRetention,
+ kind: LogFileKind,
+}
+
+impl RetentionAppender {
+ fn new(inner: RollingFileAppender, retention: DirectoryRetention, kind: LogFileKind) -> Self {
+ Self {
+ inner,
+ retention,
+ kind,
+ }
+ }
+}
+
+impl Write for RetentionAppender {
+ fn write(&mut self, buf: &[u8]) -> io::Result {
+ self.retention.reclaim(buf.len() as u64);
+ let written = self.inner.write(buf)?;
+ self.retention.track(self.kind, written as u64);
+ Ok(written)
+ }
+
+ fn flush(&mut self) -> io::Result<()> {
+ self.inner.flush()
+ }
+}
+
+/// Builds the existing hourly appender, optionally wrapped with retention.
+pub(crate) fn build_file_appender(
+ opts: &LoggingOptions,
+ kind: LogFileKind,
+ retention: Option<&DirectoryRetention>,
+) -> Box {
+ // Directory retention owns size and count pruning so its index remains authoritative.
+ let upstream_max_log_files = if retention.is_some() {
+ 0
+ } else {
+ opts.max_log_files
+ };
+ let mut builder = RollingFileAppender::builder()
+ .rotation(Rotation::HOURLY)
+ .filename_prefix(kind.prefix())
+ .max_log_files(upstream_max_log_files);
+ if retention.is_some() {
+ builder = builder.latest_symlink(kind.current_link_name());
+ }
+
+ let appender = builder.build(&opts.dir).unwrap_or_else(|error| {
+ panic!(
+ "initializing rolling file appender at {} failed: {}",
+ opts.dir, error
+ )
+ });
+
+ if let Some(retention) = retention {
+ retention.register(kind);
+ Box::new(RetentionAppender::new(appender, retention.clone(), kind))
+ } else {
+ Box::new(appender)
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use std::fs::{self, File};
+ use std::io::Write;
+ #[cfg(unix)]
+ use std::os::unix::fs::symlink;
+ use std::path::{Path, PathBuf};
+ use std::time::SystemTime;
+
+ use common_base::readable_size::ReadableSize;
+ use tempfile::TempDir;
+
+ use super::FileIndex;
+ use crate::logging::LoggingOptions;
+ use crate::logging::file_retention::{DirectoryRetention, LogFileKind, build_file_appender};
+
+ fn write_file(path: &Path, contents: &[u8]) {
+ let mut file = File::create(path).unwrap();
+ file.write_all(contents).unwrap();
+ }
+
+ #[cfg(unix)]
+ fn register_default_kind(retention: &DirectoryRetention, directory: &Path, active: &Path) {
+ retention.register(LogFileKind::Default);
+ symlink(active, directory.join(".greptimedb.current")).unwrap();
+ }
+
+ #[test]
+ fn test_disabled_retention() {
+ assert!(DirectoryRetention::new("/tmp", ReadableSize::default(), 0).is_none());
+ }
+
+ #[test]
+ fn test_recognizes_managed_file_name() {
+ assert!(LogFileKind::kind_from_file_name("greptimedb.2026-01-01-00").is_some());
+ assert!(LogFileKind::kind_from_file_name("greptimedb-err.2026-01-01-00").is_some());
+ assert!(
+ LogFileKind::kind_from_file_name("greptimedb-slow-queries.2026-01-01-00").is_some()
+ );
+ assert!(LogFileKind::kind_from_file_name("greptimedb.2026-01-01-00.1").is_none());
+ }
+
+ #[test]
+ fn test_file_index_ignores_missing_closed_file() {
+ let mut files = FileIndex::default();
+
+ files.remove(&(SystemTime::UNIX_EPOCH, PathBuf::from("missing")));
+
+ assert_eq!(files.total_size, 0);
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_file_appender_creates_current_link_before_first_write() {
+ let directory = TempDir::new().unwrap();
+ let opts = LoggingOptions {
+ dir: directory.path().display().to_string(),
+ ..Default::default()
+ };
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize::mb(1), 0).unwrap();
+
+ let _appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
+
+ let symlink = directory.path().join(".greptimedb.current");
+ assert!(symlink.is_symlink());
+ assert!(fs::read_link(symlink).unwrap().exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_file_appender_prunes_closed_files_by_size() {
+ let directory = TempDir::new().unwrap();
+ let oldest = directory.path().join("greptimedb.2026-01-01-00");
+ let old = directory.path().join("greptimedb.2026-01-01-01");
+ write_file(&oldest, &vec![b'a'; 2 * 1024]);
+ write_file(&old, &vec![b'b'; 2 * 1024]);
+
+ let opts = LoggingOptions {
+ dir: directory.path().display().to_string(),
+ ..Default::default()
+ };
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1024), 0).unwrap();
+ let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
+
+ retention.initialize();
+ appender.write_all(b"current").unwrap();
+
+ assert!(!oldest.exists());
+ assert!(!old.exists());
+ assert!(
+ fs::read_link(directory.path().join(".greptimedb.current"))
+ .unwrap()
+ .exists()
+ );
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_file_appender_prunes_oldest_closed_file_by_size() {
+ let directory = TempDir::new().unwrap();
+ let oldest = directory.path().join("greptimedb.2026-01-01-00");
+ let old = directory.path().join("greptimedb.2026-01-01-01");
+ write_file(&oldest, &vec![b'a'; 2 * 1024]);
+ write_file(&old, &vec![b'b'; 512]);
+
+ let opts = LoggingOptions {
+ dir: directory.path().display().to_string(),
+ ..Default::default()
+ };
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1024), 0).unwrap();
+ let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
+
+ retention.initialize();
+ appender.write_all(b"current").unwrap();
+
+ assert!(!oldest.exists());
+ assert!(old.exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_file_appender_prunes_closed_files_by_count() {
+ let directory = TempDir::new().unwrap();
+ let oldest = directory.path().join("greptimedb.2026-01-01-00");
+ let old = directory.path().join("greptimedb.2026-01-01-01");
+ write_file(&oldest, b"oldest");
+ write_file(&old, b"old");
+
+ let opts = LoggingOptions {
+ dir: directory.path().display().to_string(),
+ ..Default::default()
+ };
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
+ let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
+
+ retention.initialize();
+ appender.write_all(b"current").unwrap();
+
+ assert!(!oldest.exists());
+ assert!(old.exists());
+ assert!(
+ fs::read_link(directory.path().join(".greptimedb.current"))
+ .unwrap()
+ .exists()
+ );
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_file_appender_keeps_oversized_active_file() {
+ let directory = TempDir::new().unwrap();
+ let opts = LoggingOptions {
+ dir: directory.path().display().to_string(),
+ ..Default::default()
+ };
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
+ let mut appender = build_file_appender(&opts, LogFileKind::Default, Some(&retention));
+
+ retention.initialize();
+ appender.write_all(b"current").unwrap();
+
+ let active = fs::metadata(directory.path().join(".greptimedb.current")).unwrap();
+ assert!(active.len() > 1);
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_removes_closed_files() {
+ let directory = TempDir::new().unwrap();
+ let old = directory.path().join("greptimedb.2026-01-01-00");
+ let active = directory.path().join("greptimedb.2026-01-01-01");
+ write_file(&old, b"old");
+ write_file(&active, b"new");
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(4), 0).unwrap();
+ register_default_kind(&retention, directory.path(), &active);
+ retention.initialize();
+
+ assert!(!old.exists());
+ assert!(active.exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_keeps_active_file() {
+ let directory = TempDir::new().unwrap();
+ let active = directory.path().join("greptimedb.2026-01-01-00");
+ write_file(&active, b"active");
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
+ register_default_kind(&retention, directory.path(), &active);
+ retention.initialize();
+
+ assert!(active.exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_tracks_rotated_file_in_memory() {
+ let directory = TempDir::new().unwrap();
+ let old = directory.path().join("greptimedb.2026-01-01-00");
+ let active = directory.path().join("greptimedb.2026-01-01-01");
+ write_file(&old, b"old");
+ write_file(&active, b"");
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(3), 0).unwrap();
+ register_default_kind(&retention, directory.path(), &old);
+ retention.initialize();
+ fs::remove_file(directory.path().join(".greptimedb.current")).unwrap();
+ symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
+ retention.track(LogFileKind::Default, 1);
+
+ assert!(!old.exists());
+ assert!(active.exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_retries_initialization() {
+ let directory = TempDir::new().unwrap();
+ let old = directory.path().join("greptimedb.2026-01-01-00");
+ let active = directory.path().join("greptimedb.2026-01-01-01");
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(5), 0).unwrap();
+
+ retention.register(LogFileKind::Default);
+ retention.initialize();
+ assert!(!retention.state.lock().initialized);
+
+ write_file(&old, b"old");
+ write_file(&active, b"new");
+ symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
+ retention.reclaim(1);
+
+ assert!(!old.exists());
+ assert!(active.exists());
+ assert!(retention.state.lock().initialized);
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_track_retries_initialization() {
+ let directory = TempDir::new().unwrap();
+ let active = directory.path().join("greptimedb.2026-01-01-00");
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
+
+ retention.register(LogFileKind::Default);
+ retention.initialize();
+ assert!(!retention.state.lock().initialized);
+
+ write_file(&active, b"");
+ symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
+ retention.track(LogFileKind::Default, 0);
+
+ let state = retention.state.lock();
+ assert!(state.initialized);
+ assert_eq!(state.files.active[&LogFileKind::Default].path, active);
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_reconciles_after_remove_failure() {
+ let directory = TempDir::new().unwrap();
+ let old = directory.path().join("greptimedb.2026-01-01-00");
+ let active = directory.path().join("greptimedb.2026-01-01-01");
+ write_file(&old, b"old");
+ write_file(&active, b"new");
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(6), 0).unwrap();
+ register_default_kind(&retention, directory.path(), &active);
+ retention.initialize();
+ fs::remove_file(&old).unwrap();
+
+ retention.reclaim(1);
+
+ let state = retention.state.lock();
+ assert!(state.files.closed.is_empty());
+ assert_eq!(state.files.total_size, 3);
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_enforces_max_log_files() {
+ let directory = TempDir::new().unwrap();
+ let oldest = directory.path().join("greptimedb.2026-01-01-00");
+ let old = directory.path().join("greptimedb.2026-01-01-01");
+ let active = directory.path().join("greptimedb.2026-01-01-02");
+ write_file(&oldest, b"oldest");
+ write_file(&old, b"old");
+ write_file(&active, b"active");
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
+ register_default_kind(&retention, directory.path(), &active);
+ retention.initialize();
+
+ assert!(!oldest.exists());
+ assert!(old.exists());
+ assert!(active.exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_retention_enforces_max_log_files_after_rotation() {
+ let directory = TempDir::new().unwrap();
+ let oldest = directory.path().join("greptimedb.2026-01-01-00");
+ let old = directory.path().join("greptimedb.2026-01-01-01");
+ let active = directory.path().join("greptimedb.2026-01-01-02");
+ write_file(&oldest, b"oldest");
+ write_file(&old, b"old");
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize::gb(1), 2).unwrap();
+ register_default_kind(&retention, directory.path(), &old);
+ retention.initialize();
+ write_file(&active, b"");
+ fs::remove_file(directory.path().join(".greptimedb.current")).unwrap();
+ symlink(&active, directory.path().join(".greptimedb.current")).unwrap();
+ retention.track(LogFileKind::Default, 1);
+
+ assert!(!oldest.exists());
+ assert!(old.exists());
+ assert!(active.exists());
+ }
+
+ #[test]
+ fn test_unmanaged_files_are_ignored() {
+ let directory = TempDir::new().unwrap();
+ let unmanaged = directory.path().join("keep-me");
+ write_file(&unmanaged, b"unmanaged");
+ write_file(
+ &directory.path().join("greptimedb.2026-01-01-00"),
+ b"managed",
+ );
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
+ retention.initialize();
+
+ assert!(unmanaged.exists());
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn test_managed_file_symlink_is_ignored() {
+ let directory = TempDir::new().unwrap();
+ let target = directory.path().join("target");
+ write_file(&target, b"target");
+
+ let link = directory.path().join("greptimedb.2026-01-01-00");
+ symlink(&target, &link).unwrap();
+
+ let retention = DirectoryRetention::new(directory.path(), ReadableSize(1), 0).unwrap();
+ retention.initialize();
+
+ assert!(link.exists());
+ }
+}
diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs
index 8312d0d5aa..bf9c5ea488 100644
--- a/tests-integration/tests/http.rs
+++ b/tests-integration/tests/http.rs
@@ -2317,6 +2317,7 @@ read_preference = "Leader"
[logging]
max_log_files = 720
+max_log_dir_size = "0KiB"
append_stdout = true
enable_file_logging = true
enable_otlp_tracing = false