use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use agentplane::case::CaseStore;
use agentplane::core::{Outcome, Skill, SkillDescriptor, SkillError, Tainted};
use agentplane::journal::JournalStore;
use agentplane::runtime::effects::Recorded;
use agentplane::runtime::telemetry;
use agentplane::runtime::{Mode, RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use serde_json::{Value, json};
use tracing::field::{Field, Visit};
use tracing_subscriber::Layer;
use tracing_subscriber::layer::{Context, SubscriberExt};
use tracing_subscriber::registry::LookupSpan;
use tracing_subscriber::util::SubscriberInitExt;
#[derive(Debug, Default)]
struct Collected {
live_effects: Vec<String>,
replayed_effects: usize,
metrics: BTreeMap<String, BTreeMap<String, f64>>,
}
type Sink = Arc<Mutex<Collected>>;
struct DropReplayedEffects(Sink);
impl<S> Layer<S> for DropReplayedEffects
where
S: tracing::Subscriber + for<'a> LookupSpan<'a>,
{
fn on_new_span(
&self,
attrs: &tracing::span::Attributes<'_>,
_id: &tracing::span::Id,
_ctx: Context<'_, S>,
) {
if attrs.metadata().name() != telemetry::EFFECT_SPAN {
return;
}
let mut fields = Fields::default();
attrs.record(&mut fields);
if fields.bools.get(telemetry::EFFECT_REPLAYED).copied() == Some(true) {
self.0.lock().expect("collector").replayed_effects += 1;
return;
}
let kind = fields
.strings
.get(telemetry::EFFECT_KIND)
.cloned()
.unwrap_or_else(|| "unknown".to_owned());
self.0.lock().expect("collector").live_effects.push(kind);
}
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
if event.metadata().target() != telemetry::EFFECT_SPAN {
return;
}
let mut fields = Fields::default();
event.record(&mut fields);
if fields.bools.get("replayed").copied() == Some(true) {
self.0.lock().expect("collector").replayed_effects += 1;
}
}
}
struct MetricBridge(Sink);
impl<S: tracing::Subscriber> Layer<S> for MetricBridge {
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
if event.metadata().target() != agentplane::runtime::metrics::METRIC {
return;
}
let mut fields = Fields::default();
event.record(&mut fields);
let Some(metric) = fields.strings.get("metric") else {
return;
};
let dim = fields
.strings
.get("dim")
.cloned()
.unwrap_or_else(|| "-".to_owned());
let value = fields.numbers.get("value").copied().unwrap_or_default();
let mut out = self.0.lock().expect("collector");
let series = out.metrics.entry(metric.clone()).or_default();
let is_gauge = agentplane::runtime::metrics::CATALOGUE
.iter()
.find(|i| i.name == metric.as_str())
.is_some_and(|i| matches!(i.kind, agentplane::runtime::metrics::Kind::Gauge));
let slot = series.entry(dim).or_default();
if is_gauge {
*slot = value;
} else {
*slot += value;
}
}
}
#[derive(Default)]
struct Fields {
strings: BTreeMap<String, String>,
bools: BTreeMap<String, bool>,
numbers: BTreeMap<String, f64>,
}
#[allow(clippy::cast_precision_loss)]
impl Visit for Fields {
fn record_str(&mut self, field: &Field, value: &str) {
self.strings
.insert(field.name().to_owned(), value.to_owned());
}
fn record_bool(&mut self, field: &Field, value: bool) {
self.bools.insert(field.name().to_owned(), value);
}
fn record_u64(&mut self, field: &Field, value: u64) {
self.numbers.insert(field.name().to_owned(), value as f64);
}
fn record_i64(&mut self, field: &Field, value: i64) {
self.numbers.insert(field.name().to_owned(), value as f64);
}
fn record_f64(&mut self, field: &Field, value: f64) {
self.numbers.insert(field.name().to_owned(), value);
}
fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
self.strings.insert(
field.name().to_owned(),
format!("{value:?}").trim_matches('"').to_owned(),
);
}
}
#[derive(Debug)]
struct Post;
#[async_trait::async_trait]
impl Skill for Post {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("ledger.post").provides("ledger.post")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
_input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let arguments = Tainted::trusted(json!(null));
let posted = cx.sink(Recorded::new("ledger.post"), &arguments).await?;
Ok(Outcome::done(posted))
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let collected: Sink = Arc::default();
tracing_subscriber::registry()
.with(DropReplayedEffects(Arc::clone(&collected)))
.with(MetricBridge(Arc::clone(&collected)))
.init();
let store = Arc::new(RedbStore::open_in_memory()?);
let plane = Runtime::builder(Arc::clone(&store) as Arc<dyn JournalStore>)
.cases(Arc::clone(&store) as Arc<dyn CaseStore>)
.skill(Post)
.build();
let out = plane
.run("ledger.post", Tainted::trusted(json!({})))
.await?;
assert_eq!(out.status, RunStatus::Succeeded);
let replayed = plane.replay(out.run_id, Mode::Strict).await?;
assert_eq!(replayed.status, RunStatus::Succeeded);
{
let seen = collected.lock().expect("collector");
println!("1. latency series");
println!(" live effect spans : {:?}", seen.live_effects);
println!(" replays, not timed : {}", seen.replayed_effects);
assert_eq!(
seen.live_effects.len(),
1,
"exactly one effect reached the world, and only it belongs in a \
latency histogram"
);
assert!(
seen.replayed_effects >= 1,
"the replay must be visible as a replay: 'how much of this plane's \
work is recovery' is a question an operator asks, and the answer \
is not zero just because it is not latency"
);
}
{
let seen = collected.lock().expect("collector");
assert!(
!seen
.metrics
.contains_key(agentplane::runtime::metrics::OPEN_CASES.name),
"gauges must not appear before something queries the stores"
);
}
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let report = plane
.sweep(now, std::time::Duration::from_secs(3600))
.await?;
println!("\n2. the census, from the sweep");
println!(" open cases : {}", report.census.open_cases);
println!(" open tasks : {}", report.census.open_tasks);
println!(" due deadlines: {}", report.census.due_deadlines);
{
let seen = collected.lock().expect("collector");
assert!(
seen.metrics
.contains_key(agentplane::runtime::metrics::OPEN_CASES.name),
"the sweep must emit the census, or a plane's gauges never exist"
);
}
println!("\n3. alerting");
println!(" needs_attention(): {}", report.needs_attention());
println!(" is_quiet() : {}", report.is_quiet());
if report.needs_attention() {
println!(" → page somebody: {report:?}");
}
println!(
"\n census_unavailable is part of needs_attention() on purpose: gauges\n\
\x20 that could not be read are a blind spot wearing a default, and\n\
\x20 'zero open cases' and 'I could not count them' must not look alike."
);
println!("\n4. what a meter would hold");
for (metric, series) in &collected.lock().expect("collector").metrics {
for (dim, value) in series {
println!(" {metric:<34} {dim:<12} {value}");
}
}
Ok(())
}