use std::{
collections::{HashMap, HashSet},
env,
fs::{self, File, OpenOptions},
io::{self, BufWriter, Read, Seek, SeekFrom, Write},
path::{Path, PathBuf},
sync::{
Arc, Mutex, OnceLock,
atomic::{AtomicBool, AtomicU64, Ordering},
mpsc,
mpsc::RecvTimeoutError,
},
thread,
time::Duration,
};
#[cfg(any(target_os = "linux", target_os = "macos"))]
use std::{
os::unix::{
io::{AsRawFd, FromRawFd, IntoRawFd, RawFd},
net::UnixStream,
},
process::{ChildStderr, ChildStdout, Command, Stdio},
};
use terminal_size::Width;
use tracing::debug;
use crate::{
config::EffectiveLogsConfig, error::LogsManagerError, runtime,
upgrade::HandoffLogPipe,
};
const SERVICE_LOG_THREAD: &str = "sysg-service-log";
const SERVICE_STDOUT_THREAD: &str = "sysg-service-stdout";
const SERVICE_STDERR_THREAD: &str = "sysg-service-stderr";
const CHILD_LOG_THREAD: &str = "sysg-child-log";
const SPAWN_LOG_DIR: &str = "spawn";
const MIN_HANDOFF_FD: RawFd = 3;
const LOG_HANDOFF_TIMEOUT: Duration = Duration::from_secs(5);
const LOG_HANDOFF_POLL_INTERVAL: Duration = Duration::from_millis(5);
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LogSection {
Running,
Offline,
}
impl LogSection {
pub fn label(self) -> &'static str {
match self {
Self::Running => "Running Services",
Self::Offline => "Offline Services",
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
enum LogStream {
Stdout,
Stderr,
Combined,
}
impl LogStream {
fn as_str(self) -> &'static str {
match self {
Self::Stdout => "stdout",
Self::Stderr => "stderr",
Self::Combined => "combined",
}
}
fn from_filter(kind: &str) -> Option<Self> {
match kind {
"stdout" => Some(Self::Stdout),
"stderr" => Some(Self::Stderr),
_ => None,
}
}
}
impl From<LogStream> for &'static str {
fn from(stream: LogStream) -> Self {
stream.as_str()
}
}
pub fn get_service_log_path(project: &str, service: &str) -> PathBuf {
resolve_combined_log_path(project, service)
}
pub fn supervisor_log_path() -> PathBuf {
runtime::log_dir().join("supervisor.log")
}
pub fn tail_service_log(project: &str, service: &str, n: usize) -> Vec<String> {
tail_service_log_after(project, service, n, None)
}
pub fn tail_service_log_since(
project: &str,
service: &str,
n: usize,
since: chrono::DateTime<chrono::Utc>,
) -> Vec<String> {
tail_service_log_after(project, service, n, Some(since))
}
fn tail_service_log_after(
project: &str,
service: &str,
n: usize,
since: Option<chrono::DateTime<chrono::Utc>>,
) -> Vec<String> {
use std::io::{Read, Seek, SeekFrom};
let path = get_service_log_path(project, service);
let Ok(mut file) = fs::File::open(&path) else {
return Vec::new();
};
let len = file.metadata().map(|m| m.len()).unwrap_or(0);
let window = 16 * 1024;
let start = len.saturating_sub(window);
if file.seek(SeekFrom::Start(start)).is_err() {
return Vec::new();
}
let mut buf = Vec::with_capacity(window as usize);
if file.read_to_end(&mut buf).is_err() {
return Vec::new();
}
let text = String::from_utf8_lossy(&strip_ansi(&buf)).into_owned();
diagnostic_log_lines(&text, n, since)
}
fn diagnostic_log_lines(
text: &str,
n: usize,
since: Option<chrono::DateTime<chrono::Utc>>,
) -> Vec<String> {
text.lines()
.filter(|line| !line.trim().is_empty())
.filter(|line| {
since.is_none_or(|since| {
line.split_once(' ')
.and_then(|(timestamp, _)| timestamp.parse().ok())
.is_some_and(|timestamp: chrono::DateTime<chrono::Utc>| {
timestamp >= since
})
})
})
.map(strip_log_line_prefix)
.collect::<Vec<_>>()
.into_iter()
.rev()
.take(n)
.rev()
.collect()
}
fn strip_log_line_prefix(line: &str) -> String {
let mut parts = line.splitn(3, ' ');
let (Some(first), Some(second), Some(rest)) =
(parts.next(), parts.next(), parts.next())
else {
return line.to_string();
};
let looks_like_timestamp = first.len() >= 20
&& first.ends_with('Z')
&& first.contains('T')
&& first.starts_with(|c: char| c.is_ascii_digit());
if looks_like_timestamp && matches!(second, "stdout" | "stderr") {
rest.to_string()
} else {
line.to_string()
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct PruneSummary {
pub removed_files: usize,
pub reclaimed_bytes: u64,
}
pub fn parse_byte_size(value: &str) -> Result<u64, LogsManagerError> {
let trimmed = value.trim();
let split = trimmed
.find(|c: char| !c.is_ascii_digit() && c != '.')
.unwrap_or(trimmed.len());
let (number, unit) = trimmed.split_at(split);
let number: f64 = number
.parse()
.map_err(|_| LogsManagerError::InvalidPruneArg(value.to_string()))?;
let multiplier = match unit.trim().to_ascii_lowercase().as_str() {
"" | "b" => 1.0,
"k" | "kb" => 1024.0,
"m" | "mb" => 1024.0 * 1024.0,
"g" | "gb" => 1024.0 * 1024.0 * 1024.0,
"t" | "tb" => 1024.0 * 1024.0 * 1024.0 * 1024.0,
_ => return Err(LogsManagerError::InvalidPruneArg(value.to_string())),
};
Ok((number * multiplier) as u64)
}
pub fn parse_age_seconds(value: &str) -> Result<u64, LogsManagerError> {
let trimmed = value.trim();
let split = trimmed
.find(|c: char| !c.is_ascii_digit())
.unwrap_or(trimmed.len());
let (number, unit) = trimmed.split_at(split);
let number: u64 = number
.parse()
.map_err(|_| LogsManagerError::InvalidPruneArg(value.to_string()))?;
let multiplier = match unit.trim().to_ascii_lowercase().as_str() {
"" | "s" => 1,
"m" => 60,
"h" => 60 * 60,
"d" => 24 * 60 * 60,
"w" => 7 * 24 * 60 * 60,
_ => return Err(LogsManagerError::InvalidPruneArg(value.to_string())),
};
Ok(number * multiplier)
}
fn contained_project_log_dir(
root: &Path,
project: &str,
) -> Result<PathBuf, LogsManagerError> {
let reject = || {
io::Error::new(
io::ErrorKind::InvalidInput,
format!("'{project}' does not name a single project"),
)
};
let mut components = Path::new(project).components();
let dir = match (components.next(), components.next()) {
(Some(std::path::Component::Normal(name)), None) => root.join(name),
_ => return Err(reject().into()),
};
if dir.exists() {
let resolved_root = root.canonicalize().unwrap_or_else(|_| root.to_path_buf());
let resolved_dir = dir.canonicalize().map_err(|_| reject())?;
if !resolved_dir.starts_with(&resolved_root) {
return Err(reject().into());
}
}
Ok(dir)
}
fn is_rotated_backup(file_name: &str) -> bool {
file_name.rsplit_once('.').is_some_and(|(stem, suffix)| {
stem.ends_with(".log") && suffix.parse::<usize>().is_ok()
})
}
fn managed_log_files(log_dir: &Path) -> Result<Vec<PathBuf>, LogsManagerError> {
let mut files = Vec::new();
for entry in fs::read_dir(log_dir)? {
let entry = entry?;
let file_type = entry.file_type()?;
if file_type.is_file() {
files.push(entry.path());
continue;
}
if !file_type.is_dir() || entry.file_name() == SPAWN_LOG_DIR {
continue;
}
for project_entry in fs::read_dir(entry.path())? {
let project_entry = project_entry?;
if project_entry.file_type()?.is_file() {
files.push(project_entry.path());
}
}
}
Ok(files)
}
pub fn prune_logs(
max_size: Option<&str>,
max_age: Option<&str>,
) -> Result<PruneSummary, LogsManagerError> {
let max_bytes = max_size.map(parse_byte_size).transpose()?;
let max_age_secs = max_age.map(parse_age_seconds).transpose()?;
let log_dir = runtime::log_dir();
if !log_dir.exists() {
return Ok(PruneSummary::default());
}
let mut backups: Vec<(PathBuf, u64, std::time::SystemTime)> = Vec::new();
for path in managed_log_files(&log_dir)? {
let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else {
continue;
};
if !is_rotated_backup(file_name) {
continue;
}
let metadata = fs::metadata(&path)?;
let modified = metadata.modified().unwrap_or(std::time::UNIX_EPOCH);
backups.push((path, metadata.len(), modified));
}
backups.sort_by_key(|(_, _, modified)| *modified);
let mut summary = PruneSummary::default();
if let Some(max_age_secs) = max_age_secs {
let now = std::time::SystemTime::now();
backups.retain(|(path, len, modified)| {
let age = now
.duration_since(*modified)
.map(|d| d.as_secs())
.unwrap_or(0);
if age > max_age_secs {
if fs::remove_file(path).is_ok() {
summary.removed_files += 1;
summary.reclaimed_bytes += len;
}
false
} else {
true
}
});
}
if let Some(max_bytes) = max_bytes {
let mut total: u64 = backups.iter().map(|(_, len, _)| *len).sum();
for (path, len, _) in &backups {
if total <= max_bytes {
break;
}
if fs::remove_file(path).is_ok() {
summary.removed_files += 1;
summary.reclaimed_bytes += len;
total = total.saturating_sub(*len);
}
}
}
Ok(summary)
}
pub fn parse_time_bound(
value: &str,
now: chrono::DateTime<chrono::Utc>,
) -> Result<chrono::DateTime<chrono::Utc>, LogsManagerError> {
let trimmed = value.trim();
if trimmed.is_empty() {
return Err(LogsManagerError::InvalidTimeBound(value.to_string()));
}
if let Ok(parsed) = chrono::DateTime::parse_from_rfc3339(trimmed) {
return Ok(parsed.with_timezone(&chrono::Utc));
}
if let Ok(date) = chrono::NaiveDate::parse_from_str(trimmed, "%Y-%m-%d")
&& let Some(midnight) = date.and_hms_opt(0, 0, 0)
{
return Ok(chrono::DateTime::from_naive_utc_and_offset(
midnight,
chrono::Utc,
));
}
let seconds = parse_age_seconds(trimmed)
.map_err(|_| LogsManagerError::InvalidTimeBound(value.to_string()))?;
Ok(now - chrono::Duration::seconds(seconds as i64))
}
#[derive(Clone, Default)]
pub struct LogFilter {
pub since: Option<chrono::DateTime<chrono::Utc>>,
pub until: Option<chrono::DateTime<chrono::Utc>>,
pub grep: Option<regex::Regex>,
pub all: bool,
}
impl LogFilter {
pub fn from_parts(
since: Option<&str>,
until: Option<&str>,
grep: Option<&str>,
all: bool,
now: chrono::DateTime<chrono::Utc>,
) -> Result<Self, LogsManagerError> {
let since = since
.map(|value| parse_time_bound(value, now))
.transpose()?;
let until = until
.map(|value| parse_time_bound(value, now))
.transpose()?;
let grep = grep
.map(|pattern| {
regex::Regex::new(pattern)
.map_err(|err| LogsManagerError::InvalidGrep(err.to_string()))
})
.transpose()?;
Ok(Self {
since,
until,
grep,
all,
})
}
pub fn has_content_filter(&self) -> bool {
self.since.is_some() || self.until.is_some() || self.grep.is_some()
}
pub fn is_noop(&self) -> bool {
!self.has_content_filter() && !self.all
}
fn matches(&self, line: &[u8]) -> bool {
if let Some(ts) = captured_line_timestamp(line) {
if let Some(since) = self.since
&& ts < since
{
return false;
}
if let Some(until) = self.until
&& ts > until
{
return false;
}
} else if self.since.is_some() || self.until.is_some() {
return false;
}
if let Some(pattern) = &self.grep {
let text = String::from_utf8_lossy(line);
if !pattern.is_match(&text) {
return false;
}
}
true
}
pub fn apply(&self, bytes: &[u8]) -> Vec<u8> {
if !self.has_content_filter() {
return bytes.to_vec();
}
bytes
.split_inclusive(|byte| *byte == b'\n')
.filter(|line| self.matches(line.trim_ascii_end()))
.flat_map(|line| line.iter().copied())
.collect()
}
}
fn captured_line_timestamp(line: &[u8]) -> Option<chrono::DateTime<chrono::Utc>> {
let text = std::str::from_utf8(line).ok()?;
let first = text.split(' ').next()?;
chrono::DateTime::parse_from_rfc3339(first)
.ok()
.map(|parsed| parsed.with_timezone(&chrono::Utc))
}
pub fn rotated_history_paths(active: &Path) -> Vec<PathBuf> {
let Some(parent) = active.parent() else {
return vec![active.to_path_buf()];
};
let Some(base_name) = active.file_name().and_then(|name| name.to_str()) else {
return vec![active.to_path_buf()];
};
let prefix = format!("{base_name}.");
let mut backups: Vec<(usize, PathBuf)> = Vec::new();
if let Ok(entries) = fs::read_dir(parent) {
for entry in entries.flatten() {
let path = entry.path();
if !path.is_file() {
continue;
}
let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else {
continue;
};
if let Some(suffix) = file_name.strip_prefix(&prefix)
&& let Ok(index) = suffix.parse::<usize>()
{
backups.push((index, path));
}
}
}
backups.sort_by_key(|(index, _)| std::cmp::Reverse(*index));
let mut paths: Vec<PathBuf> = backups.into_iter().map(|(_, path)| path).collect();
if active.exists() {
paths.push(active.to_path_buf());
}
if paths.is_empty() {
paths.push(active.to_path_buf());
}
paths
}
fn read_full_history(active: &Path) -> Result<Vec<u8>, LogsManagerError> {
let mut bytes = Vec::new();
for path in rotated_history_paths(active) {
match fs::read(&path) {
Ok(mut chunk) => bytes.append(&mut chunk),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => continue,
Err(err) => return Err(err.into()),
}
}
Ok(bytes)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
pub enum LogFormat {
#[default]
Text,
Raw,
Json,
}
pub const SERVICE_MARKER_PREFIX: &str = "\u{1e}sysg-service ";
pub fn service_marker_line(service: &str) -> Vec<u8> {
format!("{SERVICE_MARKER_PREFIX}{service}\n").into_bytes()
}
fn parse_service_marker(line: &str) -> Option<&str> {
line.strip_prefix(SERVICE_MARKER_PREFIX)
}
pub fn strip_ansi(bytes: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(bytes.len());
let mut iter = bytes.iter().copied().peekable();
while let Some(byte) = iter.next() {
if byte != 0x1b {
out.push(byte);
continue;
}
match iter.peek().copied() {
Some(b'[') => {
iter.next();
for follow in iter.by_ref() {
if (0x40..=0x7e).contains(&follow) {
break;
}
}
}
Some(b']') => {
iter.next();
while let Some(follow) = iter.next() {
if follow == 0x07 {
break;
}
if follow == 0x1b && iter.peek() == Some(&b'\\') {
iter.next();
break;
}
}
}
Some(_) => {
iter.next();
}
None => {}
}
}
out
}
struct CapturedLine<'a> {
timestamp: &'a str,
stream: &'a str,
message: &'a str,
}
fn parse_captured_line(line: &str) -> Option<CapturedLine<'_>> {
let mut parts = line.splitn(3, ' ');
let timestamp = parts.next()?;
chrono::DateTime::parse_from_rfc3339(timestamp).ok()?;
let stream = parts.next()?;
if !matches!(stream, "stdout" | "stderr" | "combined") {
return None;
}
let message = parts.next().unwrap_or("");
Some(CapturedLine {
timestamp,
stream,
message,
})
}
fn json_escape(value: &str) -> String {
let mut escaped = String::with_capacity(value.len() + 2);
for ch in value.chars() {
match ch {
'"' => escaped.push_str("\\\""),
'\\' => escaped.push_str("\\\\"),
'\n' => escaped.push_str("\\n"),
'\r' => escaped.push_str("\\r"),
'\t' => escaped.push_str("\\t"),
c if (c as u32) < 0x20 => {
escaped.push_str(&format!("\\u{:04x}", c as u32));
}
c => escaped.push(c),
}
}
escaped
}
pub struct LogWriter<W: Write> {
inner: W,
format: LogFormat,
strip_ansi: bool,
service: Option<String>,
pending: Vec<u8>,
}
impl<W: Write> LogWriter<W> {
pub fn new(
inner: W,
format: LogFormat,
strip_ansi: bool,
service: Option<String>,
) -> Self {
Self {
inner,
format,
strip_ansi,
service,
pending: Vec::new(),
}
}
fn render_line(&mut self, raw: &[u8]) -> std::io::Result<()> {
let bytes = if self.strip_ansi {
strip_ansi(raw)
} else {
raw.to_vec()
};
if let Ok(text) = std::str::from_utf8(&bytes)
&& let Some(service) = parse_service_marker(text)
{
self.service = Some(service.to_string());
return Ok(());
}
if matches!(self.format, LogFormat::Text) {
self.inner.write_all(&bytes)?;
self.inner.write_all(b"\n")?;
return Ok(());
}
let text = String::from_utf8_lossy(&bytes);
let parsed = parse_captured_line(&text);
match self.format {
LogFormat::Text => unreachable!(),
LogFormat::Raw => {
if let Some(parsed) = parsed {
self.inner.write_all(parsed.message.as_bytes())?;
self.inner.write_all(b"\n")?;
}
}
LogFormat::Json => {
if let Some(parsed) = parsed {
let service = self.service.as_deref().unwrap_or("");
let json = format!(
"{{\"ts\":\"{}\",\"stream\":\"{}\",\"service\":\"{}\",\"line\":\"{}\"}}\n",
json_escape(parsed.timestamp),
json_escape(parsed.stream),
json_escape(service),
json_escape(parsed.message),
);
self.inner.write_all(json.as_bytes())?;
}
}
}
Ok(())
}
}
impl<W: Write> Write for LogWriter<W> {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
if matches!(self.format, LogFormat::Text) && !self.strip_ansi {
return self.inner.write(buf);
}
self.pending.extend_from_slice(buf);
while let Some(pos) = self.pending.iter().position(|byte| *byte == b'\n') {
let line: Vec<u8> = self.pending.drain(..=pos).collect();
let line = &line[..line.len() - 1];
self.render_line(line)?;
}
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
if !self.pending.is_empty() {
let line = std::mem::take(&mut self.pending);
self.render_line(&line)?;
}
self.inner.flush()
}
}
pub fn get_log_path(project: &str, service: &str, kind: &str) -> PathBuf {
resolve_log_path(project, service, kind)
}
pub fn validate_service_name(service: &str) -> Result<(), LogsManagerError> {
let invalid = service.is_empty()
|| service.len() > 255
|| service == "."
|| service == ".."
|| service
.chars()
.any(|c| c == '/' || c == '\\' || c == '\0' || std::path::is_separator(c));
if invalid {
return Err(LogsManagerError::InvalidServiceName(service.to_string()));
}
Ok(())
}
fn assert_within_log_dir(path: &Path) -> Result<(), LogsManagerError> {
let log_dir = runtime::log_dir();
let parent = path.parent().unwrap_or(&log_dir);
let base = log_dir.canonicalize().unwrap_or(log_dir.clone());
let resolved = parent
.canonicalize()
.unwrap_or_else(|_| parent.to_path_buf());
if resolved == base || resolved.parent() == Some(base.as_path()) {
Ok(())
} else {
Err(LogsManagerError::InvalidServiceName(
path.display().to_string(),
))
}
}
fn project_log_dir(project: &str) -> PathBuf {
runtime::log_dir().join(project)
}
fn canonical_log_path(project: &str, service: &str, kind: &str) -> PathBuf {
let mut path = project_log_dir(project);
path.push(format!("{service}_{kind}.log"));
path
}
fn canonical_combined_log_path(project: &str, service: &str) -> PathBuf {
let mut path = project_log_dir(project);
path.push(format!("{service}.log"));
path
}
const LIVE_LOG_BUFFER_LIMIT: usize = 256 * 1024;
const PROJECT_LOG_CHANNEL_CAPACITY: usize = 4096;
const LOG_FOLLOW_POLL_INTERVAL: Duration = Duration::from_millis(250);
struct LiveLogEntry {
buffer: Vec<u8>,
}
impl LiveLogEntry {
fn new() -> Self {
Self { buffer: Vec::new() }
}
fn append(&mut self, chunk: &[u8]) {
self.buffer.extend_from_slice(chunk);
if self.buffer.len() > LIVE_LOG_BUFFER_LIMIT {
let overflow = self.buffer.len() - LIVE_LOG_BUFFER_LIMIT;
let cutoff = self.buffer[overflow..]
.iter()
.position(|byte| *byte == b'\n')
.map(|offset| overflow + offset + 1)
.unwrap_or(self.buffer.len());
self.buffer.drain(..cutoff);
}
}
}
type LiveLogKey = (String, String, String);
#[derive(Clone)]
struct ProjectLogChunk {
project: String,
service: String,
bytes: Vec<u8>,
}
struct ProjectLogSub {
id: u64,
projects: Vec<String>,
receiver: mpsc::Receiver<ProjectLogChunk>,
}
type ProjectLogSender = (u64, mpsc::SyncSender<ProjectLogChunk>);
type ProjectLogSubscribers = std::collections::HashMap<String, Vec<ProjectLogSender>>;
impl Drop for ProjectLogSub {
fn drop(&mut self) {
let mut subscribers = project_log_subscribers()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for project in &self.projects {
let empty = if let Some(project_subscribers) = subscribers.get_mut(project) {
project_subscribers.retain(|(id, _)| *id != self.id);
project_subscribers.is_empty()
} else {
false
};
if empty {
subscribers.remove(project);
}
}
}
}
fn live_log_registry()
-> &'static Mutex<std::collections::HashMap<LiveLogKey, LiveLogEntry>> {
static REGISTRY: OnceLock<
Mutex<std::collections::HashMap<LiveLogKey, LiveLogEntry>>,
> = OnceLock::new();
REGISTRY.get_or_init(|| Mutex::new(std::collections::HashMap::new()))
}
fn project_log_subscribers() -> &'static Mutex<ProjectLogSubscribers> {
static SUBSCRIBERS: OnceLock<Mutex<ProjectLogSubscribers>> = OnceLock::new();
SUBSCRIBERS.get_or_init(|| Mutex::new(std::collections::HashMap::new()))
}
fn append_live_log_chunk(project: &str, service: &str, stream: LogStream, chunk: &[u8]) {
let key = (
project.to_string(),
service.to_string(),
stream.as_str().to_string(),
);
let mut registry = live_log_registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let entry = registry.entry(key).or_insert_with(LiveLogEntry::new);
entry.append(chunk);
if stream == LogStream::Combined {
let mut subscribers = project_log_subscribers()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut failed = HashSet::new();
if let Some(project_subscribers) = subscribers.get_mut(project) {
let event = ProjectLogChunk {
project: project.to_string(),
service: service.to_string(),
bytes: chunk.to_vec(),
};
project_subscribers.retain(|(id, subscriber)| {
if subscriber.try_send(event.clone()).is_ok() {
true
} else {
failed.insert(*id);
false
}
});
}
if !failed.is_empty() {
for project_subscribers in subscribers.values_mut() {
project_subscribers.retain(|(id, _)| !failed.contains(id));
}
subscribers.retain(|_, project_subscribers| !project_subscribers.is_empty());
}
}
}
fn subscribe_project_logs(projects: &[String]) -> (Vec<ProjectLogChunk>, ProjectLogSub) {
static NEXT_ID: AtomicU64 = AtomicU64::new(1);
let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
let (tx, receiver) = mpsc::sync_channel(PROJECT_LOG_CHANNEL_CAPACITY);
let mut snapshot = Vec::new();
let registry = live_log_registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut subscribers = project_log_subscribers()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for ((project, service, stream), entry) in registry.iter() {
if stream == LogStream::Combined.as_str()
&& projects.iter().any(|candidate| candidate == project)
&& !entry.buffer.is_empty()
{
snapshot.push(ProjectLogChunk {
project: project.clone(),
service: service.clone(),
bytes: entry.buffer.clone(),
});
}
}
for project in projects {
subscribers
.entry(project.clone())
.or_default()
.push((id, tx.clone()));
}
snapshot.sort_unstable_by(|left, right| {
(&left.project, &left.service).cmp(&(&right.project, &right.service))
});
(
snapshot,
ProjectLogSub {
id,
projects: projects.to_vec(),
receiver,
},
)
}
pub fn clear_live_log(project: &str, service: &str) {
if let Ok(mut registry) = live_log_registry().lock() {
registry.retain(|(proj, name, _), _| proj != project || name != service);
}
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn socket_peer_disconnected(stream: &UnixStream) -> bool {
let fd = stream.as_raw_fd();
let mut byte = 0_u8;
let result = unsafe {
libc::recv(
fd,
&mut byte as *mut u8 as *mut libc::c_void,
1,
libc::MSG_PEEK | libc::MSG_DONTWAIT,
)
};
if result == 0 {
return true;
}
if result < 0 {
let err = std::io::Error::last_os_error();
return !matches!(
err.raw_os_error(),
Some(code) if code == libc::EAGAIN || code == libc::EWOULDBLOCK
);
}
false
}
fn normalize(name: &str) -> String {
name.chars()
.filter(|c| c.is_ascii_alphanumeric())
.flat_map(|c| c.to_lowercase())
.collect()
}
fn locate_existing_log(project: &str, service: &str, kind: &str) -> Option<PathBuf> {
let canonical = canonical_log_path(project, service, kind);
let directory = canonical.parent()?;
let needle = normalize(service);
let suffix = format!("_{kind}.log");
let entries = fs::read_dir(directory).ok()?;
for entry in entries.flatten() {
let path = entry.path();
if !path.is_file() {
continue;
}
let file_name = path.file_name()?.to_str()?;
if !file_name.ends_with(&suffix) {
continue;
}
if let Some(service_name) = file_name.strip_suffix(&suffix)
&& normalize(service_name) == needle
{
return Some(path);
}
}
None
}
fn locate_existing_combined_log(project: &str, service: &str) -> Option<PathBuf> {
let canonical = canonical_combined_log_path(project, service);
let directory = canonical.parent()?;
let needle = normalize(service);
let entries = fs::read_dir(directory).ok()?;
for entry in entries.flatten() {
let path = entry.path();
if !path.is_file() {
continue;
}
let file_name = path.file_name()?.to_str()?;
if file_name == "supervisor.log"
|| file_name.ends_with("_stdout.log")
|| file_name.ends_with("_stderr.log")
|| !file_name.ends_with(".log")
{
continue;
}
if let Some(service_name) = file_name.strip_suffix(".log")
&& normalize(service_name) == needle
{
return Some(path);
}
}
None
}
pub fn resolve_log_path(project: &str, service: &str, kind: &str) -> PathBuf {
let canonical = canonical_log_path(project, service, kind);
if canonical.exists() {
return canonical;
}
locate_existing_log(project, service, kind).unwrap_or(canonical)
}
fn resolve_combined_log_path(project: &str, service: &str) -> PathBuf {
let canonical = canonical_combined_log_path(project, service);
if canonical.exists() {
return canonical;
}
locate_existing_combined_log(project, service).unwrap_or(canonical)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum TailMode {
Follow,
OneShot,
}
fn resolve_tail_mode(mode: TailMode, filter: &LogFilter) -> TailMode {
if filter.all || filter.since.is_some() || filter.until.is_some() {
TailMode::OneShot
} else {
mode
}
}
impl TailMode {
fn current() -> Self {
match env::var("SYSTEMG_TAIL_MODE") {
Ok(value) if value.eq_ignore_ascii_case("oneshot") => TailMode::OneShot,
_ => TailMode::Follow,
}
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn configure_command(
self,
cmd: &mut Command,
lines: usize,
stdout_path: &Path,
stderr_path: &Path,
combined_path: &Path,
kind: Option<&str>,
) {
if lines == usize::MAX {
cmd.arg("-n").arg("+1");
} else {
cmd.arg("-n").arg(lines.to_string());
}
if matches!(self, TailMode::Follow) {
cmd.arg("-F");
}
match kind {
Some("stdout") => {
if combined_path.exists() {
cmd.arg(combined_path);
} else {
cmd.arg(stdout_path);
}
}
Some("stderr") => {
if combined_path.exists() {
cmd.arg(combined_path);
} else {
cmd.arg(stderr_path);
}
}
_ => {
if combined_path.exists() {
cmd.arg(combined_path);
} else {
cmd.arg(stdout_path).arg(stderr_path);
}
}
}
}
}
fn touch_log_file(path: &Path) {
if let Some(parent) = path.parent() {
let _ = fs::create_dir_all(parent);
}
let _ = OpenOptions::new().create(true).append(true).open(path);
}
fn truncate_log_file(path: &Path) -> Result<(), LogsManagerError> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(path)?;
Ok(())
}
fn remove_rotated_log_files(path: &Path) -> Result<(), LogsManagerError> {
let Some(parent) = path.parent() else {
return Ok(());
};
let Some(base_name) = path.file_name().and_then(|name| name.to_str()) else {
return Ok(());
};
let prefix = format!("{base_name}.");
if !parent.exists() {
return Ok(());
}
for entry in fs::read_dir(parent)? {
let entry_path = entry?.path();
if !entry_path.is_file() {
continue;
}
let Some(file_name) = entry_path.file_name().and_then(|name| name.to_str())
else {
continue;
};
if file_name
.strip_prefix(&prefix)
.is_some_and(|suffix| suffix.parse::<usize>().is_ok())
{
fs::remove_file(entry_path)?;
}
}
Ok(())
}
fn rotated_log_path(path: &Path, index: usize) -> PathBuf {
let mut rotated = path.as_os_str().to_os_string();
rotated.push(format!(".{index}"));
PathBuf::from(rotated)
}
fn rotate_log_file(path: &Path, max_files: usize) -> std::io::Result<()> {
if max_files == 0 {
match fs::remove_file(path) {
Ok(()) => {}
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
Err(err) => return Err(err),
}
return Ok(());
}
let oldest = rotated_log_path(path, max_files);
match fs::remove_file(&oldest) {
Ok(()) => {}
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
Err(err) => return Err(err),
}
for index in (1..max_files).rev() {
let from = rotated_log_path(path, index);
let to = rotated_log_path(path, index + 1);
if from.exists() {
fs::rename(from, to)?;
}
}
if path.exists() {
fs::rename(path, rotated_log_path(path, 1))?;
}
Ok(())
}
struct ActiveLogFile {
path: PathBuf,
file: BufWriter<File>,
active_len: u64,
settings: EffectiveLogsConfig,
}
impl ActiveLogFile {
fn open(path: PathBuf, settings: EffectiveLogsConfig) -> std::io::Result<Self> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let raw_file = OpenOptions::new().create(true).append(true).open(&path)?;
let active_len = raw_file.metadata().map(|meta| meta.len()).unwrap_or(0);
Ok(Self {
path,
file: BufWriter::new(raw_file),
active_len,
settings,
})
}
fn write_line(&mut self, line: &[u8]) -> std::io::Result<()> {
if self.settings.max_bytes > 0
&& self.active_len > 0
&& self.active_len.saturating_add(line.len() as u64) > self.settings.max_bytes
{
self.file.flush()?;
rotate_log_file(&self.path, self.settings.max_files)?;
let raw_file = OpenOptions::new()
.create(true)
.append(true)
.open(&self.path)?;
self.file = BufWriter::new(raw_file);
self.active_len = 0;
}
self.file.write_all(line)?;
self.active_len = self.active_len.saturating_add(line.len() as u64);
Ok(())
}
fn flush(&mut self) -> std::io::Result<()> {
self.file.flush()
}
}
#[derive(Clone)]
pub struct RotatingLogWriter {
inner: Arc<Mutex<ActiveLogFile>>,
}
impl RotatingLogWriter {
pub fn open(path: PathBuf, settings: EffectiveLogsConfig) -> std::io::Result<Self> {
Ok(Self {
inner: Arc::new(Mutex::new(ActiveLogFile::open(path, settings)?)),
})
}
}
impl Write for RotatingLogWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let payload = truncate_log_payload(buf);
let mut file = self
.inner
.lock()
.map_err(|_| std::io::Error::other("supervisor log writer poisoned"))?;
file.write_line(&payload)?;
file.flush()?;
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
let mut file = self
.inner
.lock()
.map_err(|_| std::io::Error::other("supervisor log writer poisoned"))?;
file.flush()
}
}
impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for RotatingLogWriter {
type Writer = RotatingLogWriter;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}
const MAX_LOG_LINE_BYTES: usize = 16 * 1024;
fn capture_timestamp() -> String {
chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Micros, true)
}
fn truncate_log_payload(line: &[u8]) -> Vec<u8> {
if line.len() <= MAX_LOG_LINE_BYTES {
return line.to_vec();
}
let dropped = line.len() - MAX_LOG_LINE_BYTES;
let mut boundary = MAX_LOG_LINE_BYTES;
while boundary > 0 && (line[boundary] & 0b1100_0000) == 0b1000_0000 {
boundary -= 1;
}
let mut truncated = line[..boundary].to_vec();
truncated.extend_from_slice(format!("…[truncated {dropped} bytes]").as_bytes());
truncated
}
fn format_captured_log_line(kind: &str, line: &[u8]) -> Vec<u8> {
let line = truncate_log_payload(line);
let line = String::from_utf8_lossy(&line);
format!("{} {} {}\n", capture_timestamp(), kind, line).into_bytes()
}
#[cfg(target_os = "linux")]
fn process_fds_present(pid: u32) -> bool {
let stdout_fd_path = format!("/proc/{pid}/fd/1");
let stderr_fd_path = format!("/proc/{pid}/fd/2");
let stdout_fd = Path::new(&stdout_fd_path);
let stderr_fd = Path::new(&stderr_fd_path);
stdout_fd.exists() || stderr_fd.exists()
}
fn resolve_tail_targets(
project: &str,
service_name: &str,
pid: Option<u32>,
) -> Result<(PathBuf, PathBuf, PathBuf), LogsManagerError> {
validate_service_name(service_name)?;
let stdout_path = resolve_log_path(project, service_name, "stdout");
let stderr_path = resolve_log_path(project, service_name, "stderr");
let combined_path = resolve_combined_log_path(project, service_name);
let stdout_exists = stdout_path.exists();
let stderr_exists = stderr_path.exists();
if !combined_path.exists() && !stdout_exists {
assert_within_log_dir(&stdout_path)?;
touch_log_file(&stdout_path);
}
if !combined_path.exists() && !stderr_exists {
assert_within_log_dir(&stderr_path)?;
touch_log_file(&stderr_path);
}
#[cfg(target_os = "linux")]
{
if let Some(pid_value) = pid
&& !(combined_path.exists()
|| stdout_exists
|| stderr_exists
|| process_fds_present(pid_value))
{
return Err(LogsManagerError::LogUnavailable(pid_value));
}
}
#[cfg(not(target_os = "linux"))]
let _ = pid;
Ok((stdout_path, stderr_path, combined_path))
}
fn write_log_header(
mut writer: impl Write,
service_name: &str,
pid: Option<u32>,
) -> Result<(), LogsManagerError> {
write_boxed_log_title(
&mut writer,
&LogManager::format_log_title(service_name, pid),
)
}
pub fn prefix_lines_with_service(chunk: &[u8], service_name: &str) -> Vec<u8> {
let prefix = format!("{service_name} | ");
let mut out = Vec::with_capacity(chunk.len() + prefix.len());
let mut at_line_start = true;
for &byte in chunk {
if at_line_start && byte != b'\n' {
out.extend_from_slice(prefix.as_bytes());
}
out.push(byte);
at_line_start = byte == b'\n';
}
out
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn write_project_log_event(
socket: &mut UnixStream,
event: &ProjectLogChunk,
lines: usize,
kind: Option<&str>,
filter: &LogFilter,
structured: bool,
backlog: bool,
) -> Result<(), LogsManagerError> {
let bytes = match kind {
Some(kind_name) => filter_captured_log_bytes(&event.bytes, kind_name),
None => event.bytes.clone(),
};
let bytes = filter.apply(&bytes);
let bytes = if backlog {
tail_log_bytes(&bytes, lines).to_vec()
} else {
bytes
};
if bytes.is_empty() {
return Ok(());
}
if structured {
socket.write_all(&service_marker_line(&event.service))?;
socket.write_all(&bytes)?;
} else {
socket.write_all(&prefix_lines_with_service(&bytes, &event.service))?;
}
socket.flush()?;
Ok(())
}
pub fn write_log_section_header(
mut writer: impl Write,
section: LogSection,
) -> Result<(), LogsManagerError> {
write_boxed_log_title(&mut writer, section.label())?;
writer.flush()?;
Ok(())
}
fn detect_log_terminal_width(default_width: usize) -> usize {
terminal_size::terminal_size()
.map(|(Width(width), _)| width as usize)
.unwrap_or(default_width)
.max(24)
}
fn truncate_log_title(title: &str, max_width: usize) -> String {
let title_width = title.chars().count();
if title_width <= max_width {
return title.to_string();
}
if max_width <= 3 {
return ".".repeat(max_width);
}
let visible_width = max_width.saturating_sub(3);
let mut truncated = title.chars().take(visible_width).collect::<String>();
truncated.push_str("...");
truncated
}
fn write_boxed_log_title(
mut writer: impl Write,
title: &str,
) -> Result<(), LogsManagerError> {
let terminal_width = detect_log_terminal_width(100);
write!(writer, "{}", format_boxed_log_title(title, terminal_width))?;
writer.flush()?;
Ok(())
}
fn format_boxed_log_title(title: &str, terminal_width: usize) -> String {
let inner_width = terminal_width.saturating_sub(2).max(1);
let title = truncate_log_title(title, inner_width);
let title_width = title.chars().count();
let left_padding = inner_width.saturating_sub(title_width) / 2;
let right_padding = inner_width.saturating_sub(title_width + left_padding);
format!(
"\n┌{}┐\n│{}{}{}│\n└{}┘\n\n",
"─".repeat(inner_width),
" ".repeat(left_padding),
title,
" ".repeat(right_padding),
"─".repeat(inner_width)
)
}
const LOG_TAIL_CHUNK_SIZE: u64 = 8192;
fn tail_log_bytes(bytes: &[u8], lines: usize) -> Vec<u8> {
if lines == 0 || bytes.is_empty() {
return Vec::new();
}
let mut index = bytes.len();
if bytes.last() == Some(&b'\n') {
index = index.saturating_sub(1);
}
let mut newlines_seen = 0usize;
while index > 0 {
index -= 1;
if bytes[index] == b'\n' {
newlines_seen += 1;
if newlines_seen == lines {
return bytes[index + 1..].to_vec();
}
}
}
bytes.to_vec()
}
fn tail_log_file(path: &Path, lines: usize) -> Result<Vec<u8>, LogsManagerError> {
if lines == 0 {
return Ok(Vec::new());
}
let mut file = File::open(path)?;
let mut remaining = file.metadata()?.len();
let mut bytes = Vec::new();
while remaining > 0 {
let chunk_len = remaining.min(LOG_TAIL_CHUNK_SIZE);
remaining -= chunk_len;
file.seek(SeekFrom::Start(remaining))?;
let mut chunk = vec![0_u8; chunk_len as usize];
file.read_exact(&mut chunk)?;
chunk.extend_from_slice(&bytes);
bytes = chunk;
if tail_log_bytes(&bytes, lines).len() < bytes.len() {
break;
}
}
Ok(tail_log_bytes(&bytes, lines))
}
fn captured_log_line_matches_kind(line: &[u8], kind: &str) -> bool {
let Some(stream) = LogStream::from_filter(kind) else {
return false;
};
let Some(first_space) = line.iter().position(|byte| *byte == b' ') else {
return false;
};
let rest = &line[first_space + 1..];
rest.strip_prefix(stream.as_str().as_bytes())
.is_some_and(|remaining| remaining.first() == Some(&b' '))
}
fn tail_log_file_filtered(
path: &Path,
lines: usize,
kind: &str,
) -> Result<Vec<u8>, LogsManagerError> {
if lines == 0 {
return Ok(Vec::new());
}
let mut file = File::open(path)?;
let mut remaining = file.metadata()?.len();
let mut bytes = Vec::new();
while remaining > 0 {
let chunk_len = remaining.min(LOG_TAIL_CHUNK_SIZE);
remaining -= chunk_len;
file.seek(SeekFrom::Start(remaining))?;
let mut chunk = vec![0_u8; chunk_len as usize];
file.read_exact(&mut chunk)?;
chunk.extend_from_slice(&bytes);
bytes = chunk;
let matching_count = bytes
.split(|byte| *byte == b'\n')
.filter(|line| captured_log_line_matches_kind(line, kind))
.count();
if matching_count > lines {
break;
}
}
let mut matching = bytes
.split_inclusive(|byte| *byte == b'\n')
.filter(|line| captured_log_line_matches_kind(line.trim_ascii_end(), kind))
.map(Vec::from)
.collect::<Vec<_>>();
if matching.len() > lines {
matching.drain(..matching.len() - lines);
}
Ok(matching.concat())
}
fn filter_captured_log_bytes(bytes: &[u8], kind: &str) -> Vec<u8> {
bytes
.split_inclusive(|byte| *byte == b'\n')
.filter(|line| captured_log_line_matches_kind(line.trim_ascii_end(), kind))
.flat_map(|line| line.iter().copied())
.collect()
}
fn line_matches_stream(line: &[u8], stream: Option<LogStream>) -> bool {
match stream {
Some(stream) => captured_log_line_matches_kind(line, stream.as_str()),
None => true,
}
}
fn follow_filtered_log_file(
mut writer: impl Write,
path: &Path,
lines: usize,
stream: Option<LogStream>,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
let initial = match stream {
Some(stream) => tail_log_file_filtered(path, lines, stream.as_str())?,
None => tail_log_file(path, lines)?,
};
writer.write_all(&filter.apply(&initial))?;
writer.flush()?;
let mut offset = fs::metadata(path)?.len();
let mut pending = Vec::new();
loop {
thread::sleep(Duration::from_millis(250));
let current_len = match fs::metadata(path) {
Ok(metadata) => metadata.len(),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
offset = 0;
pending.clear();
continue;
}
Err(err) => return Err(err.into()),
};
if current_len < offset {
offset = 0;
pending.clear();
}
if current_len == offset {
continue;
}
let mut file = File::open(path)?;
file.seek(SeekFrom::Start(offset))?;
let mut chunk = Vec::with_capacity((current_len - offset) as usize);
file.read_to_end(&mut chunk)?;
offset = current_len;
pending.extend_from_slice(&chunk);
while let Some(newline_pos) = pending.iter().position(|byte| *byte == b'\n') {
let line = pending.drain(..=newline_pos).collect::<Vec<_>>();
let trimmed = line.trim_ascii_end();
if line_matches_stream(trimmed, stream) && filter.matches(trimmed) {
writer.write_all(&line)?;
writer.flush()?;
}
}
}
}
fn write_log_file_tail(
mut writer: impl Write,
stdout_path: &Path,
stderr_path: &Path,
combined_path: &Path,
lines: usize,
kind: Option<&str>,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
for bytes in
collect_log_tail(stdout_path, stderr_path, combined_path, lines, kind, filter)?
{
writer.write_all(&bytes)?;
}
writer.flush()?;
Ok(())
}
fn collect_log_tail(
stdout_path: &Path,
stderr_path: &Path,
combined_path: &Path,
lines: usize,
kind: Option<&str>,
filter: &LogFilter,
) -> Result<Vec<Vec<u8>>, LogsManagerError> {
let stream_kind = kind.and_then(LogStream::from_filter);
if filter.all {
let mut chunks = Vec::new();
if combined_path.exists() {
let raw = read_full_history(combined_path)?;
let selected = match stream_kind {
Some(stream) => filter_captured_log_bytes(&raw, stream.as_str()),
None => raw,
};
chunks.push(filter.apply(&selected));
} else {
match stream_kind {
Some(LogStream::Stdout) => {
chunks.push(filter.apply(&read_full_history(stdout_path)?))
}
Some(LogStream::Stderr) => {
chunks.push(filter.apply(&read_full_history(stderr_path)?))
}
_ => {
chunks.push(filter.apply(&read_full_history(stdout_path)?));
chunks.push(filter.apply(&read_full_history(stderr_path)?));
}
}
}
return Ok(chunks);
}
let read_lines = if filter.has_content_filter() {
usize::MAX
} else {
lines
};
let mut chunks: Vec<Vec<u8>> = Vec::new();
match stream_kind {
Some(stream) => {
if combined_path.exists() {
chunks.push(tail_log_file_filtered(
combined_path,
read_lines,
stream.as_str(),
)?);
} else {
let single = match stream {
LogStream::Stdout => stdout_path,
_ => stderr_path,
};
chunks.push(tail_log_file(single, read_lines)?);
}
}
None => {
if combined_path.exists() {
chunks.push(tail_log_file(combined_path, read_lines)?);
} else {
chunks.push(tail_log_file(stdout_path, read_lines)?);
chunks.push(tail_log_file(stderr_path, read_lines)?);
}
}
}
if filter.has_content_filter() {
for chunk in &mut chunks {
let filtered = filter.apply(chunk);
*chunk = tail_log_bytes(&filtered, lines);
}
}
Ok(chunks)
}
fn write_forwarded_console_line(
mut writer: impl Write,
prefix: &str,
line: &str,
) -> std::io::Result<()> {
writeln!(writer, "{prefix}{line}")
}
fn forward_prefixed_line(service_label: &str, line: &[u8], echo_to_terminal: bool) {
let line = String::from_utf8_lossy(line);
if echo_to_terminal {
if let Err(err) = write_forwarded_console_line(
std::io::stderr(),
&format!("[{service_label}] "),
&line,
) {
eprintln!(
"Warning: Failed to write forwarded log for [{}]: {}",
service_label, err
);
}
} else {
debug!("[{service_label}] {line}");
}
}
fn flush_forwarded_lines(
pending: &mut Vec<u8>,
service_label: &str,
echo_to_terminal: bool,
) {
while let Some(newline_pos) = pending.iter().position(|byte| *byte == b'\n') {
let mut line = pending.drain(..=newline_pos).collect::<Vec<_>>();
if matches!(line.last(), Some(b'\n')) {
line.pop();
}
if matches!(line.last(), Some(b'\r')) {
line.pop();
}
forward_prefixed_line(service_label, &line, echo_to_terminal);
}
}
fn flush_remaining_forwarded_line(
pending: &mut Vec<u8>,
service_label: &str,
echo_to_terminal: bool,
) {
if pending.is_empty() {
return;
}
let line = std::mem::take(pending);
forward_prefixed_line(service_label, &line, echo_to_terminal);
}
struct ServiceLogLine {
stream: LogStream,
line: Vec<u8>,
}
enum ServiceLogMessage {
Line(ServiceLogLine),
Flush(mpsc::SyncSender<std::io::Result<()>>),
}
fn read_service_log_stream(
service_label: &str,
stream: LogStream,
mut reader: impl Read,
sender: mpsc::Sender<ServiceLogMessage>,
) -> std::io::Result<()> {
let mut buffer = [0_u8; 8192];
let mut pending = Vec::new();
let mut forward_pending = Vec::new();
let echo_to_terminal = false;
loop {
let bytes_read = reader.read(&mut buffer)?;
if bytes_read == 0 {
break;
}
let chunk = &buffer[..bytes_read];
pending.extend_from_slice(chunk);
forward_pending.extend_from_slice(chunk);
while let Some(newline_pos) = pending.iter().position(|byte| *byte == b'\n') {
let mut line = pending.drain(..=newline_pos).collect::<Vec<_>>();
if matches!(line.last(), Some(b'\n')) {
line.pop();
}
if matches!(line.last(), Some(b'\r')) {
line.pop();
}
let _ = sender.send(ServiceLogMessage::Line(ServiceLogLine { stream, line }));
}
flush_forwarded_lines(&mut forward_pending, service_label, echo_to_terminal);
}
if !pending.is_empty() {
let _ = sender.send(ServiceLogMessage::Line(ServiceLogLine {
stream,
line: pending.clone(),
}));
}
flush_remaining_forwarded_line(&mut forward_pending, service_label, echo_to_terminal);
Ok(())
}
fn write_service_log(
project: &str,
service_label: &str,
path: PathBuf,
receiver: mpsc::Receiver<ServiceLogMessage>,
settings: EffectiveLogsConfig,
) -> std::io::Result<()> {
let mut file = ActiveLogFile::open(path, settings)?;
for message in receiver {
match message {
ServiceLogMessage::Line(line) => {
let formatted =
format_captured_log_line(line.stream.as_str(), &line.line);
file.write_line(&formatted)?;
file.flush()?;
append_live_log_chunk(
project,
service_label,
LogStream::Combined,
&formatted,
);
}
ServiceLogMessage::Flush(reply) => match file.flush() {
Ok(()) => {
let _ = reply.send(Ok(()));
}
Err(err) => {
let reported = io::Error::new(err.kind(), err.to_string());
let _ = reply.send(Err(reported));
return Err(err);
}
},
}
}
file.flush()
}
fn stream_dynamic_child_log(
path: &Path,
owner_label: Option<&str>,
child_label: &str,
mut reader: impl Read,
echo_to_console: bool,
) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let mut file = OpenOptions::new().create(true).append(true).open(path)?;
let mut buffer = [0_u8; 8192];
let mut pending = Vec::new();
loop {
let bytes_read = reader.read(&mut buffer)?;
if bytes_read == 0 {
break;
}
let chunk = &buffer[..bytes_read];
file.write_all(chunk)?;
if echo_to_console {
pending.extend_from_slice(chunk);
while let Some(newline_pos) = pending.iter().position(|byte| *byte == b'\n') {
let mut line = pending.drain(..=newline_pos).collect::<Vec<_>>();
if matches!(line.last(), Some(b'\n')) {
line.pop();
}
if matches!(line.last(), Some(b'\r')) {
line.pop();
}
let owner = owner_label.unwrap_or("spawn");
println!(
"[{}:{}] {}",
owner,
child_label,
String::from_utf8_lossy(&line)
);
}
}
}
if echo_to_console && !pending.is_empty() {
let owner = owner_label.unwrap_or("spawn");
println!(
"[{}:{}] {}",
owner,
child_label,
String::from_utf8_lossy(&pending)
);
}
file.flush()
}
struct LogReaderState {
paused: AtomicBool,
pending: Mutex<Vec<u8>>,
}
struct RegisteredLogPipe {
id: u64,
writer_id: u64,
project: String,
service: String,
stream: LogStream,
settings: EffectiveLogsConfig,
owner: File,
state: Arc<LogReaderState>,
writer: mpsc::Sender<ServiceLogMessage>,
}
fn registered_log_pipes() -> &'static Mutex<Vec<RegisteredLogPipe>> {
static REGISTRY: OnceLock<Mutex<Vec<RegisteredLogPipe>>> = OnceLock::new();
REGISTRY.get_or_init(|| Mutex::new(Vec::new()))
}
fn log_handoff_paused() -> &'static AtomicBool {
static PAUSED: AtomicBool = AtomicBool::new(false);
&PAUSED
}
fn next_log_handoff_id() -> u64 {
static NEXT_ID: AtomicU64 = AtomicU64::new(1);
NEXT_ID.fetch_add(1, Ordering::Relaxed)
}
fn set_close_on_exec(fd: RawFd, enabled: bool) -> io::Result<()> {
let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
if flags < 0 {
return Err(io::Error::last_os_error());
}
let next = if enabled {
flags | libc::FD_CLOEXEC
} else {
flags & !libc::FD_CLOEXEC
};
if unsafe { libc::fcntl(fd, libc::F_SETFD, next) } < 0 {
return Err(io::Error::last_os_error());
}
Ok(())
}
fn set_nonblocking(fd: RawFd) -> io::Result<()> {
let flags = unsafe { libc::fcntl(fd, libc::F_GETFL) };
if flags < 0 {
return Err(io::Error::last_os_error());
}
if unsafe { libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) } < 0 {
return Err(io::Error::last_os_error());
}
Ok(())
}
fn duplicate_for_handoff(fd: RawFd) -> io::Result<File> {
let duplicate = unsafe { libc::fcntl(fd, libc::F_DUPFD_CLOEXEC, MIN_HANDOFF_FD) };
if duplicate < 0 {
return Err(io::Error::last_os_error());
}
Ok(unsafe { File::from_raw_fd(duplicate) })
}
fn spawn_canonical_service_writer(
project: &str,
service: &str,
settings: EffectiveLogsConfig,
) -> io::Result<(u64, mpsc::Sender<ServiceLogMessage>)> {
let path = get_service_log_path(project, service);
let project_label = project.to_string();
let service_label = service.to_string();
let (sender, receiver) = mpsc::channel();
thread::Builder::new()
.name(SERVICE_LOG_THREAD.into())
.spawn(move || {
if let Err(err) = write_service_log(
&project_label,
&service_label,
path.clone(),
receiver,
settings,
) {
eprintln!(
"Warning: Unable to write service log file at {:?}: {}",
path, err
);
}
})?;
Ok((next_log_handoff_id(), sender))
}
fn remove_registered_log_pipe(id: u64) {
if let Ok(mut registry) = registered_log_pipes().lock() {
registry.retain(|entry| entry.id != id);
}
}
fn read_registered_log_stream(
service_label: &str,
stream: LogStream,
mut reader: impl Read,
sender: mpsc::Sender<ServiceLogMessage>,
state: &LogReaderState,
) -> io::Result<()> {
let mut buffer = [0_u8; 8192];
let mut pending = state
.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
loop {
if log_handoff_paused().load(Ordering::Acquire) {
state.paused.store(true, Ordering::Release);
while log_handoff_paused().load(Ordering::Acquire) {
thread::sleep(LOG_HANDOFF_POLL_INTERVAL);
}
state.paused.store(false, Ordering::Release);
continue;
}
match reader.read(&mut buffer) {
Ok(0) => break,
Ok(bytes_read) => {
pending.extend_from_slice(&buffer[..bytes_read]);
while let Some(newline_pos) =
pending.iter().position(|byte| *byte == b'\n')
{
let mut line = pending.drain(..=newline_pos).collect::<Vec<_>>();
if matches!(line.last(), Some(b'\n')) {
line.pop();
}
if matches!(line.last(), Some(b'\r')) {
line.pop();
}
forward_prefixed_line(service_label, &line, false);
if sender
.send(ServiceLogMessage::Line(ServiceLogLine { stream, line }))
.is_err()
{
return Ok(());
}
}
*state
.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = pending.clone();
}
Err(err) if err.kind() == io::ErrorKind::WouldBlock => {
thread::sleep(LOG_HANDOFF_POLL_INTERVAL);
}
Err(err) if err.kind() == io::ErrorKind::Interrupted => continue,
Err(err) => return Err(err),
}
}
if !pending.is_empty() {
forward_prefixed_line(service_label, &pending, false);
let _ = sender.send(ServiceLogMessage::Line(ServiceLogLine {
stream,
line: pending,
}));
}
state
.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clear();
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn spawn_registered_log_reader<R>(
project: &str,
service: &str,
stream: LogStream,
reader: R,
pending: Vec<u8>,
settings: EffectiveLogsConfig,
writer_id: u64,
writer: mpsc::Sender<ServiceLogMessage>,
) -> io::Result<()>
where
R: Read + AsRawFd + Send + 'static,
{
set_nonblocking(reader.as_raw_fd())?;
set_close_on_exec(reader.as_raw_fd(), true)?;
let owner = duplicate_for_handoff(reader.as_raw_fd())?;
let id = next_log_handoff_id();
let state = Arc::new(LogReaderState {
paused: AtomicBool::new(false),
pending: Mutex::new(pending),
});
registered_log_pipes()
.lock()
.map_err(|_| io::Error::other("managed log pipe registry is poisoned"))?
.push(RegisteredLogPipe {
id,
writer_id,
project: project.to_string(),
service: service.to_string(),
stream,
settings,
owner,
state: Arc::clone(&state),
writer: writer.clone(),
});
let service_label = service.to_string();
let thread_name = match stream {
LogStream::Stdout => SERVICE_STDOUT_THREAD,
LogStream::Stderr => SERVICE_STDERR_THREAD,
LogStream::Combined => SERVICE_LOG_THREAD,
};
if let Err(err) = thread::Builder::new()
.name(thread_name.into())
.spawn(move || {
if let Err(err) =
read_registered_log_stream(&service_label, stream, reader, writer, &state)
{
eprintln!(
"Warning: Unable to read {} for [{}]: {}",
stream.as_str(),
service_label,
err
);
}
remove_registered_log_pipe(id);
})
{
remove_registered_log_pipe(id);
return Err(err);
}
Ok(())
}
pub fn spawn_managed_service_log_writers(
project: &str,
service: &str,
stdout: Option<ChildStdout>,
stderr: Option<ChildStderr>,
settings: EffectiveLogsConfig,
) -> io::Result<()> {
let (writer_id, writer) = spawn_canonical_service_writer(project, service, settings)?;
if let Some(stdout) = stdout {
spawn_registered_log_reader(
project,
service,
LogStream::Stdout,
stdout,
Vec::new(),
settings,
writer_id,
writer.clone(),
)?;
}
if let Some(stderr) = stderr {
spawn_registered_log_reader(
project,
service,
LogStream::Stderr,
stderr,
Vec::new(),
settings,
writer_id,
writer,
)?;
}
Ok(())
}
pub fn service_log_handoff_ready(project: &str, service: &str) -> bool {
let Ok(registry) = registered_log_pipes().lock() else {
return false;
};
[LogStream::Stdout, LogStream::Stderr]
.into_iter()
.all(|stream| {
registry.iter().any(|entry| {
entry.project == project
&& entry.service == service
&& entry.stream == stream
})
})
}
pub fn prepare_log_pipe_handoff() -> io::Result<Vec<HandoffLogPipe>> {
log_handoff_paused().store(true, Ordering::Release);
let deadline = std::time::Instant::now() + LOG_HANDOFF_TIMEOUT;
loop {
let all_paused = registered_log_pipes()
.lock()
.map_err(|_| io::Error::other("managed log pipe registry is poisoned"))?
.iter()
.all(|entry| entry.state.paused.load(Ordering::Acquire));
if all_paused {
break;
}
if std::time::Instant::now() >= deadline {
cancel_log_pipe_handoff();
return Err(io::Error::new(
io::ErrorKind::TimedOut,
"managed log readers did not pause for supervisor handoff",
));
}
thread::sleep(LOG_HANDOFF_POLL_INTERVAL);
}
let writers = registered_log_pipes()
.lock()
.map_err(|_| io::Error::other("managed log pipe registry is poisoned"))?
.iter()
.map(|entry| (entry.writer_id, entry.writer.clone()))
.collect::<HashMap<_, _>>();
for writer in writers.into_values() {
let (reply, response) = mpsc::sync_channel(1);
if writer.send(ServiceLogMessage::Flush(reply)).is_err() {
cancel_log_pipe_handoff();
return Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"managed service log writer stopped before handoff",
));
}
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
match response.recv_timeout(remaining) {
Ok(result) => result?,
Err(_) => {
cancel_log_pipe_handoff();
return Err(io::Error::new(
io::ErrorKind::TimedOut,
"managed service logs did not flush before handoff",
));
}
}
}
let mut registry = registered_log_pipes()
.lock()
.map_err(|_| io::Error::other("managed log pipe registry is poisoned"))?;
let mut handoff = Vec::with_capacity(registry.len());
for entry in registry.iter_mut() {
if let Err(err) = set_close_on_exec(entry.owner.as_raw_fd(), false) {
for retained in registry.iter() {
let _ = set_close_on_exec(retained.owner.as_raw_fd(), true);
}
log_handoff_paused().store(false, Ordering::Release);
return Err(err);
}
handoff.push(HandoffLogPipe {
project: entry.project.clone(),
service: entry.service.clone(),
stream: entry.stream.as_str().to_string(),
fd: entry.owner.as_raw_fd(),
pending: entry
.state
.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone(),
settings: entry.settings,
});
}
Ok(handoff)
}
pub fn cancel_log_pipe_handoff() {
if let Ok(registry) = registered_log_pipes().lock() {
for entry in registry.iter() {
let _ = set_close_on_exec(entry.owner.as_raw_fd(), true);
}
}
log_handoff_paused().store(false, Ordering::Release);
}
fn validate_log_pipe_handoff(pipes: &[HandoffLogPipe]) -> io::Result<Vec<LogStream>> {
let mut descriptors = HashSet::new();
let mut streams = HashSet::new();
let mut settings = HashMap::new();
let mut parsed = Vec::with_capacity(pipes.len());
for pipe in pipes {
if pipe.project.is_empty() || pipe.service.is_empty() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"handed-off log pipe has an empty project or service",
));
}
if pipe.fd < MIN_HANDOFF_FD || unsafe { libc::fcntl(pipe.fd, libc::F_GETFD) } < 0
{
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("handed-off log descriptor {} is invalid", pipe.fd),
));
}
if !descriptors.insert(pipe.fd) {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("handed-off log descriptor {} is duplicated", pipe.fd),
));
}
let stream = match pipe.stream.as_str() {
"stdout" => LogStream::Stdout,
"stderr" => LogStream::Stderr,
other => {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("unsupported handed-off log stream `{other}`"),
));
}
};
if !streams.insert((pipe.project.as_str(), pipe.service.as_str(), stream)) {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!(
"handed-off {} stream for `{}/{}` is duplicated",
pipe.stream, pipe.project, pipe.service
),
));
}
let key = (pipe.project.as_str(), pipe.service.as_str());
if let Some(previous) = settings.insert(key, pipe.settings)
&& previous != pipe.settings
{
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!(
"handed-off log settings for `{}/{}` are inconsistent",
pipe.project, pipe.service
),
));
}
parsed.push(stream);
}
Ok(parsed)
}
pub fn resume_log_pipe_handoff(pipes: &[HandoffLogPipe]) -> io::Result<()> {
if !registered_log_pipes()
.lock()
.map_err(|_| io::Error::other("managed log pipe registry is poisoned"))?
.is_empty()
{
return Err(io::Error::other(
"managed log pipe registry was not empty during resume",
));
}
let streams = validate_log_pipe_handoff(pipes)?;
let mut writers: HashMap<(String, String), (u64, mpsc::Sender<ServiceLogMessage>)> =
HashMap::new();
for (pipe, stream) in pipes.iter().zip(streams) {
let key = (pipe.project.clone(), pipe.service.clone());
let (writer_id, writer) = match writers.get(&key) {
Some((writer_id, writer)) => (*writer_id, writer.clone()),
None => {
let created = spawn_canonical_service_writer(
&pipe.project,
&pipe.service,
pipe.settings,
)?;
writers.insert(key, created.clone());
created
}
};
let reader = duplicate_for_handoff(pipe.fd)?;
spawn_registered_log_reader(
&pipe.project,
&pipe.service,
stream,
reader,
pipe.pending.clone(),
pipe.settings,
writer_id,
writer,
)?;
}
for pipe in pipes {
if unsafe { libc::close(pipe.fd) } < 0 {
return Err(io::Error::last_os_error());
}
}
Ok(())
}
pub fn spawn_log_writer(
project: &str,
service: &str,
reader: impl Read + Send + 'static,
kind: &str,
) -> io::Result<()> {
spawn_log_writer_with_config(
project,
service,
reader,
kind,
EffectiveLogsConfig::default(),
)
}
pub fn spawn_log_writer_with_config(
project: &str,
service: &str,
reader: impl Read + Send + 'static,
kind: &str,
settings: EffectiveLogsConfig,
) -> io::Result<()> {
let reader = Box::new(reader) as Box<dyn Read + Send>;
match LogStream::from_filter(kind) {
Some(LogStream::Stdout) => {
spawn_service_log_writers(project, service, Some(reader), None, settings)
}
Some(LogStream::Stderr) => {
spawn_service_log_writers(project, service, None, Some(reader), settings)
}
_ => spawn_service_log_writers(project, service, Some(reader), None, settings),
}
}
pub fn spawn_service_log_writers(
project: &str,
service: &str,
stdout: Option<Box<dyn Read + Send>>,
stderr: Option<Box<dyn Read + Send>>,
settings: EffectiveLogsConfig,
) -> io::Result<()> {
let path = get_service_log_path(project, service);
let project_label = project.to_string();
let service_label = service.to_string();
let (sender, receiver) = mpsc::channel();
{
let project_label = project_label.clone();
let service_label = service_label.clone();
let path = path.clone();
thread::Builder::new()
.name(SERVICE_LOG_THREAD.into())
.spawn(move || {
if let Err(err) = write_service_log(
&project_label,
&service_label,
path.clone(),
receiver,
settings,
) {
eprintln!(
"Warning: Unable to write service log file at {:?}: {}",
path, err
);
}
})?;
}
if let Some(stdout) = stdout {
let service_label = service_label.clone();
let sender = sender.clone();
thread::Builder::new()
.name(SERVICE_STDOUT_THREAD.into())
.spawn(move || {
if let Err(err) = read_service_log_stream(
&service_label,
LogStream::Stdout,
stdout,
sender,
) {
eprintln!(
"Warning: Unable to read stdout for [{}]: {}",
service_label, err
);
}
})?;
}
if let Some(stderr) = stderr {
let service_label = service_label.clone();
thread::Builder::new()
.name(SERVICE_STDERR_THREAD.into())
.spawn(move || {
if let Err(err) = read_service_log_stream(
&service_label,
LogStream::Stderr,
stderr,
sender,
) {
eprintln!(
"Warning: Unable to read stderr for [{}]: {}",
service_label, err
);
}
})?;
}
Ok(())
}
pub fn spawn_dynamic_child_log_writer(
root_service: Option<&str>,
child_name: &str,
pid: u32,
reader: impl Read + Send + 'static,
kind: &str,
echo_to_console: bool,
) -> io::Result<()> {
let owner_component = root_service
.map(normalize)
.filter(|s| !s.is_empty())
.unwrap_or_else(|| "dynamic".to_string());
let child_component = normalize(child_name);
let child_component = if child_component.is_empty() {
"child".to_string()
} else {
child_component
};
let mut path = runtime::log_dir();
path.push(SPAWN_LOG_DIR);
let file_name = format!(
"{}_{}_{}_{}.log",
owner_component, child_component, pid, kind
);
path.push(file_name);
let owner_label = root_service.map(str::to_string);
let child_label = child_name.to_string();
thread::Builder::new()
.name(CHILD_LOG_THREAD.into())
.spawn(move || {
if let Err(err) = stream_dynamic_child_log(
&path,
owner_label.as_deref(),
&child_label,
reader,
echo_to_console,
) {
eprintln!("Warning: Unable to write spawn log {:?}: {}", path, err);
}
})?;
Ok(())
}
#[derive(Default)]
pub struct LogManager {}
impl LogManager {
pub fn new() -> Self {
Self {}
}
pub fn collect_service_log(
&self,
project: &str,
service_name: &str,
lines: usize,
kind: Option<&str>,
filter: &LogFilter,
) -> Result<Vec<u8>, LogsManagerError> {
let stdout_path = resolve_log_path(project, service_name, "stdout");
let stderr_path = resolve_log_path(project, service_name, "stderr");
let combined_path = resolve_combined_log_path(project, service_name);
let mut bytes = Vec::new();
for chunk in collect_log_tail(
&stdout_path,
&stderr_path,
&combined_path,
lines,
kind,
filter,
)? {
bytes.extend_from_slice(&chunk);
}
Ok(bytes)
}
#[allow(clippy::too_many_arguments)]
pub fn show_log(
&self,
project: &str,
service_name: &str,
pid: u32,
lines: usize,
kind: Option<&str>,
show_header: bool,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
self.show_logs_platform_with_mode(
project,
service_name,
Some(pid),
lines,
kind,
resolve_tail_mode(TailMode::current(), filter),
show_header,
filter,
)
}
#[allow(clippy::too_many_arguments)]
pub fn show_log_snapshot(
&self,
project: &str,
service_name: &str,
pid: u32,
lines: usize,
kind: Option<&str>,
show_header: bool,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
self.show_logs_platform_with_mode(
project,
service_name,
Some(pid),
lines,
kind,
TailMode::OneShot,
show_header,
filter,
)
}
pub fn show_inactive_log(
&self,
project: &str,
service_name: &str,
lines: usize,
kind: Option<&str>,
show_header: bool,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
self.show_logs_platform_with_mode(
project,
service_name,
None,
lines,
kind,
resolve_tail_mode(TailMode::current(), filter),
show_header,
filter,
)
}
pub fn show_inactive_log_snapshot(
&self,
project: &str,
service_name: &str,
lines: usize,
kind: Option<&str>,
show_header: bool,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
self.show_logs_platform_with_mode(
project,
service_name,
None,
lines,
kind,
TailMode::OneShot,
show_header,
filter,
)
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[allow(clippy::too_many_arguments)]
pub fn stream_log_to_socket(
&self,
project: &str,
service_name: &str,
pid: Option<u32>,
lines: usize,
kind: Option<&str>,
follow: bool,
show_header: bool,
filter: &LogFilter,
stream: &UnixStream,
) -> Result<(), LogsManagerError> {
let mode = resolve_tail_mode(
if follow {
TailMode::Follow
} else {
TailMode::OneShot
},
filter,
);
self.stream_logs_platform_with_mode(
project,
service_name,
pid,
lines,
kind,
mode,
show_header,
filter,
stream,
)
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[allow(clippy::too_many_arguments)]
pub fn stream_project_logs_to_socket(
&self,
projects: &[String],
backlog_services: &std::collections::HashSet<(String, String)>,
lines: usize,
kind: Option<&str>,
follow: bool,
filter: &LogFilter,
stream: &UnixStream,
structured: bool,
) -> Result<(), LogsManagerError> {
let mode = resolve_tail_mode(
if follow {
TailMode::Follow
} else {
TailMode::OneShot
},
filter,
);
let (snapshot, subscription) = subscribe_project_logs(projects);
let mut socket = stream.try_clone()?;
for event in &snapshot {
if backlog_services.contains(&(event.project.clone(), event.service.clone()))
{
write_project_log_event(
&mut socket,
event,
lines,
kind,
filter,
structured,
true,
)?;
}
}
if mode == TailMode::OneShot {
return Ok(());
}
loop {
match subscription.receiver.recv_timeout(LOG_FOLLOW_POLL_INTERVAL) {
Ok(event) => write_project_log_event(
&mut socket,
&event,
lines,
kind,
filter,
structured,
false,
)?,
Err(RecvTimeoutError::Timeout) => {
if socket_peer_disconnected(&socket) {
break;
}
}
Err(RecvTimeoutError::Disconnected) => break,
}
}
Ok(())
}
pub fn clear_service_logs(
&self,
project: &str,
service_name: &str,
) -> Result<(), LogsManagerError> {
validate_service_name(service_name)?;
contained_project_log_dir(&runtime::log_dir(), project)?;
let stdout_path = resolve_log_path(project, service_name, "stdout");
let stderr_path = resolve_log_path(project, service_name, "stderr");
let combined_path = resolve_combined_log_path(project, service_name);
truncate_log_file(&stdout_path)?;
truncate_log_file(&stderr_path)?;
truncate_log_file(&combined_path)?;
remove_rotated_log_files(&stdout_path)?;
remove_rotated_log_files(&stderr_path)?;
remove_rotated_log_files(&combined_path)?;
Ok(())
}
pub fn clear_project_logs(&self, project: &str) -> Result<(), LogsManagerError> {
let root = runtime::log_dir();
let project_dir = contained_project_log_dir(&root, project)?;
let entries = match fs::read_dir(&project_dir) {
Ok(entries) => entries,
Err(err) if err.kind() == io::ErrorKind::NotFound => return Ok(()),
Err(err) => return Err(err.into()),
};
let mut files: Vec<PathBuf> = Vec::new();
for entry in entries {
let entry = entry?;
if entry.file_type()?.is_file() {
files.push(entry.path());
}
}
for path in files {
let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else {
continue;
};
if file_name.ends_with(".log") {
truncate_log_file(&path)?;
remove_rotated_log_files(&path)?;
} else if is_rotated_backup(file_name) && path.exists() {
fs::remove_file(&path)?;
}
}
Ok(())
}
pub fn clear_all_logs(&self) -> Result<(), LogsManagerError> {
let log_dir = runtime::log_dir();
runtime::create_private_dir(&log_dir)?;
for path in managed_log_files(&log_dir)? {
let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else {
continue;
};
if file_name.ends_with(".log") {
truncate_log_file(&path)?;
remove_rotated_log_files(&path)?;
} else if is_rotated_backup(file_name) {
fs::remove_file(&path)?;
}
}
Ok(())
}
#[cfg(target_os = "linux")]
#[allow(clippy::too_many_arguments)]
fn show_logs_platform_with_mode(
&self,
project: &str,
service_name: &str,
pid: Option<u32>,
lines: usize,
kind: Option<&str>,
mode: TailMode,
show_header: bool,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
if show_header {
println!(
"\n+{:-^33}+\n\
| {:^31} |\n\
+{:-^33}+\n",
"-",
Self::format_log_title(service_name, pid),
"-"
);
}
let (stdout_path, stderr_path, combined_path) =
resolve_tail_targets(project, service_name, pid)?;
debug!("Streaming logs via tail for '{}'", service_name);
if matches!(mode, TailMode::OneShot) {
return write_log_file_tail(
std::io::stdout().lock(),
&stdout_path,
&stderr_path,
&combined_path,
lines,
kind,
filter,
);
}
if combined_path.exists()
&& (filter.grep.is_some() || kind.and_then(LogStream::from_filter).is_some())
{
return follow_filtered_log_file(
std::io::stdout().lock(),
&combined_path,
lines,
kind.and_then(LogStream::from_filter),
filter,
);
}
let mut cmd = Command::new("tail");
#[cfg(target_os = "linux")]
{
if let Some(pid_value) = pid
&& !process_fds_present(pid_value)
{
debug!(
"Falling back to log files for '{}' because /proc/{pid_value} fds are unavailable",
service_name
);
}
}
mode.configure_command(
&mut cmd,
lines,
&stdout_path,
&stderr_path,
&combined_path,
kind,
);
cmd.stdout(std::process::Stdio::inherit())
.stderr(std::process::Stdio::inherit());
let status = cmd.status()?;
if !status.success() {
return Err(LogsManagerError::TailCommandFailed(status.code()));
}
Ok(())
}
#[cfg(target_os = "linux")]
#[allow(clippy::too_many_arguments)]
fn stream_logs_platform_with_mode(
&self,
project: &str,
service_name: &str,
pid: Option<u32>,
lines: usize,
kind: Option<&str>,
mode: TailMode,
show_header: bool,
filter: &LogFilter,
stream: &UnixStream,
) -> Result<(), LogsManagerError> {
if show_header {
write_log_header(stream.try_clone()?, service_name, pid)?;
}
let (stdout_path, stderr_path, combined_path) =
resolve_tail_targets(project, service_name, pid)?;
debug!("Streaming logs via supervisor tail for '{}'", service_name);
if matches!(mode, TailMode::OneShot) {
return write_log_file_tail(
stream.try_clone()?,
&stdout_path,
&stderr_path,
&combined_path,
lines,
kind,
filter,
);
}
if combined_path.exists()
&& (filter.grep.is_some() || kind.and_then(LogStream::from_filter).is_some())
{
return follow_filtered_log_file(
stream.try_clone()?,
&combined_path,
lines,
kind.and_then(LogStream::from_filter),
filter,
);
}
let mut cmd = Command::new("tail");
if let Some(pid_value) = pid
&& !process_fds_present(pid_value)
{
debug!(
"Falling back to log files for '{}' because /proc/{pid_value} fds are unavailable",
service_name
);
}
mode.configure_command(
&mut cmd,
lines,
&stdout_path,
&stderr_path,
&combined_path,
kind,
);
let stdout_stream = stream.try_clone()?;
let stderr_stream = stream.try_clone()?;
let stdout_fd = stdout_stream.into_raw_fd();
let stderr_fd = stderr_stream.into_raw_fd();
unsafe {
cmd.stdout(Stdio::from_raw_fd(stdout_fd));
cmd.stderr(Stdio::from_raw_fd(stderr_fd));
}
let status = cmd.status()?;
if !status.success() {
return Err(LogsManagerError::TailCommandFailed(status.code()));
}
Ok(())
}
#[cfg(target_os = "macos")]
#[allow(clippy::too_many_arguments)]
fn show_logs_platform_with_mode(
&self,
project: &str,
service_name: &str,
pid: Option<u32>,
lines: usize,
kind: Option<&str>,
mode: TailMode,
show_header: bool,
filter: &LogFilter,
) -> Result<(), LogsManagerError> {
if show_header {
println!(
"\n+{:-^33}+\n\
| {:^31} |\n\
+{:-^33}+\n",
"-",
Self::format_log_title(service_name, pid),
"-"
);
}
let (stdout_path, stderr_path, combined_path) =
resolve_tail_targets(project, service_name, pid)?;
debug!("Streaming logs via tail for '{}'", service_name);
if matches!(mode, TailMode::OneShot) {
return write_log_file_tail(
std::io::stdout().lock(),
&stdout_path,
&stderr_path,
&combined_path,
lines,
kind,
filter,
);
}
if combined_path.exists()
&& (filter.grep.is_some() || kind.and_then(LogStream::from_filter).is_some())
{
return follow_filtered_log_file(
std::io::stdout().lock(),
&combined_path,
lines,
kind.and_then(LogStream::from_filter),
filter,
);
}
let mut cmd = Command::new("tail");
mode.configure_command(
&mut cmd,
lines,
&stdout_path,
&stderr_path,
&combined_path,
kind,
);
cmd.stdout(std::process::Stdio::inherit())
.stderr(std::process::Stdio::inherit());
let status = cmd.status()?;
if !status.success() {
return Err(LogsManagerError::TailCommandFailed(status.code()));
}
Ok(())
}
#[cfg(target_os = "macos")]
#[allow(clippy::too_many_arguments)]
fn stream_logs_platform_with_mode(
&self,
project: &str,
service_name: &str,
pid: Option<u32>,
lines: usize,
kind: Option<&str>,
mode: TailMode,
show_header: bool,
filter: &LogFilter,
stream: &UnixStream,
) -> Result<(), LogsManagerError> {
if show_header {
write_log_header(stream.try_clone()?, service_name, pid)?;
}
let (stdout_path, stderr_path, combined_path) =
resolve_tail_targets(project, service_name, pid)?;
debug!("Streaming logs via supervisor tail for '{}'", service_name);
if matches!(mode, TailMode::OneShot) {
return write_log_file_tail(
stream.try_clone()?,
&stdout_path,
&stderr_path,
&combined_path,
lines,
kind,
filter,
);
}
if combined_path.exists()
&& (filter.grep.is_some() || kind.and_then(LogStream::from_filter).is_some())
{
return follow_filtered_log_file(
stream.try_clone()?,
&combined_path,
lines,
kind.and_then(LogStream::from_filter),
filter,
);
}
let mut cmd = Command::new("tail");
mode.configure_command(
&mut cmd,
lines,
&stdout_path,
&stderr_path,
&combined_path,
kind,
);
let stdout_stream = stream.try_clone()?;
let stderr_stream = stream.try_clone()?;
let stdout_fd = stdout_stream.into_raw_fd();
let stderr_fd = stderr_stream.into_raw_fd();
unsafe {
cmd.stdout(Stdio::from_raw_fd(stdout_fd));
cmd.stderr(Stdio::from_raw_fd(stderr_fd));
}
let status = cmd.status()?;
if !status.success() {
return Err(LogsManagerError::TailCommandFailed(status.code()));
}
Ok(())
}
fn format_log_title(service_name: &str, pid: Option<u32>) -> String {
match pid {
Some(pid) => format!("{service_name} [pid {pid}]"),
None => format!("{service_name} [offline]"),
}
}
pub fn show_supervisor_log(&self, lines: usize) -> Result<(), LogsManagerError> {
let supervisor_log = runtime::log_dir().join("supervisor.log");
if !supervisor_log.exists() {
return Ok(());
}
println!(
"\n+{:-^33}+\n\
| {:^31} |\n\
+{:-^33}+\n",
"-", "Supervisor", "-"
);
let tail = tail_log_file(&supervisor_log, lines)?;
let mut stdout = std::io::stdout().lock();
stdout.write_all(&tail)?;
stdout.flush()?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use std::{
fs::{self, File},
io::Cursor,
path::Path,
thread,
time::Duration,
};
use tempfile::tempdir_in;
use super::*;
#[test]
fn contained_project_log_dir_accepts_only_a_direct_child() {
let root = Path::new("/logs");
assert_eq!(
contained_project_log_dir(root, "alpha").unwrap(),
root.join("alpha")
);
for bad in ["", ".", "..", "a/b", "/etc", "a/../../etc"] {
assert!(
contained_project_log_dir(root, bad).is_err(),
"expected '{bad}' to be refused"
);
}
}
#[test]
fn contained_project_log_dir_refuses_a_symlink_out_of_the_root() {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("logs");
let outside = temp.path().join("outside");
fs::create_dir_all(&root).unwrap();
fs::create_dir_all(&outside).unwrap();
std::os::unix::fs::symlink(&outside, root.join("evil")).unwrap();
assert!(contained_project_log_dir(&root, "evil").is_err());
fs::create_dir_all(root.join("real")).unwrap();
assert!(contained_project_log_dir(&root, "real").is_ok());
}
#[test]
fn validate_service_name_accepts_plain_names() {
for name in ["api", "web-1", "worker_2", "svc.v1", "A.B_c-3"] {
assert!(validate_service_name(name).is_ok(), "rejected {name}");
}
}
#[test]
fn validate_service_name_rejects_traversal() {
for name in [
"",
".",
"..",
"../etc/passwd",
"../../../../etc/cron.d/x",
"a/b",
"a\\b",
"nul\0byte",
] {
assert!(
validate_service_name(name).is_err(),
"accepted traversal name {name:?}"
);
}
}
#[test]
fn tail_log_bytes_returns_last_lines_with_trailing_newline() {
let bytes = b"line 1\nline 2\nline 3\nline 4\n";
assert_eq!(tail_log_bytes(bytes, 2), b"line 3\nline 4\n");
}
#[test]
fn tail_log_bytes_returns_last_lines_without_trailing_newline() {
let bytes = b"line 1\nline 2\nline 3\nline 4";
assert_eq!(tail_log_bytes(bytes, 2), b"line 3\nline 4");
}
#[test]
fn tail_log_bytes_returns_all_bytes_when_line_count_fits() {
let bytes = b"line 1\nline 2\n";
assert_eq!(tail_log_bytes(bytes, 2), bytes);
assert_eq!(tail_log_bytes(bytes, 5), bytes);
}
#[test]
fn tail_log_bytes_returns_empty_when_zero_lines_requested() {
assert_eq!(tail_log_bytes(b"line 1\nline 2\n", 0), b"");
}
#[test]
fn diagnostic_log_lines_exclude_prior_generations() {
let cutoff = "2026-07-20T18:33:14.000000Z".parse().unwrap();
let text = "2026-07-20T18:00:00.000000Z stderr Address already in use\n\
2026-07-20T18:33:14.200000Z stderr candidate failed\n";
assert_eq!(
diagnostic_log_lines(text, 8, Some(cutoff)),
vec!["candidate failed"]
);
}
#[test]
fn resolve_log_path_matches_hyphenated_files() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
let log_dir = canonical_log_path("__loose__", "dummy", "stdout")
.parent()
.map(Path::to_path_buf)
.unwrap();
fs::create_dir_all(&log_dir).unwrap();
let target = log_dir.join("arb-rs_stdout.log");
File::create(&target).unwrap();
let resolved = resolve_log_path("__loose__", "arb_rs", "stdout");
assert_eq!(resolved, target);
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn spawn_dynamic_child_log_writer_persists_output() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
let reader = Cursor::new(b"hello\nworld\n".to_vec());
super::spawn_dynamic_child_log_writer(
Some("alpha"),
"beta",
123,
reader,
"stdout",
false,
)
.expect("spawn child log writer");
thread::sleep(Duration::from_millis(100));
let log_path = crate::runtime::log_dir()
.join("spawn")
.join("alpha_beta_123_stdout.log");
let contents =
fs::read_to_string(&log_path).expect("spawn log should be written");
assert!(contents.contains("hello"));
assert!(contents.contains("world"));
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn spawn_log_writer_persists_unterminated_output() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
super::spawn_log_writer(
"__loose__",
"svc",
Cursor::new(b"partial line".to_vec()),
"stdout",
)
.expect("spawn service log writer");
thread::sleep(Duration::from_millis(100));
let log_path = get_service_log_path("__loose__", "svc");
let contents =
fs::read_to_string(&log_path).expect("service log should be written");
assert!(contents.contains(" stdout partial line\n"));
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn spawn_log_writer_persists_non_utf8_output() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
super::spawn_log_writer(
"__loose__",
"svc",
Cursor::new(vec![0xff, b'a', b'\n']),
"stderr",
)
.expect("spawn service log writer");
thread::sleep(Duration::from_millis(100));
let log_path = get_service_log_path("__loose__", "svc");
let contents =
fs::read_to_string(&log_path).expect("service log should be written");
assert!(contents.contains(" stderr "));
assert!(contents.contains("a\n"));
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn spawn_log_writer_rotates_when_active_file_exceeds_limit() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
let settings = EffectiveLogsConfig {
sink: crate::config::LogSink::File,
max_bytes: 6,
max_files: 1,
};
let log_path = get_service_log_path("__loose__", "svc");
fs::create_dir_all(log_path.parent().expect("log parent")).unwrap();
fs::write(&log_path, "first\n").unwrap();
super::spawn_log_writer_with_config(
"__loose__",
"svc",
Cursor::new(b"second\n".to_vec()),
"stdout",
settings,
)
.expect("spawn rotating service log writer");
thread::sleep(Duration::from_millis(100));
let active = fs::read_to_string(&log_path).expect("active log exists");
let rotated = fs::read_to_string(rotated_log_path(&log_path, 1))
.expect("rotated log exists");
assert_eq!(rotated, "first\n");
assert!(active.contains(" stdout second\n"));
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn truncate_log_payload_leaves_small_lines_untouched() {
let line = b"short line";
assert_eq!(truncate_log_payload(line), line);
}
#[test]
fn truncate_log_payload_caps_oversized_lines() {
let line = vec![b'a'; MAX_LOG_LINE_BYTES + 4096];
let truncated = truncate_log_payload(&line);
assert!(truncated.len() < line.len());
let text = String::from_utf8(truncated).expect("valid utf8");
assert!(text.contains("[truncated 4096 bytes]"));
}
#[test]
fn parse_byte_size_handles_units() {
assert_eq!(parse_byte_size("1024").unwrap(), 1024);
assert_eq!(parse_byte_size("1kb").unwrap(), 1024);
assert_eq!(parse_byte_size("2MB").unwrap(), 2 * 1024 * 1024);
assert_eq!(parse_byte_size("1g").unwrap(), 1024 * 1024 * 1024);
assert!(parse_byte_size("nonsense").is_err());
}
#[test]
fn parse_age_seconds_handles_units() {
assert_eq!(parse_age_seconds("30").unwrap(), 30);
assert_eq!(parse_age_seconds("5m").unwrap(), 300);
assert_eq!(parse_age_seconds("2h").unwrap(), 7200);
assert_eq!(parse_age_seconds("7d").unwrap(), 7 * 24 * 60 * 60);
assert!(parse_age_seconds("12x").is_err());
}
#[test]
fn is_rotated_backup_matches_numbered_files_only() {
assert!(is_rotated_backup("supervisor.log.1"));
assert!(is_rotated_backup("api.log.12"));
assert!(!is_rotated_backup("api.log"));
assert!(!is_rotated_backup("api.log.bak"));
}
#[test]
fn prune_logs_trims_backups_by_size() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
let log_dir = crate::runtime::log_dir();
let project_log_dir = log_dir.join("demo");
let spawn_log_dir = log_dir.join(SPAWN_LOG_DIR);
fs::create_dir_all(&project_log_dir).unwrap();
fs::create_dir_all(&spawn_log_dir).unwrap();
fs::write(project_log_dir.join("svc.log"), vec![b'a'; 100]).unwrap();
fs::write(project_log_dir.join("svc.log.1"), vec![b'a'; 100]).unwrap();
fs::write(project_log_dir.join("svc.log.2"), vec![b'a'; 100]).unwrap();
fs::write(spawn_log_dir.join("child.log.1"), vec![b'a'; 100]).unwrap();
let summary = super::prune_logs(Some("150"), None).unwrap();
assert!(summary.removed_files >= 1);
assert!(project_log_dir.join("svc.log").exists());
assert!(spawn_log_dir.join("child.log.1").exists());
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn rotating_log_writer_rotates_supervisor_output() {
let _guard = crate::test_utils::env_lock();
let base = std::env::current_dir()
.expect("current_dir")
.join("target/tmp-home");
fs::create_dir_all(&base).unwrap();
let temp = tempdir_in(&base).unwrap();
let home = temp.path();
let original_home = std::env::var("HOME").ok();
unsafe {
std::env::set_var("HOME", home);
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
let path = supervisor_log_path();
fs::create_dir_all(path.parent().expect("log parent")).unwrap();
let settings = EffectiveLogsConfig {
sink: crate::config::LogSink::File,
max_bytes: 8,
max_files: 1,
};
let mut writer = RotatingLogWriter::open(path.clone(), settings).unwrap();
writer.write_all(b"first\n").unwrap();
writer.write_all(b"second\n").unwrap();
writer.flush().unwrap();
let active = fs::read_to_string(&path).expect("active log exists");
let rotated =
fs::read_to_string(rotated_log_path(&path, 1)).expect("rotated log exists");
assert_eq!(rotated, "first\n");
assert_eq!(active, "second\n");
unsafe {
if let Some(home) = original_home {
std::env::set_var("HOME", home);
} else {
std::env::remove_var("HOME");
}
}
crate::runtime::init(crate::runtime::RuntimeMode::User);
crate::runtime::set_drop_privileges(false);
}
#[test]
fn forwarded_console_line_preserves_ansi_bytes() {
let mut output = Vec::new();
let line = "\u{1b}[34mDEBUG\u{1b}[0m child log";
write_forwarded_console_line(&mut output, "[svc] ", line)
.expect("console line should write");
assert_eq!(
String::from_utf8(output).expect("valid utf8"),
format!("[svc] {line}\n")
);
}
fn utc(text: &str) -> chrono::DateTime<chrono::Utc> {
chrono::DateTime::parse_from_rfc3339(text)
.expect("valid rfc3339")
.with_timezone(&chrono::Utc)
}
#[test]
fn parse_time_bound_accepts_rfc3339_and_date() {
let now = utc("2026-07-07T12:00:00Z");
assert_eq!(
parse_time_bound("2026-07-07T09:30:00Z", now).unwrap(),
utc("2026-07-07T09:30:00Z")
);
assert_eq!(
parse_time_bound("2026-07-07", now).unwrap(),
utc("2026-07-07T00:00:00Z")
);
}
#[test]
fn parse_time_bound_accepts_relative_age() {
let now = utc("2026-07-07T12:00:00Z");
assert_eq!(
parse_time_bound("2h", now).unwrap(),
utc("2026-07-07T10:00:00Z")
);
}
#[test]
fn parse_time_bound_rejects_garbage() {
let now = utc("2026-07-07T12:00:00Z");
assert!(parse_time_bound("not-a-time", now).is_err());
}
#[test]
fn log_filter_applies_time_window() {
let bytes = b"2026-07-07T09:00:00Z stdout early\n\
2026-07-07T10:30:00Z stdout middle\n\
2026-07-07T12:00:00Z stdout late\n";
let filter = LogFilter {
since: Some(utc("2026-07-07T10:00:00Z")),
until: Some(utc("2026-07-07T11:00:00Z")),
..LogFilter::default()
};
let out = String::from_utf8(filter.apply(bytes)).unwrap();
assert_eq!(out, "2026-07-07T10:30:00Z stdout middle\n");
}
#[test]
fn log_filter_applies_grep() {
let bytes = b"2026-07-07T09:00:00Z stdout hello world\n\
2026-07-07T09:00:01Z stderr ERROR boom\n\
2026-07-07T09:00:02Z stdout all good\n";
let filter = LogFilter {
grep: Some(regex::Regex::new("ERROR|good").unwrap()),
..LogFilter::default()
};
let out = String::from_utf8(filter.apply(bytes)).unwrap();
assert!(out.contains("ERROR boom"));
assert!(out.contains("all good"));
assert!(!out.contains("hello world"));
}
#[test]
fn collect_all_ignores_default_lines_cap() {
let dir = std::env::temp_dir().join(format!(
"sysg_all_{}_{}",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
fs::create_dir_all(&dir).unwrap();
let combined = dir.join("svc.log");
let mut body = String::new();
for index in 0..200 {
body.push_str(&format!(
"2026-07-08T09:00:{index:02}Z stdout line {index}\n"
));
}
fs::write(&combined, body).unwrap();
let missing = dir.join("svc_stdout.log");
let filter = LogFilter {
all: true,
..LogFilter::default()
};
let chunks =
collect_log_tail(&missing, &missing, &combined, 50, None, &filter).unwrap();
let text = String::from_utf8(chunks.concat()).unwrap();
assert_eq!(text.lines().count(), 200);
fs::remove_dir_all(&dir).ok();
}
#[test]
fn collect_all_time_window_spans_rotated_history() {
let dir = std::env::temp_dir().join(format!(
"sysg_rot_{}_{}",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
fs::create_dir_all(&dir).unwrap();
let combined = dir.join("svc.log");
fs::write(
dir.join("svc.log.1"),
"2026-07-04T09:00:00Z stdout old rotated\n\
2026-07-08T09:00:00Z stdout kept rotated\n",
)
.unwrap();
fs::write(
&combined,
"2026-07-08T10:00:00Z stdout kept active\n\
2026-07-09T09:00:00Z stdout too new\n",
)
.unwrap();
let missing = dir.join("svc_stdout.log");
let filter = LogFilter {
since: Some(utc("2026-07-08T00:00:00Z")),
until: Some(utc("2026-07-09T00:00:00Z")),
all: true,
..LogFilter::default()
};
let chunks =
collect_log_tail(&missing, &missing, &combined, 50, None, &filter).unwrap();
let text = String::from_utf8(chunks.concat()).unwrap();
assert!(text.contains("kept rotated"), "{text}");
assert!(text.contains("kept active"), "{text}");
assert!(!text.contains("old rotated"), "{text}");
assert!(!text.contains("too new"), "{text}");
fs::remove_dir_all(&dir).ok();
}
#[test]
fn log_filter_noop_returns_input() {
let bytes = b"line without leading timestamp\n";
let filter = LogFilter::default();
assert_eq!(filter.apply(bytes), bytes);
assert!(filter.is_noop());
}
#[test]
fn strip_ansi_removes_color_codes() {
let input = b"\x1b[1;31mERROR\x1b[0m boom \x1b[34mblue\x1b[0m";
assert_eq!(strip_ansi(input), b"ERROR boom blue");
}
#[test]
fn strip_ansi_leaves_plain_text() {
let input = b"plain line 42";
assert_eq!(strip_ansi(input), input);
}
#[test]
fn parse_captured_line_extracts_fields() {
let parsed =
parse_captured_line("2026-07-07T09:00:00Z stdout hello world").unwrap();
assert_eq!(parsed.timestamp, "2026-07-07T09:00:00Z");
assert_eq!(parsed.stream, "stdout");
assert_eq!(parsed.message, "hello world");
}
#[test]
fn parse_captured_line_rejects_banner() {
assert!(parse_captured_line("┌─────────┐").is_none());
assert!(parse_captured_line("Project: arbitration").is_none());
}
#[test]
fn log_writer_json_emits_one_object_per_line() {
let mut out = Vec::new();
{
let mut writer =
LogWriter::new(&mut out, LogFormat::Json, true, Some("api".into()));
writer
.write_all(b"2026-07-07T09:00:00Z stdout \x1b[31mhello\x1b[0m\n")
.unwrap();
writer.write_all(b"\xe2\x94\x8c banner line\n").unwrap();
writer.flush().unwrap();
}
let text = String::from_utf8(out).unwrap();
assert_eq!(
text,
"{\"ts\":\"2026-07-07T09:00:00Z\",\"stream\":\"stdout\",\"service\":\"api\",\"line\":\"hello\"}\n"
);
}
#[test]
fn log_writer_json_service_follows_marker_lines() {
let mut out = Vec::new();
{
let mut writer = LogWriter::new(&mut out, LogFormat::Json, true, None);
writer
.write_all(&service_marker_line("arb_rs__server"))
.unwrap();
writer
.write_all(b"2026-07-08T09:00:00Z stdout openai_call\n")
.unwrap();
writer
.write_all(&service_marker_line("arb_py__curator"))
.unwrap();
writer
.write_all(b"2026-07-08T09:00:01Z stderr ingest done\n")
.unwrap();
writer.flush().unwrap();
}
let text = String::from_utf8(out).unwrap();
assert_eq!(
text,
"{\"ts\":\"2026-07-08T09:00:00Z\",\"stream\":\"stdout\",\"service\":\"arb_rs__server\",\"line\":\"openai_call\"}\n\
{\"ts\":\"2026-07-08T09:00:01Z\",\"stream\":\"stderr\",\"service\":\"arb_py__curator\",\"line\":\"ingest done\"}\n"
);
}
#[test]
fn log_writer_drops_marker_lines_in_text_mode() {
let mut out = Vec::new();
{
let mut writer = LogWriter::new(&mut out, LogFormat::Text, true, None);
writer.write_all(&service_marker_line("svc")).unwrap();
writer
.write_all(b"2026-07-08T09:00:00Z stdout hello\n")
.unwrap();
writer.flush().unwrap();
}
assert_eq!(
String::from_utf8(out).unwrap(),
"2026-07-08T09:00:00Z stdout hello\n"
);
}
#[test]
fn log_writer_raw_strips_prefix_and_banner() {
let mut out = Vec::new();
{
let mut writer = LogWriter::new(&mut out, LogFormat::Raw, true, None);
writer
.write_all(b"2026-07-07T09:00:00Z stderr actual message\n")
.unwrap();
writer.write_all(b"Running Services\n").unwrap();
writer.flush().unwrap();
}
assert_eq!(String::from_utf8(out).unwrap(), "actual message\n");
}
#[test]
fn log_writer_text_strip_ansi_keeps_prefix() {
let mut out = Vec::new();
{
let mut writer = LogWriter::new(&mut out, LogFormat::Text, true, None);
writer
.write_all(b"2026-07-07T09:00:00Z stdout \x1b[32mok\x1b[0m\n")
.unwrap();
writer.flush().unwrap();
}
assert_eq!(
String::from_utf8(out).unwrap(),
"2026-07-07T09:00:00Z stdout ok\n"
);
}
#[test]
fn rotated_history_paths_orders_oldest_to_newest() {
let dir = std::env::temp_dir().join(format!(
"sysg_hist_{}_{}",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
));
fs::create_dir_all(&dir).unwrap();
let active = dir.join("svc.log");
for name in ["svc.log", "svc.log.1", "svc.log.2", "svc.log.10"] {
fs::write(dir.join(name), b"x").unwrap();
}
let paths = rotated_history_paths(&active);
let names: Vec<_> = paths
.iter()
.map(|path| path.file_name().unwrap().to_string_lossy().into_owned())
.collect();
assert_eq!(names, ["svc.log.10", "svc.log.2", "svc.log.1", "svc.log"]);
fs::remove_dir_all(&dir).ok();
}
}