Start at Time Delta
Subscribe to KubeMQ Events Store from a relative time offset such as 30 seconds ago using the Rust SDK.
Overview
A time-delta subscription starts replay from a relative offset — "30 seconds ago" — instead of a fixed timestamp or sequence number. It's the right tool when a consumer knows how long it was offline but not the exact moment it disconnected: a worker restarting after a deploy, a dashboard reconnecting after a blip, or a batch job that only cares about "recent" history. Computing an absolute cutoff yourself is bookkeeping the broker can do for you.
EventsStoreSubscription::StartAtTimeDelta(Duration::from_secs(30)) passes the duration to the broker, which resolves it to now - delta at subscription time, replays every stored event from that point forward, then hands off to live delivery — the same replay-to-live transition as an absolute-time or sequence-based start.
Gotchas: the delta is evaluated once, server-side, at subscription creation — it does not "slide" as time passes. The server interprets the Duration in whole seconds, so sub-second precision is truncated. And since the window is wall-clock based, clock skew between producers and the broker can shift which events land inside or outside the boundary.
Prerequisites
- KubeMQ server running on
localhost:50000 - Rust SDK installed (
cargo add kubemq)
Code
use kubemq::prelude::*;
use kubemq::EventsStoreSubscription;
use std::time::Duration;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.build()
.await?;
let channel = "rust-events-store.start-at-time-delta";
let sub = client.subscribe_to_events_store(
channel, "",
EventsStoreSubscription::StartAtTimeDelta(Duration::from_secs(30)),
|event| Box::pin(async move {
println!("Seq {}: {}", event.sequence, String::from_utf8_lossy(&event.body));
}),
None,
).await?;
tokio::time::sleep(Duration::from_secs(5)).await;
sub.unsubscribe().await;
client.close().await?;
Ok(())
}How It Works
StartAtTimeDeltaspecifies a relative offset from the current time.- The server interprets the value in whole seconds (sub-second precision is truncated).
- Review timeouts, channel names, and client IDs before running against shared environments.
- Run the program while the server from the prerequisites is available.
Related
Was this page helpful?