use std::path::Path;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use tracing::{info, warn};
use crate::github_path::derive_remote_repo;
use crate::log_drain::{
DEFAULT_MAX_FILE_BYTES, DEFAULT_MAX_WIRE_BYTES, DestinationUri, DrainConfig, DrainTarget,
Level, LogSource, ObjectStoreDestination, run_once,
};
pub const DEFAULT_INTERVAL_SECS: u64 = 900;
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct SingleSourceSection {
#[serde(default)]
pub enabled: Option<bool>,
#[serde(default)]
pub destination: Option<String>,
#[serde(default)]
pub interval_secs: Option<u64>,
#[serde(default)]
pub max_file_bytes: Option<u64>,
#[serde(default)]
pub max_wire_bytes: Option<u64>,
#[serde(default)]
pub secrets: Vec<String>,
#[serde(default)]
pub owner: Option<String>,
#[serde(default)]
pub project: Option<String>,
}
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
#[non_exhaustive]
pub enum SingleSourceError {
#[error("log_drain.enabled is true but log_drain.destination is unset")]
MissingDestination,
#[error("log_drain.destination is invalid: {reason}")]
Destination {
reason: String,
},
#[error("log_drain.{field} must be greater than zero")]
NonPositive {
field: &'static str,
},
#[error(
"log_drain has no owner/project: {reason} — set owner/project in the \
log-drain config, or run this daemon inside a git checkout whose \
origin names the project"
)]
MissingIdentity {
reason: String,
},
}
#[derive(Debug, Clone)]
pub enum LogDrainSetting {
Disabled,
Enabled(Box<ResolvedLogDrain>),
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct ResolvedLogDrain {
pub destination: DestinationUri,
pub destination_display: String,
pub target: DrainTarget,
pub source: LogSource,
pub interval: Duration,
pub max_file_bytes: u64,
pub max_wire_bytes: u64,
pub secrets: Vec<String>,
}
pub fn resolve_single_source(
section: &SingleSourceSection,
log_root: &Path,
project_root: Option<&Path>,
crate_name: &str,
default_include: &str,
) -> Result<LogDrainSetting, SingleSourceError> {
let enabled = section.enabled.unwrap_or(false);
let destination = match section.destination.as_deref().map(str::trim) {
Some(raw) if !raw.is_empty() => Some((
raw.to_string(),
DestinationUri::parse(raw).map_err(|e| SingleSourceError::Destination {
reason: e.to_string(),
})?,
)),
_ => None,
};
let interval_secs = section.interval_secs.unwrap_or(DEFAULT_INTERVAL_SECS);
if interval_secs == 0 {
return Err(SingleSourceError::NonPositive {
field: "interval_secs",
});
}
let max_file_bytes = section.max_file_bytes.unwrap_or(DEFAULT_MAX_FILE_BYTES);
if max_file_bytes == 0 {
return Err(SingleSourceError::NonPositive {
field: "max_file_bytes",
});
}
let max_wire_bytes = section.max_wire_bytes.unwrap_or(DEFAULT_MAX_WIRE_BYTES);
if max_wire_bytes == 0 {
return Err(SingleSourceError::NonPositive {
field: "max_wire_bytes",
});
}
if !enabled {
return Ok(LogDrainSetting::Disabled);
}
let Some((display, uri)) = destination else {
return Err(SingleSourceError::MissingDestination);
};
let target = resolve_identity(section, project_root)?;
let source = LogSource {
crate_name: crate_name.to_string(),
root: log_root.to_path_buf(),
include: vec![default_include.to_string()],
level_filter: Some(Level::Info),
};
Ok(LogDrainSetting::Enabled(Box::new(ResolvedLogDrain {
destination: uri,
destination_display: display,
target,
source,
interval: Duration::from_secs(interval_secs),
max_file_bytes,
max_wire_bytes,
secrets: section.secrets.clone(),
})))
}
fn resolve_identity(
section: &SingleSourceSection,
project_root: Option<&Path>,
) -> Result<DrainTarget, SingleSourceError> {
match (
non_empty(section.owner.as_deref()),
non_empty(section.project.as_deref()),
) {
(Some(owner), Some(project)) => return Ok(DrainTarget { owner, project }),
(Some(_), None) | (None, Some(_)) => {
return Err(SingleSourceError::MissingIdentity {
reason: "`owner` is set but `project` is not (or vice versa)".to_string(),
});
}
(None, None) => {}
}
let Some(root) = project_root else {
return Err(SingleSourceError::MissingIdentity {
reason: "no project is bound to this daemon".to_string(),
});
};
derive_remote_repo(root)
.map(|remote| DrainTarget {
owner: remote.owner,
project: remote.repo,
})
.map_err(|e| SingleSourceError::MissingIdentity {
reason: e.to_string(),
})
}
fn non_empty(raw: Option<&str>) -> Option<String> {
raw.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_string)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DrainOutcome {
Success,
SkippedDisabled,
Failed,
}
pub async fn run_tick(
section: &SingleSourceSection,
log_root: &Path,
project_root: Option<&Path>,
state_dir: &Path,
crate_name: &str,
default_include: &str,
) -> (DrainOutcome, Duration) {
let default_interval = Duration::from_secs(DEFAULT_INTERVAL_SECS);
match resolve_single_source(section, log_root, project_root, crate_name, default_include) {
Err(e) => {
warn!("log_drain: {e}");
(DrainOutcome::Failed, default_interval)
}
Ok(LogDrainSetting::Disabled) => (DrainOutcome::SkippedDisabled, default_interval),
Ok(LogDrainSetting::Enabled(plan)) => {
let interval = plan.interval;
let dest = match ObjectStoreDestination::connect(&plan.destination).await {
Ok(dest) => dest,
Err(e) => {
warn!(
destination = %plan.destination_display,
"log_drain: cannot reach destination: {e}"
);
return (DrainOutcome::Failed, interval);
}
};
let drain_cfg = DrainConfig::new(state_dir.to_path_buf())
.with_secrets(plan.secrets.clone())
.with_max_file_bytes(plan.max_file_bytes)
.with_max_wire_bytes(plan.max_wire_bytes);
let outcome = match run_once(
&drain_cfg,
&dest,
&plan.target,
std::slice::from_ref(&plan.source),
)
.await
{
Ok(report) if report.errors.is_empty() => {
info!(
uploaded = report.uploaded,
skipped_unchanged = report.skipped_unchanged,
"log_drain: pass complete"
);
DrainOutcome::Success
}
Ok(report) => {
warn!(
errors = report.errors.len(),
first = ?report.errors.first(),
"log_drain: pass completed with per-file errors"
);
DrainOutcome::Failed
}
Err(e) => {
warn!("log_drain: run_once failed: {e}");
DrainOutcome::Failed
}
};
(outcome, interval)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
const TEST_CRATE: &str = "test-crate";
const TEST_INCLUDE: &str = "test.log*";
fn base_section() -> SingleSourceSection {
SingleSourceSection::default()
}
#[test]
fn absent_section_is_disabled() {
let dir = tempfile::tempdir().expect("tempdir");
let setting =
resolve_single_source(&base_section(), dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.expect("resolve");
assert!(matches!(setting, LogDrainSetting::Disabled));
}
#[test]
fn validates_a_malformed_destination_even_while_disabled() {
let dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
destination: Some("ftp://nope".to_string()),
..base_section()
};
let err = resolve_single_source(§ion, dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.unwrap_err();
assert!(matches!(err, SingleSourceError::Destination { .. }));
}
#[test]
fn rejects_a_zero_interval() {
let dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
interval_secs: Some(0),
..base_section()
};
let err = resolve_single_source(§ion, dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.unwrap_err();
assert_eq!(
err,
SingleSourceError::NonPositive {
field: "interval_secs"
}
);
}
#[test]
fn enabled_with_no_destination_is_missing_destination() {
let dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
enabled: Some(true),
..base_section()
};
let err = resolve_single_source(§ion, dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.unwrap_err();
assert_eq!(err, SingleSourceError::MissingDestination);
}
#[test]
fn enabled_with_no_identity_and_no_project_is_missing_identity() {
let dir = tempfile::tempdir().expect("tempdir");
let dest_dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
enabled: Some(true),
destination: Some(format!("file://{}", dest_dir.path().display())),
..base_section()
};
let err = resolve_single_source(§ion, dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.unwrap_err();
assert!(matches!(err, SingleSourceError::MissingIdentity { .. }));
}
#[test]
fn enabled_with_explicit_owner_and_project_resolves() {
let dir = tempfile::tempdir().expect("tempdir");
let dest_dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
enabled: Some(true),
destination: Some(format!("file://{}", dest_dir.path().display())),
owner: Some("acme".to_string()),
project: Some("widgets".to_string()),
..base_section()
};
let setting = resolve_single_source(§ion, dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.expect("resolve");
let LogDrainSetting::Enabled(plan) = setting else {
panic!("expected Enabled");
};
assert_eq!(plan.target.owner, "acme");
assert_eq!(plan.target.project, "widgets");
assert_eq!(plan.source.crate_name, TEST_CRATE);
assert_eq!(plan.source.root, dir.path());
}
#[test]
fn half_an_identity_is_refused_not_filled_from_git() {
let dir = tempfile::tempdir().expect("tempdir");
let dest_dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
enabled: Some(true),
destination: Some(format!("file://{}", dest_dir.path().display())),
owner: Some("acme".to_string()),
..base_section()
};
let err = resolve_single_source(§ion, dir.path(), None, TEST_CRATE, TEST_INCLUDE)
.unwrap_err();
assert!(matches!(err, SingleSourceError::MissingIdentity { .. }));
}
#[tokio::test]
async fn opt_out_skips_the_pass() {
let log_root = tempfile::tempdir().expect("tempdir");
std::fs::write(log_root.path().join("test.log"), "hello\n").expect("seed log");
let state_dir = tempfile::tempdir().expect("tempdir");
let (outcome, interval) = run_tick(
&SingleSourceSection::default(),
log_root.path(),
None,
state_dir.path(),
TEST_CRATE,
TEST_INCLUDE,
)
.await;
assert_eq!(outcome, DrainOutcome::SkippedDisabled);
assert_eq!(interval, Duration::from_secs(DEFAULT_INTERVAL_SECS));
}
#[tokio::test]
async fn an_enabled_tick_uploads_into_the_injected_state_dir_and_is_idempotent() {
let log_root = tempfile::tempdir().expect("tempdir");
std::fs::write(log_root.path().join("test.log"), "hello\n").expect("seed log");
let dest_dir = tempfile::tempdir().expect("tempdir");
let state_dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
enabled: Some(true),
destination: Some(format!("file://{}", dest_dir.path().display())),
owner: Some("acme".to_string()),
project: Some("widgets".to_string()),
..SingleSourceSection::default()
};
let (first, _) = run_tick(
§ion,
log_root.path(),
None,
state_dir.path(),
TEST_CRATE,
TEST_INCLUDE,
)
.await;
assert_eq!(first, DrainOutcome::Success);
let uploaded = std::fs::read_dir(dest_dir.path().join("acme").join("widgets"))
.expect("read the target's key prefix")
.count();
assert!(
uploaded > 0,
"run_once must have written under acme/widgets"
);
assert!(
std::fs::read_dir(state_dir.path())
.expect("read state_dir")
.count()
> 0,
"the manifest cache must land under the injected state_dir"
);
let (second, _) = run_tick(
§ion,
log_root.path(),
None,
state_dir.path(),
TEST_CRATE,
TEST_INCLUDE,
)
.await;
assert_eq!(
second,
DrainOutcome::Success,
"an unchanged file must still report success on the next tick"
);
}
#[tokio::test]
async fn an_unreachable_destination_reports_failed() {
let log_root = tempfile::tempdir().expect("tempdir");
let state_dir = tempfile::tempdir().expect("tempdir");
let blocker = tempfile::NamedTempFile::new().expect("tempfile");
let section = SingleSourceSection {
enabled: Some(true),
destination: Some(format!(
"file://{}",
blocker.path().join("nested").display()
)),
owner: Some("acme".to_string()),
project: Some("widgets".to_string()),
..SingleSourceSection::default()
};
let (outcome, _) = run_tick(
§ion,
log_root.path(),
None,
state_dir.path(),
TEST_CRATE,
TEST_INCLUDE,
)
.await;
assert_eq!(outcome, DrainOutcome::Failed);
}
#[tokio::test]
async fn a_config_error_reports_failed() {
let log_root = tempfile::tempdir().expect("tempdir");
let state_dir = tempfile::tempdir().expect("tempdir");
let section = SingleSourceSection {
destination: Some("ftp://nope".to_string()),
..SingleSourceSection::default()
};
let (outcome, interval) = run_tick(
§ion,
log_root.path(),
None,
state_dir.path(),
TEST_CRATE,
TEST_INCLUDE,
)
.await;
assert_eq!(outcome, DrainOutcome::Failed);
assert_eq!(interval, Duration::from_secs(DEFAULT_INTERVAL_SECS));
}
}