use async_nats::jetstream::stream::Config;
use futures_util::StreamExt;
use nats_counters::CounterExt;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let client = async_nats::connect("localhost:4222").await?;
let js = async_nats::jetstream::new(client);
let config = Config {
name: "METRICS".to_string(),
subjects: vec!["metrics.>".to_string()],
allow_message_counter: true,
allow_direct: true,
..Default::default()
};
let _stream = js.get_or_create_stream(config).await?;
let counter = js.get_counter("METRICS").await?;
println!("=== Adding Counter Values ===\n");
let services = ["auth", "api", "db", "cache"];
let metrics = ["requests", "errors", "latency_ms"];
for service in &services {
for metric in &metrics {
let subject = format!("metrics.{}.{}", service, metric);
let value = match *metric {
"requests" => 100 + (service.len() as i64 * 10),
"errors" => service.len() as i64,
"latency_ms" => 50 + (service.len() as i64 * 5),
_ => 0,
};
let result = counter.add(subject.clone(), value).await?;
println!(" {} = {}", subject, result);
}
}
println!("\n=== Fetching Multiple Counters (Batch) ===\n");
let mut subjects = Vec::new();
for service in &services {
for metric in &metrics {
subjects.push(format!("metrics.{}.{}", service, metric));
}
}
let mut entries = counter.get_multiple(subjects.clone()).await?;
println!("Fetched {} subjects in batch:\n", subjects.len());
while let Some(entry) = entries.next().await {
let entry = entry?;
println!(" {} = {}", entry.subject, entry.value);
if let Some(increment) = entry.increment {
println!(" (last increment: {})", increment);
}
}
println!("\n=== Service Totals ===\n");
for service in &services {
let service_subjects: Vec<String> = metrics
.iter()
.map(|m| format!("metrics.{}.{}", service, m))
.collect();
let mut service_entries = counter.get_multiple(service_subjects).await?;
let mut total_requests = 0i64;
let mut total_errors = 0i64;
let mut total_latency = 0i64;
while let Some(entry) = service_entries.next().await {
let entry = entry?;
if entry.subject.ends_with(".requests") {
total_requests = entry.value.try_into().unwrap_or(0);
} else if entry.subject.ends_with(".errors") {
total_errors = entry.value.try_into().unwrap_or(0);
} else if entry.subject.ends_with(".latency_ms") {
total_latency = entry.value.try_into().unwrap_or(0);
}
}
let error_rate = if total_requests > 0 {
(total_errors as f64 / total_requests as f64) * 100.0
} else {
0.0
};
println!(
" Service '{}': {} requests, {} errors ({:.2}%), {}ms avg latency",
service, total_requests, total_errors, error_rate, total_latency
);
}
println!("\n=== Partial Batch Fetch ===\n");
let error_subjects: Vec<String> = services
.iter()
.map(|s| format!("metrics.{}.errors", s))
.collect();
let mut error_entries = counter.get_multiple(error_subjects).await?;
println!("Error counts by service:");
while let Some(entry) = error_entries.next().await {
let entry = entry?;
println!(" {} = {}", entry.subject, entry.value);
}
Ok(())
}