use observation::{Event, EventSource, Level, TaskRunId};
use std::collections::HashMap;
use workload_spec::{ImageRef, MeshIdent};
pub mod pino;
pub mod tracing_json;
pub mod unstructured;
pub mod vanilla;
pub use pino::{PinoBeholder, PinoFactory};
pub use tracing_json::{TracingJsonBeholder, TracingJsonFactory};
pub use unstructured::UnstructuredBeholder;
pub use vanilla::VanillaBeholder;
#[derive(Debug, Clone)]
pub struct LogLine {
pub line: String,
pub offset_ms: u32,
}
#[derive(Debug, Clone, Default)]
pub struct ServiceHints {
pub image_labels: HashMap<String, String>,
pub env: HashMap<String, String>,
}
impl ServiceHints {
pub fn forced_beholder(&self) -> Option<&str> {
self.image_labels.get("yah.beholder").map(|s| s.as_str())
}
}
pub struct BeholderCtx {
pub seq: u32,
pub run_id_anchor: TaskRunId,
}
impl BeholderCtx {
pub fn new() -> Self {
Self { seq: 0, run_id_anchor: TaskRunId::new() }
}
pub fn next_seq(&mut self) -> u32 {
let s = self.seq;
self.seq += 1;
s
}
pub fn make_event(
&mut self,
line: &LogLine,
level: Level,
target: String,
msg: String,
fields: serde_json::Value,
source: EventSource,
) -> Event {
Event {
run_id: self.run_id_anchor.clone(),
seq: self.next_seq(),
offset_ms: line.offset_ms,
level,
target,
msg,
fields,
anchor: None,
source,
}
}
}
impl Default for BeholderCtx {
fn default() -> Self {
Self::new()
}
}
pub trait ServiceBeholder: Send {
fn name(&self) -> &'static str;
fn version(&self) -> &'static str;
fn matches(&self, ident: &MeshIdent, image: &ImageRef, hints: &ServiceHints) -> bool;
fn parse_line(&mut self, line: &LogLine, ctx: &mut BeholderCtx) -> Vec<Event>;
fn unknown_format_reason(&self) -> Option<&str> {
None
}
}
pub struct ServiceBeholderRegistry {
factories: Vec<Box<dyn ServiceBeholderFactory>>,
}
pub trait ServiceBeholderFactory: Send + Sync {
fn name(&self) -> &'static str;
fn version(&self) -> &'static str;
fn matches(&self, ident: &MeshIdent, image: &ImageRef, hints: &ServiceHints) -> bool;
fn create(&self) -> Box<dyn ServiceBeholder>;
}
impl ServiceBeholderRegistry {
pub fn new() -> Self {
Self { factories: Vec::new() }
}
pub fn register(&mut self, factory: Box<dyn ServiceBeholderFactory>) {
self.factories.push(factory);
}
pub fn attach(
&self,
ident: &MeshIdent,
image: &ImageRef,
hints: &ServiceHints,
) -> Option<(Box<dyn ServiceBeholder>, &'static str)> {
if let Some(forced) = hints.forced_beholder() {
if let Some(f) = self.factories.iter().find(|f| f.name() == forced) {
return Some((f.create(), f.name()));
}
}
for f in &self.factories {
if f.matches(ident, image, hints) {
return Some((f.create(), f.name()));
}
}
None
}
pub fn with_bundled() -> Self {
let mut r = Self::new();
r.register(Box::new(pino::PinoFactory));
r.register(Box::new(tracing_json::TracingJsonFactory));
r.register(Box::new(vanilla::VanillaFactory));
r.register(Box::new(unstructured::UnstructuredFactory));
r
}
}
impl Default for ServiceBeholderRegistry {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn no_hints() -> ServiceHints {
ServiceHints::default()
}
fn label_hints(beholder: &str) -> ServiceHints {
let mut h = ServiceHints::default();
h.image_labels.insert("yah.beholder".to_string(), beholder.to_string());
h
}
fn env_hints(key: &str, val: &str) -> ServiceHints {
let mut h = ServiceHints::default();
h.env.insert(key.to_string(), val.to_string());
h
}
fn ident(s: &str) -> MeshIdent {
MeshIdent(s.to_string())
}
fn image() -> ImageRef {
ImageRef {
registry: "ghcr.io".to_string(),
repository: "test/svc".to_string(),
tag: "latest".to_string(),
digest: workload_spec::testing::test_digest(),
}
}
#[test]
fn image_label_forces_pino() {
let reg = ServiceBeholderRegistry::with_bundled();
let hints = label_hints("pino");
let result = reg.attach(&ident("svc.a"), &image(), &hints);
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, pino::NAME);
}
#[test]
fn image_label_forces_tracing_json() {
let reg = ServiceBeholderRegistry::with_bundled();
let hints = label_hints("tracing-json");
let result = reg.attach(&ident("svc.a"), &image(), &hints);
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, tracing_json::NAME);
}
#[test]
fn image_label_forces_vanilla() {
let reg = ServiceBeholderRegistry::with_bundled();
let hints = label_hints("vanilla");
let result = reg.attach(&ident("svc.a"), &image(), &hints);
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, "vanilla");
}
#[test]
fn unknown_image_label_falls_through() {
let reg = ServiceBeholderRegistry::with_bundled();
let hints = label_hints("not-a-real-beholder");
let result = reg.attach(&ident("svc.a"), &image(), &hints);
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, "unstructured");
}
#[test]
fn no_hints_falls_through_to_unstructured() {
let reg = ServiceBeholderRegistry::with_bundled();
let result = reg.attach(&ident("svc.a"), &image(), &no_hints());
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, "unstructured");
}
#[test]
fn node_env_selects_pino() {
let reg = ServiceBeholderRegistry::with_bundled();
let hints = env_hints("NODE_ENV", "production");
let result = reg.attach(&ident("svc.a"), &image(), &hints);
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, pino::NAME);
}
#[test]
fn rust_log_selects_tracing_json() {
let reg = ServiceBeholderRegistry::with_bundled();
let hints = env_hints("RUST_LOG", "info");
let result = reg.attach(&ident("svc.a"), &image(), &hints);
assert!(result.is_some());
let (_, name) = result.unwrap();
assert_eq!(name, tracing_json::NAME);
}
}