KubeMQ
Client SDKsRustHow-to guidesEvents Store

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

main.rs
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

  • StartAtTimeDelta specifies 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.

Was this page helpful?

On this page