use crate::adapters::{Adapter, AdapterError};
use crate::beholders::{BeholderCtx, LogLine, ServiceBeholder};
use crate::service::Scryer;
use async_trait::async_trait;
use observation::EventScope;
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::mpsc;
use workload_spec::MeshIdent;
#[async_trait]
pub trait ContainerLogSource: Send + Sync {
async fn connect(&self, ident: &MeshIdent) -> Result<mpsc::Receiver<String>, AdapterError>;
}
pub struct ContainerdLogsAdapter {
name: String,
ident: MeshIdent,
scryer: Arc<Scryer>,
source: Arc<dyn ContainerLogSource>,
beholder: Box<dyn ServiceBeholder>,
ctx: BeholderCtx,
started_at: Instant,
}
impl ContainerdLogsAdapter {
pub fn new(
scryer: Arc<Scryer>,
ident: MeshIdent,
source: Arc<dyn ContainerLogSource>,
beholder: Box<dyn ServiceBeholder>,
) -> Self {
Self {
name: format!("containerd_logs::{}", ident.0),
ident,
scryer,
source,
beholder,
ctx: BeholderCtx::new(),
started_at: Instant::now(),
}
}
fn offset_ms(&self) -> u32 {
let elapsed = self.started_at.elapsed().as_millis();
elapsed.min(u32::MAX as u128) as u32
}
}
#[async_trait]
impl Adapter for ContainerdLogsAdapter {
fn name(&self) -> &str {
&self.name
}
fn scope(&self) -> EventScope {
EventScope::Service(self.ident.clone())
}
async fn run(&mut self) -> Result<(), AdapterError> {
let mut rx = self.source.connect(&self.ident).await?;
let scope = EventScope::Service(self.ident.clone());
while let Some(line) = rx.recv().await {
let log = LogLine { line, offset_ms: self.offset_ms() };
let events = self.beholder.parse_line(&log, &mut self.ctx);
for ev in events {
self.scryer.push(scope.clone(), ev)?;
}
}
Err(AdapterError::StreamBroken(format!(
"containerd log stream for {} closed",
self.ident.0
)))
}
}
pub struct DockerLogSource {
tail: String,
}
impl DockerLogSource {
pub fn new() -> std::sync::Arc<Self> {
std::sync::Arc::new(Self { tail: "50".to_string() })
}
pub fn with_tail(n: u32) -> std::sync::Arc<Self> {
std::sync::Arc::new(Self { tail: n.to_string() })
}
pub fn follow_only() -> std::sync::Arc<Self> {
std::sync::Arc::new(Self { tail: "0".to_string() })
}
}
#[async_trait]
impl ContainerLogSource for DockerLogSource {
async fn connect(&self, ident: &MeshIdent) -> Result<mpsc::Receiver<String>, AdapterError> {
use std::process::Stdio;
use tokio::io::{AsyncBufReadExt, BufReader};
let container = ident.0.clone();
let mut child = tokio::process::Command::new("docker")
.args(["logs", "--follow", "--tail", &self.tail, &container])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true)
.spawn()
.map_err(|e| AdapterError::Permanent(format!("docker logs spawn failed: {e}")))?;
let stdout = child.stdout.take().ok_or_else(|| {
AdapterError::Permanent("docker logs: could not capture stdout".into())
})?;
let stderr = child.stderr.take().ok_or_else(|| {
AdapterError::Permanent("docker logs: could not capture stderr".into())
})?;
let (tx, rx) = mpsc::channel::<String>(256);
let tx1 = tx.clone();
tokio::spawn(async move {
let mut lines = BufReader::new(stdout).lines();
while let Ok(Some(line)) = lines.next_line().await {
if tx1.send(line).await.is_err() {
break;
}
}
});
let tx2 = tx.clone();
tokio::spawn(async move {
let mut lines = BufReader::new(stderr).lines();
while let Ok(Some(line)) = lines.next_line().await {
if tx2.send(line).await.is_err() {
break;
}
}
});
tokio::spawn(async move {
let _ = child.wait().await;
drop(tx);
});
Ok(rx)
}
}
#[cfg(test)]
pub mod test_source {
use super::*;
use std::sync::Mutex;
pub struct ScriptedSource {
streams: Mutex<Vec<Vec<String>>>,
}
impl ScriptedSource {
pub fn new(streams: Vec<Vec<String>>) -> Arc<Self> {
Arc::new(Self { streams: Mutex::new(streams) })
}
}
#[async_trait]
impl ContainerLogSource for ScriptedSource {
async fn connect(&self, _ident: &MeshIdent) -> Result<mpsc::Receiver<String>, AdapterError> {
let mut g = self.streams.lock().unwrap();
if g.is_empty() {
return Err(AdapterError::StreamBroken("scripted source exhausted".into()));
}
let lines = g.remove(0);
drop(g);
let (tx, rx) = mpsc::channel::<String>(64);
tokio::spawn(async move {
for line in lines {
if tx.send(line).await.is_err() {
break;
}
}
});
Ok(rx)
}
}
}
#[cfg(test)]
mod tests {
use super::test_source::ScriptedSource;
use super::*;
fn docker_available() -> bool {
std::process::Command::new("docker")
.args(["info", "--format", "{{.ServerVersion}}"])
.output()
.map(|o| o.status.success())
.unwrap_or(false)
}
#[tokio::test]
async fn docker_log_source_streams_alpine_echo() {
if !docker_available() {
eprintln!("[skip] docker not reachable");
return;
}
let run = std::process::Command::new("docker")
.args(["run", "--rm", "--name", "scryer-test-echo", "-d",
"alpine:latest", "sh", "-c",
"echo 'hello scryer'; sleep 0.1; echo 'goodbye scryer'"])
.output();
let Ok(run) = run else { return; };
if !run.status.success() { return; }
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
let source = DockerLogSource::with_tail(100);
let ident = MeshIdent("scryer-test-echo".to_string());
let rx = source.connect(&ident).await;
match rx {
Ok(mut rx) => {
let mut lines = vec![];
while let Some(line) = rx.recv().await {
lines.push(line);
}
let joined = lines.join("\n");
eprintln!("[docker-source] got {} lines: {:?}", lines.len(), joined);
}
Err(AdapterError::Permanent(_)) => {
}
Err(e) => panic!("unexpected error: {e}"),
}
}
use crate::adapters::{BackoffConfig, Supervisor};
use crate::beholders::UnstructuredBeholder;
use crate::service::{EventFilter, Scryer, ScryerConfig};
use observation::Level;
use tempfile::TempDir;
fn make_scryer(dir: &TempDir) -> Arc<Scryer> {
let cfg = ScryerConfig::new(dir.path().join("events.db"));
Arc::new(Scryer::new(cfg, None).unwrap())
}
#[tokio::test]
async fn happy_path_emits_parsed_lines() {
let dir = TempDir::new().unwrap();
let scryer = make_scryer(&dir);
let ident = MeshIdent("api.pdx".to_string());
let source = ScriptedSource::new(vec![vec![
"first line".to_string(),
"second line".to_string(),
]]);
let beholder: Box<dyn ServiceBeholder> = Box::new(UnstructuredBeholder::for_ident(&ident));
let mut adapter = ContainerdLogsAdapter::new(
scryer.clone(),
ident.clone(),
source,
beholder,
);
let result = adapter.run().await;
assert!(matches!(result, Err(AdapterError::StreamBroken(_))));
scryer.flush_ring().unwrap();
let scope = EventScope::Service(ident);
let events = scryer.events(&scope, &EventFilter::default()).await.unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].msg, "first line");
assert_eq!(events[1].msg, "second line");
assert_eq!(events[0].level, Level::Info);
}
#[tokio::test]
async fn restart_emits_service_restart_event() {
let dir = TempDir::new().unwrap();
let scryer = make_scryer(&dir);
let ident = MeshIdent("api.pdx".to_string());
let source = ScriptedSource::new(vec![
vec!["before restart".to_string()],
vec!["after restart".to_string()],
]);
let beholder: Box<dyn ServiceBeholder> = Box::new(UnstructuredBeholder::for_ident(&ident));
let mut adapter = ContainerdLogsAdapter::new(
scryer.clone(),
ident.clone(),
source,
beholder,
);
let scope = EventScope::Service(ident.clone());
let mut supervisor = Supervisor::new(scryer.clone(), scope.clone(), "containerd_logs")
.with_backoff(BackoffConfig::test());
let _ = supervisor.run(&mut adapter).await;
scryer.flush_ring().unwrap();
let events = scryer.events(&scope, &EventFilter::default()).await.unwrap();
let restart_events: Vec<_> = events
.iter()
.filter(|e| e.msg == "service.restart")
.collect();
assert!(
!restart_events.is_empty(),
"expected at least one service.restart synth event"
);
assert!(
events.iter().any(|e| e.msg == "before restart"),
"expected the pre-restart line"
);
assert!(
events.iter().any(|e| e.msg == "after restart"),
"expected the post-restart line"
);
}
}