use std::path::Path;
use dbmd_core::linkmd::{self, LinkError};
use crate::cli::SubscribeArgs;
use crate::context::Context;
use crate::error::{CliError, CliResult, ExitCode};
use crate::sanitize::sanitize;
pub fn run(ctx: &Context, args: &SubscribeArgs) -> CliResult {
if args.interval == 0 {
return Err(CliError::new(
ExitCode::Runtime,
"BAD_INTERVAL",
"--interval must be at least 1 second",
));
}
let brain = args.brain.trim().trim_start_matches('@');
let cfg = linkmd::hub_config(args.hub.as_deref(), Path::new(&args.dir))?;
let head = linkmd::head(&cfg, brain)?;
let mut baseline = args.since.unwrap_or(head.seq);
emit(ctx, &head, baseline);
if args.once {
return Ok(());
}
baseline = baseline.max(head.seq);
loop {
std::thread::sleep(std::time::Duration::from_secs(args.interval));
match linkmd::head(&cfg, brain) {
Ok(h) => {
if h.seq > baseline {
emit(ctx, &h, baseline);
baseline = h.seq;
}
}
Err(LinkError::Transport { hub, message }) => {
eprintln!("dbmd: subscribe: hub unreachable at {hub} ({message}); retrying");
}
Err(e) => return Err(e.into()),
}
}
}
fn emit(ctx: &Context, head: &linkmd::Head, prev: u64) {
if ctx.json {
let mut v = serde_json::to_value(head).unwrap_or_default();
v["prev"] = serde_json::json!(prev);
println!("{v}");
} else if head.seq == prev {
println!("{} at feed seq {}", sanitize(&head.brain), head.seq);
} else {
println!(
"{} advanced: feed seq {} -> {}",
sanitize(&head.brain),
prev,
head.seq
);
}
}