use std::io::Write;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, OnceLock, PoisonError};
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use tracing_subscriber::{EnvFilter, Layer, Registry, reload};
type OtelSlot = Option<leviath_telemetry::LogLayer>;
static OTEL_HANDLE: OnceLock<reload::Handle<OtelSlot, Registry>> = OnceLock::new();
static TUI_HOLDS_TERMINAL: AtomicBool = AtomicBool::new(false);
static PARKED: Mutex<Vec<u8>> = Mutex::new(Vec::new());
struct TerminalAwareWriter;
impl Write for TerminalAwareWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
if TUI_HOLDS_TERMINAL.load(Ordering::Relaxed) {
PARKED
.lock()
.unwrap_or_else(PoisonError::into_inner)
.extend_from_slice(buf);
return Ok(buf.len());
}
std::io::stderr().write(buf)
}
fn flush(&mut self) -> std::io::Result<()> {
if TUI_HOLDS_TERMINAL.load(Ordering::Relaxed) {
return Ok(());
}
std::io::stderr().flush()
}
}
fn writer() -> TerminalAwareWriter {
TerminalAwareWriter
}
pub fn hold_for_tui() {
TUI_HOLDS_TERMINAL.store(true, Ordering::Relaxed);
}
pub fn release_from_tui() {
TUI_HOLDS_TERMINAL.store(false, Ordering::Relaxed);
let parked = std::mem::take(&mut *PARKED.lock().unwrap_or_else(PoisonError::into_inner));
if parked.is_empty() {
return;
}
let _ = std::io::stderr().write_all(&parked);
let _ = std::io::stderr().flush();
}
pub fn init(verbose: bool) {
let level = if verbose { "debug" } else { "info" };
let (otel_layer, handle) = reload::Layer::new(None as OtelSlot);
let subscriber = tracing_subscriber::registry().with(otel_layer).with(
tracing_subscriber::fmt::layer()
.with_writer(writer)
.with_filter(EnvFilter::new(level)),
);
let _ = subscriber.try_init();
let _ = OTEL_HANDLE.set(handle);
}
pub fn install_otel_layer(layer: leviath_telemetry::LogLayer) -> bool {
match OTEL_HANDLE.get() {
Some(handle) => handle.reload(Some(layer)).is_ok(),
None => false,
}
}
#[cfg(test)]
mod tests {
use super::*;
use opentelemetry_sdk::logs::{InMemoryLogExporter, SdkLoggerProvider};
fn bridge_with_exporter() -> (leviath_telemetry::LogLayer, InMemoryLogExporter) {
let exporter = InMemoryLogExporter::default();
let provider = SdkLoggerProvider::builder()
.with_simple_exporter(exporter.clone())
.build();
let sink = leviath_telemetry::OtelSink::new(
opentelemetry_sdk::trace::SdkTracerProvider::builder().build(),
opentelemetry_sdk::metrics::SdkMeterProvider::builder().build(),
provider,
);
(sink.tracing_log_layer(), exporter)
}
#[test]
fn init_parks_the_handle_and_install_forwards_events() {
let (layer, _exporter) = bridge_with_exporter();
assert!(!install_otel_layer(layer));
let (otel_layer, handle) = reload::Layer::new(None as OtelSlot);
assert!(
OTEL_HANDLE.set(handle).is_ok(),
"this test parks the handle first"
);
let subscriber = tracing_subscriber::registry().with(otel_layer);
let _guard = tracing::subscriber::set_default(subscriber);
let (layer, exporter) = bridge_with_exporter();
assert!(install_otel_layer(layer));
tracing::info!(target: "leviath::logging::test", "forwarded line");
let emitted = exporter.get_emitted_logs().unwrap();
assert!(
emitted
.iter()
.any(|log| format!("{:?}", log.record.body()).contains("forwarded line")),
"{emitted:?}"
);
tracing::info!(target: "opentelemetry_sdk", "feedback line");
let emitted = exporter.get_emitted_logs().unwrap();
assert!(
!emitted
.iter()
.any(|log| format!("{:?}", log.record.body()).contains("feedback line"))
);
init(false);
init(true);
let (layer, _exporter) = bridge_with_exporter();
assert!(install_otel_layer(layer));
}
#[test]
fn holding_the_terminal_parks_output_until_it_is_released() {
release_from_tui();
assert!(!TUI_HOLDS_TERMINAL.load(Ordering::Relaxed));
writer().write_all(b"").expect("stderr accepts a write");
writer().flush().expect("stderr accepts a flush");
hold_for_tui();
writer()
.write_all(b"parked line\n")
.expect("a held write is buffered, never refused");
writer().flush().expect("a held flush is a no-op");
assert_eq!(
PARKED.lock().expect("uncontended").as_slice(),
b"parked line\n"
);
release_from_tui();
assert!(
PARKED.lock().expect("uncontended").is_empty(),
"release hands the buffer to stderr and empties it"
);
release_from_tui();
}
}