use std::error::Error;
use std::time::Duration;
use ruststream::memory::{MemoryBroker, MemoryPosition, MemorySeeker, MemorySource};
use ruststream::runtime::{AppInfo, HandlerResult, RustStream, Seek};
use ruststream::{OutgoingMessage, Publisher, Seeker, subscriber};
use serde::{Deserialize, Serialize};
use tokio::time::sleep;
#[derive(Debug, Serialize, Deserialize)]
struct Job {
id: u64,
}
#[subscriber("audit", start_at(MemoryPosition::start()))]
async fn record(entry: &Job) -> HandlerResult {
println!("audit: entry {}", entry.id);
HandlerResult::Ack
}
#[subscriber(MemorySource::new("jobs"))]
async fn work(job: &Job, Seek(seeker): Seek<MemorySeeker>) -> HandlerResult {
if job.id == 999 {
if seeker.seek(MemoryPosition::sequence(3)).await.is_err() {
return HandlerResult::retry();
}
return HandlerResult::Ack;
}
println!("jobs: processed {}", job.id);
HandlerResult::Ack
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let broker = MemoryBroker::new();
let ingress = broker.publisher();
for id in 1..=2u64 {
let payload = serde_json::to_vec(&Job { id })?;
ingress
.publish(OutgoingMessage::new("audit", payload.as_slice()))
.await?;
}
let app = RustStream::new(AppInfo::new("seek-demo", "0.1.0")).with_broker(broker, |b| {
b.include(work);
b.include(record);
});
let running = app.start().await?;
for id in [1, 999, 3, 4] {
let payload = serde_json::to_vec(&Job { id })?;
ingress
.publish(OutgoingMessage::new("jobs", payload.as_slice()))
.await?;
}
sleep(Duration::from_millis(100)).await;
running.shutdown().await?;
println!("ok: replayed the audit history and skipped the poisoned region");
Ok(())
}