use std::env;
use std::error::Error;
use tephra_client::{AsyncClient, Event, Position, Query, SubEvent};
use tokio_stream::StreamExt;
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let addr = env::args()
.nth(1)
.unwrap_or_else(|| "127.0.0.1:9000".to_string());
let client = AsyncClient::connect(&addr).await?;
println!("connected to {addr}");
let (a, b, c) = tokio::join!(
client.append([Event::new("Enrolled", &["course:c1"], b"{}")?], None),
client.append([Event::new("Enrolled", &["course:c2"], b"{}")?], None),
client.append([Event::new("Enrolled", &["course:c3"], b"{}")?], None),
);
for range in [a?, b?, c?] {
println!("appended at position {}", range.first);
}
let (events, watermark) = client.read_all(Query::all(), Position::ZERO, None).await?;
println!("log holds {} events (watermark {watermark}):", events.len());
for sequenced in &events {
let ev = sequenced.event();
let tags: Vec<&str> = ev.tags().collect();
println!(" {} {} {tags:?}", sequenced.position(), ev.event_type());
}
println!("subscribing until caught up:");
let mut sub = client.subscribe(Query::all(), Position::ZERO).await;
while let Some(item) = sub.next().await {
match item? {
SubEvent::Event(sequenced) => {
println!(
" live {} {}",
sequenced.position(),
sequenced.event().event_type()
);
}
SubEvent::CaughtUp(watermark) => {
println!(" caught up at {watermark}");
break;
}
}
}
Ok(())
}