use crate::service::{Scryer, ScryerError};
use async_trait::async_trait;
use observation::{Event, EventScope, EventSource, Level, TaskRunId};
use serde_json::Value;
use std::sync::Arc;
use std::time::Duration;
use thiserror::Error;
use tokio::time::sleep;
use workload_spec::MeshIdent;
pub mod containerd_logs;
pub mod journald;
pub mod warden_rpc;
pub use containerd_logs::{ContainerLogSource, ContainerdLogsAdapter, DockerLogSource};
pub use journald::{JournaldAdapter, JournaldEntry, JournaldSource};
pub use warden_rpc::{WardenEvent, WardenRpcAdapter};
#[derive(Debug, Error)]
pub enum AdapterError {
#[error("stream broken: {0}")]
StreamBroken(String),
#[error("permanent: {0}")]
Permanent(String),
#[error("scryer push: {0}")]
Push(#[from] ScryerError),
}
impl AdapterError {
pub fn is_recoverable(&self) -> bool {
matches!(self, AdapterError::StreamBroken(_) | AdapterError::Push(_))
}
}
#[async_trait]
pub trait Adapter: Send {
fn name(&self) -> &str;
fn scope(&self) -> EventScope;
async fn run(&mut self) -> Result<(), AdapterError>;
}
#[derive(Debug, Clone)]
pub struct BackoffConfig {
pub initial: Duration,
pub max: Duration,
pub multiplier: f64,
pub max_attempts: Option<u32>,
}
impl Default for BackoffConfig {
fn default() -> Self {
Self {
initial: Duration::from_secs(1),
max: Duration::from_secs(30),
multiplier: 2.0,
max_attempts: None,
}
}
}
impl BackoffConfig {
pub fn test() -> Self {
Self {
initial: Duration::from_millis(0),
max: Duration::from_millis(0),
multiplier: 1.0,
max_attempts: Some(3),
}
}
}
pub struct Supervisor {
name: String,
scope: EventScope,
scryer: Arc<Scryer>,
cfg: BackoffConfig,
restart_count: u32,
synth_seq: u32,
}
impl Supervisor {
pub fn new(scryer: Arc<Scryer>, scope: EventScope, name: impl Into<String>) -> Self {
Self {
name: name.into(),
scope,
scryer,
cfg: BackoffConfig::default(),
restart_count: 0,
synth_seq: u32::MAX,
}
}
pub fn with_backoff(mut self, cfg: BackoffConfig) -> Self {
self.cfg = cfg;
self
}
pub async fn run(&mut self, adapter: &mut dyn Adapter) -> Result<(), AdapterError> {
let mut delay = self.cfg.initial;
let mut attempt: u32 = 0;
loop {
let result = adapter.run().await;
match result {
Ok(()) => return Ok(()),
Err(e) if matches!(e, AdapterError::Permanent(_)) => return Err(e),
Err(_) => {
attempt += 1;
if let Some(cap) = self.cfg.max_attempts {
if attempt > cap {
return Err(AdapterError::StreamBroken(format!(
"supervisor exhausted {cap} attempts for {}",
self.name
)));
}
}
self.emit_restart_event()?;
if delay > Duration::ZERO {
sleep(delay).await;
}
delay = (delay.mul_f64(self.cfg.multiplier)).min(self.cfg.max);
if delay == Duration::ZERO {
delay = self.cfg.initial;
}
}
}
}
}
fn emit_restart_event(&mut self) -> Result<(), AdapterError> {
self.restart_count += 1;
let seq = self.synth_seq;
self.synth_seq = self.synth_seq.saturating_sub(1);
let event = Event {
run_id: TaskRunId::new(),
seq,
offset_ms: 0,
level: Level::Info,
target: format!("scryer.{}", self.name),
msg: "service.restart".to_string(),
fields: serde_json::json!({ "attempt": self.restart_count, "adapter": self.name }),
anchor: None,
source: EventSource::Synth,
};
self.scryer.push(self.scope.clone(), event)?;
Ok(())
}
}
pub(crate) fn synth_event(
target: impl Into<String>,
level: Level,
msg: impl Into<String>,
fields: Value,
offset_ms: u32,
seq: u32,
) -> Event {
Event {
run_id: TaskRunId::new(),
seq,
offset_ms,
level,
target: target.into(),
msg: msg.into(),
fields,
anchor: None,
source: EventSource::Synth,
}
}
#[allow(dead_code)]
pub(crate) fn warden_local_scope() -> EventScope {
EventScope::Service(MeshIdent("yubaba.local".to_string()))
}