use async_nats::jetstream;
use nats_counters::CounterExt;
use num_bigint::BigInt;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let client = async_nats::connect("nats://localhost:4222").await?;
let js = jetstream::new(client);
let stream_config = jetstream::stream::Config {
name: "EVENTS_COUNTER".to_string(),
subjects: vec!["events.*".to_string()],
allow_message_counter: true,
allow_direct: true,
..Default::default()
};
let _ = js.create_stream(stream_config).await;
let counter = js.get_counter("EVENTS_COUNTER").await?;
println!("=== NATS JetStream Counter Example ===\n");
println!("Recording events...");
let clicks = counter.add("events.clicks", 5).await?;
println!(" Clicks: +5 = {}", clicks);
let views = counter.add("events.views", 100).await?;
println!(" Views: +100 = {}", views);
let errors = counter.add("events.errors", 2).await?;
println!(" Errors: +2 = {}", errors);
println!("\nRecording more events...");
let clicks = counter.add("events.clicks", 3).await?;
println!(" Clicks: +3 = {}", clicks);
let views = counter.add("events.views", 50).await?;
println!(" Views: +50 = {}", views);
let errors = counter.add("events.errors", BigInt::from(-1)).await?;
println!(" Errors: -1 = {}", errors);
println!("\nCurrent totals:");
println!(" Total clicks: {}", counter.load("events.clicks").await?);
println!(" Total views: {}", counter.load("events.views").await?);
println!(" Total errors: {}", counter.load("events.errors").await?);
println!("\nDetailed entry for clicks:");
let entry = counter.get("events.clicks").await?;
println!(" Subject: {}", entry.subject);
println!(" Value: {}", entry.value);
if let Some(increment) = entry.increment {
println!(" Last increment: {}", increment);
}
if !entry.sources.is_empty() {
println!(" Sources: {:?}", entry.sources);
}
println!("\n=== Counter operations completed ===");
Ok(())
}