Start at Time Delta
Subscribe to a KubeMQ Events Store channel and replay events from a relative time offset using the C++ SDK.
Overview
A time-delta subscription starts replay from a relative offset — "the last 30 minutes" — 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.
kubemq::SubscriptionOption::StartFromTimeDelta(std::chrono::seconds(1800)) 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. A delta of zero effectively replays nothing and delivers only future events. 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 - C++ SDK installed (vcpkg or CMake FetchContent)
- C++17 compiler (GCC 9+, Clang 9+, MSVC 2019+)
Code
// Example: events_store/start_at_time_delta
//
// Demonstrates subscribing to event store with StartFromTimeDelta.
// Events are replayed from (now - delta). For example, 30 minutes ago.
//
// Channel: cpp-events-store.start-at-time-delta
// Client ID: cpp-events-store-start-at-time-delta-client
//
// Run with a KubeMQ server on localhost:50000
// (see https://docs.kubemq.io/deploy).
#include <kubemq/kubemq.h>
#include <atomic>
#include <chrono>
#include <iostream>
#include <thread>
int main() {
std::cout << "[1] Connecting to localhost:50000" << std::endl;
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
options.set_client_id("cpp-events-store-start-at-time-delta-client");
auto client_result = kubemq::Client::Create(options);
if (!client_result.ok()) {
std::cerr << "[ERROR] Failed to create client: " << client_result.status().message()
<< std::endl;
return 1;
}
auto& client = *client_result;
std::string channel = "cpp-events-store.start-at-time-delta";
std::atomic<int> received_count{0};
// Subscribe starting from 30 minutes ago (1800 seconds).
std::cout << "[2] Subscribing with StartFromTimeDelta(30 minutes)" << std::endl;
auto sub_result = client->SubscribeToEventsStore(
channel, "", kubemq::SubscriptionOption::StartFromTimeDelta(std::chrono::seconds(1800)),
[&received_count](const kubemq::EventStoreReceive& e) {
std::cout << "[4] [StartFromTimeDelta] seq=" << e.sequence << " body=" << e.body
<< std::endl;
received_count.fetch_add(1);
},
[](const kubemq::Status& err) {
std::cerr << "[ERROR] Subscription error: " << err.message() << std::endl;
});
if (!sub_result.ok()) {
std::cerr << "[ERROR] SubscribeToEventsStore: " << sub_result.status().message()
<< std::endl;
return 1;
}
auto& sub = *sub_result;
// Allow subscription to establish.
std::this_thread::sleep_for(std::chrono::seconds(1));
// Send an event that falls within the time delta window.
std::cout << "[3] Sending event within delta window" << std::endl;
auto ev_result =
kubemq::EventStore::Builder().SetChannel(channel).SetBody("recent event").Build();
if (!ev_result.ok()) {
std::cerr << "[ERROR] Build: " << ev_result.status().message() << std::endl;
return 1;
}
auto send_result = client->SendEventStore(*ev_result);
if (!send_result.ok()) {
std::cerr << "[ERROR] SendEventStore: " << send_result.status().message() << std::endl;
return 1;
}
// Wait for events to arrive.
std::this_thread::sleep_for(std::chrono::seconds(2));
std::cout << "[5] Received " << received_count.load() << " event(s)" << std::endl;
std::cout << "[6] Start at time delta demo complete" << std::endl;
// Cancel the subscription explicitly.
// Note: The Subscription destructor also calls Cancel(), but explicit
// cleanup is shown here for clarity and to match Go's defer pattern.
sub->Cancel();
auto close_status = client->Close();
if (!close_status.ok()) {
std::cerr << "[ERROR] Close failed: " << close_status.message() << std::endl;
return 1;
}
std::cout << "[7] Client closed" << std::endl;
return 0;
}How It Works
- Subscribes with
StartFromTimeDelta(std::chrono::seconds(1800))for 30 minutes ago. - Events stored within the delta window are replayed.
- Simpler than
StartFromTime()when you need a relative offset rather than absolute time. - New events continue to be delivered after replay completes.
Related
Was this page helpful?