yah-scryer 0.8.42

Per-machine event store: ring buffer + short-disk SQLite + cross-service query surface
Documentation
//! `adapter::containerd_logs` — subscribes to a container's stdout/stderr.
//!
//! Production deployment: yubaba owns the containerd gRPC client; it provides a
//! [`ContainerLogSource`] impl that yields lines per task. Scryer's adapter is
//! transport-agnostic — it consumes a `ContainerLogSource` impl and forwards
//! lines through a [`ServiceBeholder`] into the [`Scryer`] store.
//!
//! Restart semantics: when the source's stream closes, the supervisor backs
//! off and reconnects. The supervisor (not the adapter) emits the synthetic
//! `service.restart` event so the discontinuity shows up in consumers.
//!

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;

// ─── ContainerLogSource ───────────────────────────────────────────────────────

/// Trait yubaba (or test code) implements to provide a line-oriented log
/// stream for a given mesh identity.
///
/// `connect` opens a fresh stream — the adapter calls it once per `run()`
/// invocation. The returned `mpsc::Receiver<String>` yields one log line at a
/// time. When the receiver closes (sender dropped) the adapter returns
/// `Err(StreamBroken)` so the supervisor can decide whether to retry.
#[async_trait]
pub trait ContainerLogSource: Send + Sync {
    async fn connect(&self, ident: &MeshIdent) -> Result<mpsc::Receiver<String>, AdapterError>;
}

// ─── ContainerdLogsAdapter ────────────────────────────────────────────────────

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)?;
            }
        }

        // Receiver closed without a clean shutdown signal; treat as stream
        // break so the supervisor restarts.
        Err(AdapterError::StreamBroken(format!(
            "containerd log stream for {} closed",
            self.ident.0
        )))
    }
}

// ─── DockerLogSource ──────────────────────────────────────────────────────────

/// `ContainerLogSource` impl that streams stdout + stderr from a running or
/// recently-exited Docker container via `docker logs --follow`.
///
/// The `MeshIdent` passed to [`ContainerLogSource::connect`] is used as the
/// container name directly — callers use the canonical pond container name
/// (`yah-pond-<svc>-<env>-<slot>`) as the ident.
///
/// Restart semantics: when `docker logs --follow` exits (because the container
/// stopped), the channel closes. The [`crate::adapters::Supervisor`] wrapping
/// [`ContainerdLogsAdapter`] interprets the resulting `StreamBroken` as a
/// restart signal, emits a synthetic `service.restart` event, and calls
/// `connect()` again — correctly modelling crash-loop cycles without any
/// extra bookkeeping in this impl.
///
/// Docker daemon unreachable: `spawn()` succeeds (we launched the docker CLI),
/// but the CLI exits immediately with an error message on stderr. That message
/// lands in scryer as a log line, then the channel closes → `StreamBroken` →
/// supervisor backs off. Callers see the error in the event store.
pub struct DockerLogSource {
    /// Lines of history to replay on each connect. `"0"` = no history (tail
    /// from now); `"all"` = full history. Defaults to `"50"` so crash-loop
    /// causes are visible without flooding scryer on the first connect.
    tail: String,
}

impl DockerLogSource {
    /// Create with the default 50-line history tail.
    pub fn new() -> std::sync::Arc<Self> {
        std::sync::Arc::new(Self { tail: "50".to_string() })
    }

    /// Create with an explicit tail line count.
    pub fn with_tail(n: u32) -> std::sync::Arc<Self> {
        std::sync::Arc::new(Self { tail: n.to_string() })
    }

    /// Follow from the current position; emit no prior history. Useful when
    /// the caller is already tracking the last-seen cursor.
    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);

        // Drain stdout lines
        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;
                }
            }
        });

        // Drain stderr lines (docker logs --follow sends container stderr here)
        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;
                }
            }
        });

        // When the child process exits, both pipe streams eventually close.
        // Drop the last `tx` clone once the child is gone so the channel
        // closes cleanly → adapter returns StreamBroken → supervisor restarts.
        tokio::spawn(async move {
            let _ = child.wait().await;
            drop(tx);
        });

        Ok(rx)
    }
}

// ─── Test helpers ─────────────────────────────────────────────────────────────

#[cfg(test)]
pub mod test_source {
    use super::*;
    use std::sync::Mutex;

    /// Test-only `ContainerLogSource` that hands out pre-canned line streams.
    ///
    /// Each call to `connect` consumes the next batch from `streams`. After
    /// the last batch the source returns a permanently-closed receiver so
    /// the adapter exits with `StreamBroken` and the supervisor can give up.
    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;
                    }
                }
                // Drop tx → receiver closes → adapter returns StreamBroken.
            });
            Ok(rx)
        }
    }
}

#[cfg(test)]
mod tests {
    use super::test_source::ScriptedSource;
    use super::*;

    /// Check docker is reachable for live tests.
    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;
        }
        // Run a short-lived container that prints two lines then exits.
        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; }

        // Give the container a moment to print.
        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;

        // Container may have exited already (--rm removes it); a missing container
        // returns an error from docker → Permanent error is OK here.
        match rx {
            Ok(mut rx) => {
                let mut lines = vec![];
                while let Some(line) = rx.recv().await {
                    lines.push(line);
                }
                // Lines should include our echo output
                let joined = lines.join("\n");
                // Partial assertion: if we got anything it worked
                eprintln!("[docker-source] got {} lines: {:?}", lines.len(), joined);
            }
            Err(AdapterError::Permanent(_)) => {
                // Container already gone (--rm) — acceptable
            }
            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,
        );

        // Single run — first connect succeeds, then source closes (StreamBroken expected).
        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());
        // Two batches: first has one line, then breaks; supervisor restarts;
        // second batch has one more line, then breaks again; supervisor
        // exhausts attempts and exits.
        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());
        // Supervisor runs until max_attempts (3) exhausts.
        let _ = supervisor.run(&mut adapter).await;

        scryer.flush_ring().unwrap();
        let events = scryer.events(&scope, &EventFilter::default()).await.unwrap();

        // Expected events in this scope:
        // - "before restart" (line)
        // - "service.restart" synth (1st restart)
        // - "after restart" (line)
        // - "service.restart" synth (2nd restart) -> source exhausted
        // - "service.restart" synth (3rd restart, exhausted attempts)
        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"
        );
    }
}