KubeMQ
ConnectorsCloudEventsHow-to guides

Events Store

Publish persistent CloudEvents and subscribe with replay using the six start positions and Last-Event-ID resume over SSE.

The events-store pattern persists every CloudEvent to a durable, sequenced log so a subscriber can replay history, start at a specific point, or resume after a disconnect. It is the CloudEvents-over-HTTP equivalent of KubeMQ's Events Store messaging pattern.

Overview

Where a plain events channel is fire-and-forget, an events-store channel writes each event to disk with a monotonically increasing sequence number. Subscribers choose a start position when they connect — receive only new events, replay everything from the beginning, jump to a sequence number, or seek by time. Because every Server-Sent Events (SSE) frame carries the sequence number in its id: field, a subscriber that drops its connection can reconnect with the Last-Event-ID header and the server resumes from sequence + 1, automatically replaying whatever it missed.

You send persistent events to POST /ce/send/event-store and subscribe over SSE at GET /ce/subscribe/events-store. Both accept the same CloudEvent structured and binary content modes as a non-persistent event.

How it works

A publisher persists CloudEvents to the durable channel; a subscriber connects with a start position and the connector replays stored events, then streams new ones — each frame tagged with its sequence number for resume.

Persisted events get a sequence number; subscribers choose where to start and resume from the last sequence they saw.

Start positions

The GET /ce/subscribe/events-store endpoint accepts an events_store_type query parameter (default 1). Types 4, 5, and 6 also read events_store_value.

events_store_typeNameevents_store_valueBehavior
1StartNewOnlyOnly events that arrive after subscribing (default)
2StartFromFirstReplay from the first stored event
3StartFromLastStart from the last stored event
4StartAtSequencesequence numberStart at a specific sequence number
5StartAtTimeUnix secondsStart at an absolute timestamp
6StartAtTimeDeltasecondsStart at a time delta back from now

Only events-store subscriptions emit an id: field on each SSE frame (the sequence number), because only persisted events can be replayed. Plain events subscriptions never emit id:.

Resume with Last-Event-ID

When a subscriber reconnects with the Last-Event-ID header set to the last sequence number it processed, the server resumes the stream from sequence + 1. This is the standard SSE reconnection protocol — browsers and EventSource clients send Last-Event-ID automatically.

If both Last-Event-ID and events_store_type are present, the query parameter wins. To use Last-Event-ID for reconnection, omit events_store_type from the reconnect URL.

Send a persistent event

Publish a CloudEvent to an events-store channel. The subject attribute resolves to the channel; the connector responds with HTTP 202.

curl -X POST http://localhost:9090/ce/send/event-store \
  -H "Content-Type: application/cloudevents+json" \
  -d '{
    "specversion": "1.0",
    "type": "com.example.audit.entry",
    "source": "audit-service",
    "subject": "audit-log",
    "datacontenttype": "application/json",
    "data": {"action": "user.login", "user": "admin"}
  }'
using CloudNative.CloudEvents;
using CloudNative.CloudEvents.SystemTextJson;
using System.Net.Http.Headers;

var base_ = "http://localhost:9090";
var channel = "audit-log";
var formatter = new JsonEventFormatter();
using var httpClient = new HttpClient();

var ev = new CloudEvent {
    Id = Guid.NewGuid().ToString(),
    Type = "com.kubemq.examples.eventsstore.stored",
    Source = new Uri("urn:kubemq-ce-csharp-example"),
    Subject = channel,
    DataContentType = "application/json",
    Data = new { msg = "hello events-store from C#!" },
};
var bytes = formatter.EncodeStructuredModeMessage(ev, out var ct);
using var content = new ByteArrayContent(bytes.ToArray());
content.Headers.ContentType = MediaTypeHeaderValue.Parse(ct.ToString());
var resp = await httpClient.PostAsync($"{base_}/ce/send/event-store", content);
Console.WriteLine($"Published to events-store: status={resp.StatusCode}");
import (
	"encoding/json"
	"net/http"
	"strings"

	cloudevents "github.com/cloudevents/sdk-go/v2"
)

func sendEventStore(base, channel, msg string) error {
	event := cloudevents.NewEvent()
	event.SetType("com.kubemq.examples.eventsstore.stored")
	event.SetSource("kubemq-ce-go-example")
	event.SetSubject(channel)
	_ = event.SetData(cloudevents.ApplicationJSON, map[string]string{"msg": msg})

	body, _ := json.Marshal(event)
	req, _ := http.NewRequest("POST", base+"/ce/send/event-store", strings.NewReader(string(body)))
	req.Header.Set("Content-Type", "application/cloudevents+json")
	resp, err := http.DefaultClient.Do(req)
	if err != nil {
		return err
	}
	return resp.Body.Close()
}
EventFormatProvider.getInstance().registerFormat(new JsonFormat());
EventFormat format = EventFormatProvider.getInstance().resolveFormat(JsonFormat.CONTENT_TYPE);
CloudEvent event = CloudEventBuilder.v1()
        .withId(UUID.randomUUID().toString())
        .withType("com.kubemq.examples.eventsstore.stored")
        .withSource(URI.create("kubemq-ce-java-example"))
        .withSubject(channel)
        .withDataContentType("application/json")
        .withTime(OffsetDateTime.now())
        .withData("application/json",
                MAPPER.writeValueAsBytes(Map.of("msg", "hello events-store from Java!")))
        .build();

HttpClient httpClient = HttpClient.newHttpClient();
HttpResponse<String> resp = httpClient.send(
        HttpRequest.newBuilder().uri(URI.create(base + "/ce/send/event-store"))
                .POST(HttpRequest.BodyPublishers.ofByteArray(format.serialize(event)))
                .header("Content-Type", "application/cloudevents+json").build(),
        HttpResponse.BodyHandlers.ofString());
System.out.println("Published to events-store: status=" + resp.statusCode());
import { CloudEvent, HTTP } from 'cloudevents';

const base = 'http://localhost:9090';
const channel = 'audit-log';

const event = new CloudEvent({
  type: 'com.kubemq.examples.eventsstore.stored',
  source: 'kubemq-ce-js-example',
  subject: channel,
  datacontenttype: 'application/json',
  data: { msg: 'hello events-store from JS!' },
});
const msg = HTTP.structured(event);
const resp = await fetch(`${base}/ce/send/event-store`, {
  method: 'POST',
  headers: msg.headers as Record<string, string>,
  body: msg.body as string,
});
console.log(`Published to events-store: status=${resp.status}`);
import requests
from cloudevents.v1.conversion import to_structured
from cloudevents.v1.http import CloudEvent

base = "http://localhost:9090"
channel = "audit-log"

event = CloudEvent(
    attributes={
        "type": "com.kubemq.examples.eventsstore.stored",
        "source": "kubemq-ce-python-example",
        "subject": channel,
        "datacontenttype": "application/json",
    },
    data={"msg": "hello events-store from Python!"},
)
headers, body = to_structured(event)
resp = requests.post(f"{base}/ce/send/event-store", data=body,
                     headers=dict(headers), timeout=10)
print(f"Published to events-store: status={resp.status_code}")
require "net/http"; require "uri"; require "json"; require "cloud_events"; require "securerandom"

base = "http://localhost:9090"; channel = "audit-log"
sdk = CloudEvents::HttpBinding.default

event = CloudEvents::Event::V1.new(
  id: SecureRandom.uuid, type: "com.kubemq.examples.eventsstore.stored",
  source: URI("urn:kubemq-ce-ruby-example"), subject: channel,
  spec_version: "1.0",
  data_content_type: CloudEvents::ContentType.new("application/json"),
  data: JSON.generate({ msg: "hello events-store from Ruby!" })
)
enc_h, enc_b = sdk.encode_event(event, structured_format: "json")
uri = URI("#{base}/ce/send/event-store")
Net::HTTP.start(uri.host, uri.port) do |http|
  req = Net::HTTP::Post.new(uri); enc_h.each { |k, v| req[k] = v }; req.body = enc_b
  res = http.request(req)
  puts "Published to events-store: status=#{res.code}"
end
use cloudevents::{EventBuilder, EventBuilderV10};
use reqwest::Client;
use serde_json::json;
use uuid::Uuid;

let base = "http://localhost:9090";
let channel = "audit-log";
let client = Client::new();

let event = EventBuilderV10::new()
    .id(Uuid::new_v4().to_string())
    .ty("com.kubemq.examples.eventsstore.stored")
    .source("urn:kubemq-ce-rust-example")
    .subject(channel)
    .data("application/json", json!({"msg": "hello events-store from Rust!"}))
    .build()?;
let body = serde_json::to_string(&event)?;
let resp = client.post(format!("{}/ce/send/event-store", base))
    .header("Content-Type", "application/cloudevents+json")
    .body(body).send().await?;
println!("Published to events-store: status={}", resp.status());

Subscribe with StartNewOnly

The default subscription (events_store_type=1) delivers only events that arrive after the connection opens. The SSE handler reads event: cloudevent frames and parses each data: payload as a reconstructed CloudEvent.

curl -N "http://localhost:9090/ce/subscribe/events-store?client_id=my-client&channel=audit-log&events_store_type=1"
using System.Net.Http.Headers;
using System.Text;
using System.Text.Json;

var url = $"{base_}/ce/subscribe/events-store?client_id=csharp-es-sub&channel={Uri.EscapeDataString(channel)}&events_store_type=1";
using var http = new HttpClient { Timeout = Timeout.InfiniteTimeSpan };
using var req = new HttpRequestMessage(HttpMethod.Get, url);
req.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("text/event-stream"));
using var resp = await http.SendAsync(req, HttpCompletionOption.ResponseHeadersRead);
using var reader = new StreamReader(await resp.Content.ReadAsStreamAsync(), Encoding.UTF8);

string? evType = null, data = null, line;
while ((line = await reader.ReadLineAsync()) != null)
{
    if (line == "") {
        if (evType == "cloudevent" && data != null) {
            var ce = JsonSerializer.Deserialize<JsonElement>(data);
            Console.WriteLine($"Received: type={ce.GetProperty("type")} data={ce.GetProperty("data")}");
            return;
        }
        evType = null; data = null;
    }
    else if (line.StartsWith("event:")) evType = line[6..].Trim();
    else if (line.StartsWith("data:")) data = line[5..].Trim();
}
// events_store_type=1 = StartNewOnly — only messages arriving after subscribe.
sseURL := fmt.Sprintf(
	"%s/ce/subscribe/events-store?client_id=%s&channel=%s&events_store_type=1",
	base, clientID, channel)
req, _ := http.NewRequest("GET", sseURL, nil)
req.Header.Set("Accept", "text/event-stream")

resp, err := (&http.Client{Timeout: 0}).Do(req)
if err != nil {
	log.Fatal("SSE connect:", err)
}
defer resp.Body.Close()

scanner := bufio.NewScanner(resp.Body)
var evType, data string
for scanner.Scan() {
	line := scanner.Text()
	if line == "" {
		if evType == "cloudevent" && data != "" {
			var ce map[string]interface{}
			_ = json.Unmarshal([]byte(data), &ce)
			fmt.Printf("Received: type=%v data=%v\n", ce["type"], ce["data"])
			return
		}
		evType, data = "", ""
		continue
	}
	if strings.HasPrefix(line, ":") {
		continue
	}
	if strings.HasPrefix(line, "event:") {
		evType = strings.TrimSpace(strings.TrimPrefix(line, "event:"))
	} else if strings.HasPrefix(line, "data:") {
		data = strings.TrimSpace(strings.TrimPrefix(line, "data:"))
	}
}
// events_store_type=1 = StartNewOnly
String sseUrl = base + "/ce/subscribe/events-store?client_id=java-es-sub&channel="
        + channel + "&events_store_type=1";
HttpURLConnection conn = (HttpURLConnection) new URL(sseUrl).openConnection();
conn.setRequestProperty("Accept", "text/event-stream");
conn.setReadTimeout(15_000);
try (BufferedReader reader = new BufferedReader(
        new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) {
    String line; String evType = null, data = null;
    while ((line = reader.readLine()) != null) {
        if (line.isEmpty()) {
            if ("cloudevent".equals(evType) && data != null) {
                Map<?, ?> ce = MAPPER.readValue(data, Map.class);
                System.out.println("Received: type=" + ce.get("type") + " data=" + ce.get("data"));
                return;
            }
            evType = null; data = null;
        } else if (line.startsWith("event:")) evType = line.substring(6).trim();
        else if (line.startsWith("data:")) data = line.substring(5).trim();
    }
}
import EventSource from 'eventsource';

const url = `${base}/ce/subscribe/events-store?client_id=js-es-sub&channel=${encodeURIComponent(channel)}&events_store_type=1`;
const es = new EventSource(url);

es.addEventListener('cloudevent', (evt: MessageEvent) => {
  const ce = JSON.parse(evt.data) as Record<string, unknown>;
  console.log(`Received: type=${ce.type} data=${JSON.stringify(ce.data)}`);
  es.close();
});
es.addEventListener('error', (err) => {
  console.error('SSE error:', err);
  es.close();
});
import json
import requests

# events_store_type=1 = StartNewOnly
sse_url = (f"{base}/ce/subscribe/events-store"
           f"?client_id=python-es-sub&channel={channel}&events_store_type=1")
with requests.get(sse_url, stream=True, timeout=None,
                  headers={"Accept": "text/event-stream"}) as resp:
    ev_type = data = ""
    for line in resp.iter_lines(decode_unicode=True):
        if line == "":
            if ev_type == "cloudevent" and data:
                ce = json.loads(data)
                print(f"Received: type={ce.get('type')} data={ce.get('data')}")
                break
            ev_type = data = ""
            continue
        if line.startswith(":"):
            continue
        if line.startswith("event:"):
            ev_type = line[6:].strip()
        elif line.startswith("data:"):
            data = line[5:].strip()
require "net/http"; require "uri"; require "json"

uri = URI("#{base}/ce/subscribe/events-store?client_id=ruby-es-sub&channel=#{URI.encode_www_form_component(channel)}&events_store_type=1")
Net::HTTP.start(uri.host, uri.port) do |http|
  req = Net::HTTP::Get.new(uri); req["Accept"] = "text/event-stream"
  http.request(req) do |resp|
    ev_type = nil; data = nil
    resp.read_body do |chunk|
      chunk.each_line do |line|
        line.chomp!
        if line.empty?
          if ev_type == "cloudevent" && data
            ce = JSON.parse(data)
            puts "Received: type=#{ce['type']} data=#{ce['data']}"
          end
          ev_type = nil; data = nil
        elsif line.start_with?(":") then # keepalive
        elsif line.start_with?("event:") then ev_type = line.sub("event:", "").strip
        elsif line.start_with?("data:") then data = line.sub("data:", "").strip
        end
      end
    end
  end
end
use futures_util::StreamExt;

// events_store_type=1 = StartNewOnly
let sub_url = format!(
    "{}/ce/subscribe/events-store?client_id=rust-es-sub&channel={}&events_store_type=1",
    base, channel
);
let stream = client.get(&sub_url)
    .header("Accept", "text/event-stream")
    .send().await?.bytes_stream();
let mut stream = Box::pin(stream);

let mut ev_type = String::new();
let mut data = String::new();
let mut buffer = String::new();
while let Some(chunk) = stream.next().await {
    buffer.push_str(&String::from_utf8_lossy(&chunk?));
    while let Some(pos) = buffer.find('\n') {
        let line = buffer[..pos].trim_end_matches('\r').to_string();
        buffer = buffer[pos + 1..].to_string();
        if line.is_empty() {
            if ev_type == "cloudevent" && !data.is_empty() {
                let ce: serde_json::Value = serde_json::from_str(&data)?;
                println!("Received: type={} data={}", ce["type"], ce["data"]);
                return Ok(());
            }
            ev_type.clear(); data.clear();
        } else if line.starts_with(':') {
        } else if let Some(v) = line.strip_prefix("event:") { ev_type = v.trim().to_string(); }
        else if let Some(v) = line.strip_prefix("data:") { data = v.trim().to_string(); }
    }
}

Replay from the first event

Set events_store_type=2 (StartFromFirst) to replay every stored event from the beginning of the log before streaming new ones.

curl -N "http://localhost:9090/ce/subscribe/events-store?client_id=my-client&channel=audit-log&events_store_type=2"
// events_store_type=2 = StartFromFirst — replay all stored events.
var url = $"{base_}/ce/subscribe/events-store?client_id=csharp-es-replay&channel={Uri.EscapeDataString(channel)}&events_store_type=2";
using var http = new HttpClient { Timeout = Timeout.InfiniteTimeSpan };
using var req = new HttpRequestMessage(HttpMethod.Get, url);
req.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("text/event-stream"));
using var resp = await http.SendAsync(req, HttpCompletionOption.ResponseHeadersRead);
using var reader = new StreamReader(await resp.Content.ReadAsStreamAsync(), Encoding.UTF8);

string? evType = null, data = null, line;
while ((line = await reader.ReadLineAsync()) != null)
{
    if (line == "") {
        if (evType == "cloudevent" && data != null) {
            var ce = JsonSerializer.Deserialize<JsonElement>(data);
            Console.WriteLine($"  replayed data={ce.GetProperty("data")}");
        }
        evType = null; data = null;
    }
    else if (line.StartsWith("event:")) evType = line[6..].Trim();
    else if (line.StartsWith("data:")) data = line[5..].Trim();
}
// events_store_type=2 = StartFromFirst: replay all stored events.
sseURL := fmt.Sprintf(
	"%s/ce/subscribe/events-store?client_id=%s&channel=%s&events_store_type=2",
	base, clientID, channel)
req, _ := http.NewRequest("GET", sseURL, nil)
req.Header.Set("Accept", "text/event-stream")

resp, err := (&http.Client{Timeout: 0}).Do(req)
if err != nil {
	log.Fatal("SSE connect:", err)
}
defer resp.Body.Close()

scanner := bufio.NewScanner(resp.Body)
var evType, data string
for scanner.Scan() {
	line := scanner.Text()
	if line == "" {
		if evType == "cloudevent" && data != "" {
			var ce map[string]interface{}
			_ = json.Unmarshal([]byte(data), &ce)
			fmt.Printf("  replayed data=%v\n", ce["data"])
		}
		evType, data = "", ""
		continue
	}
	if strings.HasPrefix(line, "event:") {
		evType = strings.TrimSpace(strings.TrimPrefix(line, "event:"))
	} else if strings.HasPrefix(line, "data:") {
		data = strings.TrimSpace(strings.TrimPrefix(line, "data:"))
	}
}
// events_store_type=2 = StartFromFirst
String sseUrl = base + "/ce/subscribe/events-store?client_id=java-es-replay&channel="
        + channel + "&events_store_type=2";
HttpURLConnection conn = (HttpURLConnection) new URL(sseUrl).openConnection();
conn.setRequestProperty("Accept", "text/event-stream");
conn.setReadTimeout(15_000);
try (BufferedReader reader = new BufferedReader(
        new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) {
    String line; String evType = null, data = null;
    while ((line = reader.readLine()) != null) {
        if (line.isEmpty()) {
            if ("cloudevent".equals(evType) && data != null) {
                Map<?, ?> ce = MAPPER.readValue(data, Map.class);
                System.out.println("  replayed data=" + ce.get("data"));
            }
            evType = null; data = null;
        } else if (line.startsWith("event:")) evType = line.substring(6).trim();
        else if (line.startsWith("data:")) data = line.substring(5).trim();
    }
}
import EventSource from 'eventsource';

// events_store_type=2 = StartFromFirst — replay all stored events.
const url = `${base}/ce/subscribe/events-store?client_id=js-replay-sub&channel=${encodeURIComponent(channel)}&events_store_type=2`;
const es = new EventSource(url);

es.addEventListener('cloudevent', (evt: MessageEvent) => {
  const ce = JSON.parse(evt.data) as Record<string, unknown>;
  console.log(`  replayed data=${JSON.stringify(ce.data)}`);
});
es.addEventListener('error', (err) => {
  console.error('SSE error:', err);
  es.close();
});
import json
import requests

# events_store_type=2 = StartFromFirst
sse_url = (f"{base}/ce/subscribe/events-store"
           f"?client_id=python-replay-sub&channel={channel}&events_store_type=2")
with requests.get(sse_url, stream=True, timeout=None,
                  headers={"Accept": "text/event-stream"}) as resp:
    ev_type = data = ""
    for line in resp.iter_lines(decode_unicode=True):
        if line == "":
            if ev_type == "cloudevent" and data:
                ce = json.loads(data)
                print(f"  replayed data={ce.get('data')}")
            ev_type = data = ""
            continue
        if line.startswith("event:"):
            ev_type = line[6:].strip()
        elif line.startswith("data:"):
            data = line[5:].strip()
require "net/http"; require "uri"; require "json"

# events_store_type=2 = StartFromFirst
uri = URI("#{base}/ce/subscribe/events-store?client_id=ruby-es-replay&channel=#{URI.encode_www_form_component(channel)}&events_store_type=2")
Net::HTTP.start(uri.host, uri.port) do |http|
  req = Net::HTTP::Get.new(uri); req["Accept"] = "text/event-stream"
  http.request(req) do |resp|
    ev_type = nil; data = nil
    resp.read_body do |chunk|
      chunk.each_line do |line|
        line.chomp!
        if line.empty?
          if ev_type == "cloudevent" && data
            ce = JSON.parse(data)
            ce["data"] = JSON.parse(ce["data"]) if ce["data"].is_a?(String)
            puts "  replayed seq=#{ce.dig('data', 'seq')}"
          end
          ev_type = nil; data = nil
        elsif line.start_with?("event:") then ev_type = line.sub("event:", "").strip
        elsif line.start_with?("data:") then data = line.sub("data:", "").strip
        end
      end
    end
  end
end
use futures_util::StreamExt;

// events_store_type=2 = StartFromFirst
let sub_url = format!(
    "{}/ce/subscribe/events-store?client_id=rust-es-replay&channel={}&events_store_type=2",
    base, channel
);
let stream = client.get(&sub_url)
    .header("Accept", "text/event-stream")
    .send().await?.bytes_stream();
let mut stream = Box::pin(stream);

let mut ev_type = String::new();
let mut data = String::new();
let mut buffer = String::new();
while let Some(chunk) = stream.next().await {
    buffer.push_str(&String::from_utf8_lossy(&chunk?));
    while let Some(pos) = buffer.find('\n') {
        let line = buffer[..pos].trim_end_matches('\r').to_string();
        buffer = buffer[pos + 1..].to_string();
        if line.is_empty() {
            if ev_type == "cloudevent" && !data.is_empty() {
                let ce: serde_json::Value = serde_json::from_str(&data)?;
                println!("  replayed data={}", ce["data"]);
            }
            ev_type.clear(); data.clear();
        } else if let Some(v) = line.strip_prefix("event:") { ev_type = v.trim().to_string(); }
        else if let Some(v) = line.strip_prefix("data:") { data = v.trim().to_string(); }
    }
}

Replay at a sequence number

Set events_store_type=4 (StartAtSequence) and events_store_value=<sequence> to start the stream at a specific sequence number. The example below publishes 10 events and then replays from sequence 5.

curl -N "http://localhost:9090/ce/subscribe/events-store?client_id=my-client&channel=audit-log&events_store_type=4&events_store_value=5"
// events_store_type=4 (StartAtSequence) + events_store_value=startAtSeq
long startAtSeq = 5;
var url = $"{base_}/ce/subscribe/events-store?client_id=csharp-es-seq&channel={Uri.EscapeDataString(channel)}&events_store_type=4&events_store_value={startAtSeq}";
using var http = new HttpClient { Timeout = Timeout.InfiniteTimeSpan };
using var req = new HttpRequestMessage(HttpMethod.Get, url);
req.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("text/event-stream"));
using var resp = await http.SendAsync(req, HttpCompletionOption.ResponseHeadersRead);
using var reader = new StreamReader(await resp.Content.ReadAsStreamAsync(), Encoding.UTF8);

string? evType = null, data = null, line;
while ((line = await reader.ReadLineAsync()) != null)
{
    if (line == "") {
        if (evType == "cloudevent" && data != null) {
            var ce = JsonSerializer.Deserialize<JsonElement>(data);
            Console.WriteLine($"  Received: seq={ce.GetProperty("data").GetProperty("seq")}");
        }
        evType = null; data = null;
    }
    else if (line.StartsWith("event:")) evType = line[6..].Trim();
    else if (line.StartsWith("data:")) data = line[5..].Trim();
}
// events_store_type=4 (StartAtSequence), events_store_value=5
const startSeq = 5
sseURL := fmt.Sprintf(
	"%s/ce/subscribe/events-store?client_id=go-replay-seq-sub&channel=%s&events_store_type=4&events_store_value=%d",
	base, channel, startSeq)
req, _ := http.NewRequest("GET", sseURL, nil)
req.Header.Set("Accept", "text/event-stream")

resp, err := (&http.Client{Timeout: 0}).Do(req)
if err != nil {
	log.Fatal("SSE connect:", err)
}
defer resp.Body.Close()

scanner := bufio.NewScanner(resp.Body)
var evType, data, sseID string
for scanner.Scan() {
	line := scanner.Text()
	if line == "" {
		if evType == "cloudevent" && data != "" {
			var ce map[string]interface{}
			_ = json.Unmarshal([]byte(data), &ce)
			fmt.Printf("  [seq=%s] data=%v\n", sseID, ce["data"])
		}
		evType, data, sseID = "", "", ""
		continue
	}
	if strings.HasPrefix(line, "id:") {
		sseID = strings.TrimSpace(strings.TrimPrefix(line, "id:"))
	} else if strings.HasPrefix(line, "event:") {
		evType = strings.TrimSpace(strings.TrimPrefix(line, "event:"))
	} else if strings.HasPrefix(line, "data:") {
		data = strings.TrimSpace(strings.TrimPrefix(line, "data:"))
	}
}
// events_store_type=4 = StartAtSequence, events_store_value=startAtSeq
long startAtSeq = 5;
String sseUrl = base + "/ce/subscribe/events-store?client_id=java-es-seq&channel="
        + channel + "&events_store_type=4&events_store_value=" + startAtSeq;
HttpURLConnection conn = (HttpURLConnection) new URL(sseUrl).openConnection();
conn.setRequestProperty("Accept", "text/event-stream");
conn.setReadTimeout(15_000);
try (BufferedReader reader = new BufferedReader(
        new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) {
    String line; String evType = null, data = null;
    while ((line = reader.readLine()) != null) {
        if (line.isEmpty()) {
            if ("cloudevent".equals(evType) && data != null) {
                Map<?, ?> ce = MAPPER.readValue(data, Map.class);
                System.out.println("  Received: seq=" + ((Map<?, ?>) ce.get("data")).get("seq"));
            }
            evType = null; data = null;
        } else if (line.startsWith("event:")) evType = line.substring(6).trim();
        else if (line.startsWith("data:")) data = line.substring(5).trim();
    }
}
import EventSource from 'eventsource';

// events_store_type=4 (StartAtSequence) + events_store_value
const startSeq = 5;
const url = `${base}/ce/subscribe/events-store?client_id=js-seq-sub&channel=${encodeURIComponent(channel)}&events_store_type=4&events_store_value=${startSeq}`;
const es = new EventSource(url);

es.addEventListener('cloudevent', (evt: MessageEvent & { lastEventId: string }) => {
  const ce = JSON.parse(evt.data) as Record<string, unknown>;
  console.log(`  [seq=${evt.lastEventId}] data=${JSON.stringify(ce.data)}`);
});
es.addEventListener('error', (err) => {
  console.error('SSE error:', err);
  es.close();
});
import json
import requests

# events_store_type=4 (StartAtSequence) + events_store_value
start_seq = 5
sse_url = (f"{base}/ce/subscribe/events-store"
           f"?client_id=python-seq-sub&channel={channel}"
           f"&events_store_type=4&events_store_value={start_seq}")
with requests.get(sse_url, stream=True, timeout=None,
                  headers={"Accept": "text/event-stream"}) as resp:
    ev_type = data = sse_id = ""
    for line in resp.iter_lines(decode_unicode=True):
        if line == "":
            if ev_type == "cloudevent" and data:
                ce = json.loads(data)
                print(f"  [seq={sse_id}] data={ce.get('data')}")
            ev_type = data = sse_id = ""
            continue
        if line.startswith("id:"):
            sse_id = line[3:].strip()
        elif line.startswith("event:"):
            ev_type = line[6:].strip()
        elif line.startswith("data:"):
            data = line[5:].strip()
require "net/http"; require "uri"; require "json"

# events_store_type=4 (StartAtSequence), events_store_value=start_at_seq
start_at_seq = 5
uri = URI("#{base}/ce/subscribe/events-store?client_id=ruby-es-seq&channel=#{URI.encode_www_form_component(channel)}&events_store_type=4&events_store_value=#{start_at_seq}")
Net::HTTP.start(uri.host, uri.port) do |http|
  req = Net::HTTP::Get.new(uri); req["Accept"] = "text/event-stream"
  http.request(req) do |resp|
    ev_type = nil; data = nil
    resp.read_body do |chunk|
      chunk.each_line do |line|
        line.chomp!
        if line.empty?
          if ev_type == "cloudevent" && data
            ce = JSON.parse(data)
            ce["data"] = JSON.parse(ce["data"]) if ce["data"].is_a?(String)
            puts "  Received: seq=#{ce.dig('data', 'seq')}"
          end
          ev_type = nil; data = nil
        elsif line.start_with?("event:") then ev_type = line.sub("event:", "").strip
        elsif line.start_with?("data:") then data = line.sub("data:", "").strip
        end
      end
    end
  end
end
use futures_util::StreamExt;

// events_store_type=4 (StartAtSequence) + events_store_value
let start_at_seq: u32 = 5;
let sub_url = format!(
    "{}/ce/subscribe/events-store?client_id=rust-es-seq&channel={}&events_store_type=4&events_store_value={}",
    base, channel, start_at_seq
);
let stream = client.get(&sub_url)
    .header("Accept", "text/event-stream")
    .send().await?.bytes_stream();
let mut stream = Box::pin(stream);

let mut ev_type = String::new();
let mut data = String::new();
let mut buffer = String::new();
while let Some(chunk) = stream.next().await {
    buffer.push_str(&String::from_utf8_lossy(&chunk?));
    while let Some(pos) = buffer.find('\n') {
        let line = buffer[..pos].trim_end_matches('\r').to_string();
        buffer = buffer[pos + 1..].to_string();
        if line.is_empty() {
            if ev_type == "cloudevent" && !data.is_empty() {
                let ce: serde_json::Value = serde_json::from_str(&data)?;
                println!("  Received: seq={}", ce["data"]["seq"]);
            }
            ev_type.clear(); data.clear();
        } else if let Some(v) = line.strip_prefix("event:") { ev_type = v.trim().to_string(); }
        else if let Some(v) = line.strip_prefix("data:") { data = v.trim().to_string(); }
    }
}

Reconnect and resume with Last-Event-ID

To resume after a dropped connection, capture the id: field of each frame as you process it, then reconnect with that value in the Last-Event-ID header. Omit events_store_type on the reconnect so the header takes effect — the server resumes from sequence + 1.

# Resume after the last sequence you processed (e.g. 42).
# Omit events_store_type so Last-Event-ID takes precedence.
curl -N -H "Last-Event-ID: 42" \
  "http://localhost:9090/ce/subscribe/events-store?client_id=my-client&channel=audit-log"
// Omit events_store_type when reconnecting with Last-Event-ID.
async Task<(List<JsonElement> events, string lastId)> Subscribe(
    string clientId, string? lastEventId, int maxEvents)
{
    var events = new List<JsonElement>();
    var lastId = "";
    using var http = new HttpClient { Timeout = Timeout.InfiniteTimeSpan };
    var url = lastEventId == null
        ? $"{base_}/ce/subscribe/events-store?client_id={clientId}&channel={Uri.EscapeDataString(channel)}&events_store_type=2"
        : $"{base_}/ce/subscribe/events-store?client_id={clientId}&channel={Uri.EscapeDataString(channel)}";
    using var req = new HttpRequestMessage(HttpMethod.Get, url);
    req.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("text/event-stream"));
    if (lastEventId != null) req.Headers.Add("Last-Event-ID", lastEventId);
    using var resp = await http.SendAsync(req, HttpCompletionOption.ResponseHeadersRead);
    using var reader = new StreamReader(await resp.Content.ReadAsStreamAsync(), Encoding.UTF8);
    string? evType = null, data = null, id = null, line;
    while ((line = await reader.ReadLineAsync()) != null)
    {
        if (line == "") {
            if (evType == "cloudevent" && data != null) {
                if (id != null) lastId = id;
                events.Add(JsonSerializer.Deserialize<JsonElement>(data));
                if (events.Count >= maxEvents) break;
            }
            evType = null; data = null; id = null;
        }
        else if (line.StartsWith("id:")) id = line[3..].Trim();
        else if (line.StartsWith("event:")) evType = line[6..].Trim();
        else if (line.StartsWith("data:")) data = line[5..].Trim();
    }
    return (events, lastId);
}

var (first, lastId) = await Subscribe("csharp-es-reconnect-1", null, 2);
var (second, _) = await Subscribe("csharp-es-reconnect-2", lastId, 2);
func openSSE(base, channel, clientID, lastEventID string) (*http.Response, error) {
	// When reconnecting, omit events_store_type so Last-Event-ID takes precedence.
	var sseURL string
	if lastEventID == "" {
		sseURL = fmt.Sprintf(
			"%s/ce/subscribe/events-store?client_id=%s&channel=%s&events_store_type=2",
			base, clientID, channel)
	} else {
		sseURL = fmt.Sprintf(
			"%s/ce/subscribe/events-store?client_id=%s&channel=%s",
			base, clientID, channel)
	}

	req, _ := http.NewRequest("GET", sseURL, nil)
	req.Header.Set("Accept", "text/event-stream")
	if lastEventID != "" {
		req.Header.Set("Last-Event-ID", lastEventID)
	}
	return (&http.Client{Timeout: 0}).Do(req)
}

// 1) First connection: read a batch and capture the last SSE id.
// 2) Reconnect: openSSE(base, channel, clientID, lastID) resumes from lastID+1.
// Omit events_store_type when reconnecting with Last-Event-ID.
String sseUrl = base + "/ce/subscribe/events-store?client_id=" + clientId + "&channel=" + channel;
if (lastEventId == null) {
    sseUrl += "&events_store_type=2"; // initial connection: start from first
}
HttpURLConnection conn = (HttpURLConnection) new URL(sseUrl).openConnection();
conn.setRequestProperty("Accept", "text/event-stream");
conn.setReadTimeout(10_000);
if (lastEventId != null) {
    conn.setRequestProperty("Last-Event-ID", lastEventId);
}
try (BufferedReader reader = new BufferedReader(
        new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) {
    String line; String evType = null, data = null, id = null, lastId = "";
    while ((line = reader.readLine()) != null) {
        if (line.isEmpty()) {
            if ("cloudevent".equals(evType) && data != null) {
                if (id != null) lastId = id; // capture sequence for resume
                Map<?, ?> ce = MAPPER.readValue(data, Map.class);
                System.out.println("  seq=" + ((Map<?, ?>) ce.get("data")).get("seq"));
            }
            evType = null; data = null; id = null;
        } else if (line.startsWith("id:")) id = line.substring(3).trim();
        else if (line.startsWith("event:")) evType = line.substring(6).trim();
        else if (line.startsWith("data:")) data = line.substring(5).trim();
    }
}
// EventSource sends Last-Event-ID automatically on reconnect, but custom-header
// reconnects use fetch. Capture evt.lastEventId, then resume with the header.
const headers: Record<string, string> = { 'Last-Event-ID': lastEventId };
const url = `${base}/ce/subscribe/events-store?client_id=js-reconnect-sub&channel=${encodeURIComponent(channel)}`;
const resp = await fetch(url, { headers });
const reader = resp.body!.getReader();
const decoder = new TextDecoder();
let buffer = '', evType = '', data = '', lastId = lastEventId;

for (;;) {
  const { value, done } = await reader.read();
  if (done) break;
  buffer += decoder.decode(value, { stream: true });
  const lines = buffer.split('\n');
  buffer = lines.pop() ?? '';
  for (const line of lines) {
    if (line === '') {
      if (evType === 'cloudevent' && data) {
        const ce = JSON.parse(data) as Record<string, unknown>;
        console.log(`  id=${lastId} data=${JSON.stringify(ce.data)}`);
      }
      evType = data = '';
    } else if (line.startsWith('id:')) lastId = line.slice(3).trim();
    else if (line.startsWith('event:')) evType = line.slice(6).trim();
    else if (line.startsWith('data:')) data = line.slice(5).trim();
  }
}
import json
import requests

def read_n_events(base, channel, client_id, n, last_event_id=""):
    """Open SSE, read n events, return the last SSE id seen."""
    if last_event_id:
        # Omit events_store_type so Last-Event-ID takes precedence.
        sse_url = f"{base}/ce/subscribe/events-store?client_id={client_id}&channel={channel}"
        extra_headers = {"Last-Event-ID": last_event_id}
    else:
        sse_url = (f"{base}/ce/subscribe/events-store"
                   f"?client_id={client_id}&channel={channel}&events_store_type=2")
        extra_headers = {}

    headers = {"Accept": "text/event-stream", **extra_headers}
    last_id = ""
    count = 0
    with requests.get(sse_url, stream=True, timeout=None, headers=headers) as resp:
        ev_type = data = sse_id = ""
        for line in resp.iter_lines(decode_unicode=True):
            if line == "":
                if ev_type == "cloudevent" and data:
                    ce = json.loads(data)
                    count += 1
                    last_id = sse_id
                    print(f"  [{count}] id={sse_id} data={ce.get('data')}")
                    if count == n:
                        return last_id
                ev_type = data = sse_id = ""
                continue
            if line.startswith("id:"):
                sse_id = line[3:].strip()
            elif line.startswith("event:"):
                ev_type = line[6:].strip()
            elif line.startswith("data:"):
                data = line[5:].strip()
    return last_id

# 1) First connection captures the last id; 2) reconnect resumes from it.
last_id = read_n_events(base, channel, "python-reconnect-sub", 3)
read_n_events(base, channel, "python-reconnect-sub", 3, last_id)
require "net/http"; require "uri"; require "json"

# Open one SSE connection, collect up to max_events, return [events, last_id].
def subscribe_es(base, channel, last_event_id, max_events)
  # Omit events_store_type when reconnecting with Last-Event-ID.
  query = last_event_id.nil? ? "&events_store_type=2" : ""
  uri = URI("#{base}/ce/subscribe/events-store?client_id=ruby-es-reconnect&channel=#{URI.encode_www_form_component(channel)}#{query}")
  events = []; last_id = nil
  Net::HTTP.start(uri.host, uri.port, read_timeout: 12) do |http|
    req = Net::HTTP::Get.new(uri); req["Accept"] = "text/event-stream"
    req["Last-Event-ID"] = last_event_id if last_event_id
    http.request(req) do |resp|
      ev_type = nil; data = nil; id = nil
      resp.read_body do |chunk|
        chunk.each_line do |line|
          line.chomp!
          if line.empty?
            if ev_type == "cloudevent" && data
              last_id = id if id
              ce = JSON.parse(data)
              ce["data"] = JSON.parse(ce["data"]) if ce["data"].is_a?(String)
              events << ce
              return [events, last_id] if events.size >= max_events
            end
            ev_type = nil; data = nil; id = nil
          elsif line.start_with?("id:") then id = line.sub("id:", "").strip
          elsif line.start_with?("event:") then ev_type = line.sub("event:", "").strip
          elsif line.start_with?("data:") then data = line.sub("data:", "").strip
          end
        end
      end
    end
  end
  [events, last_id]
end

first_events, last_id = subscribe_es(base, channel, nil, 2)
second_events, _ = subscribe_es(base, channel, last_id, 2)
use futures_util::StreamExt;
use reqwest::Client;
use serde_json::Value;

// Subscribe and collect up to `max` events. Returns (events, last_id_seen).
async fn subscribe_and_collect(
    client: &Client,
    url: &str,
    last_event_id: Option<&str>,
    max: usize,
) -> (Vec<Value>, String) {
    let mut req = client.get(url).header("Accept", "text/event-stream");
    if let Some(id) = last_event_id {
        req = req.header("Last-Event-ID", id);
    }
    let stream = req.send().await.expect("SSE connect").bytes_stream();
    let mut stream = Box::pin(stream);
    let mut ev_type = String::new(); let mut data = String::new();
    let mut id_field = String::new(); let mut last_id = String::new();
    let mut buffer = String::new();
    let mut events = Vec::new();

    while let Some(chunk) = stream.next().await {
        buffer.push_str(&String::from_utf8_lossy(&chunk.unwrap_or_default()));
        while let Some(pos) = buffer.find('\n') {
            let line = buffer[..pos].trim_end_matches('\r').to_string();
            buffer = buffer[pos + 1..].to_string();
            if line.is_empty() {
                if ev_type == "cloudevent" && !data.is_empty() {
                    if !id_field.is_empty() { last_id = id_field.clone(); }
                    events.push(serde_json::from_str(&data).unwrap_or(Value::Null));
                    if events.len() >= max { return (events, last_id); }
                }
                ev_type.clear(); data.clear(); id_field.clear();
            } else if let Some(v) = line.strip_prefix("id:") { id_field = v.trim().to_string(); }
            else if let Some(v) = line.strip_prefix("event:") { ev_type = v.trim().to_string(); }
            else if let Some(v) = line.strip_prefix("data:") { data = v.trim().to_string(); }
        }
    }
    (events, last_id)
}

// First connection (events_store_type=2) captures last_id; the reconnect URL
// omits events_store_type and passes Some(&last_id) to resume from sequence + 1.

Response and status codes

StatusWhen
202 AcceptedPOST /ce/send/event-store succeeded; the send result is in data
200 OK (stream)GET /ce/subscribe/events-store opened; frames follow as text/event-stream
400 Bad RequestInvalid CloudEvent, missing channel, or a reserved channel name
429 Too Many RequestsSSE connection limit reached (MaxSSEConnections exceeded)
500 Internal Server ErrorBackend messaging error or SSE setup failure

A successful send returns the standard envelope:

{
  "is_error": false,
  "message": "OK",
  "data": { ... }
}

Each delivered CloudEvent frame carries the sequence number in its id: field:

id: 42
event: cloudevent
data: {"specversion":"1.0","type":"com.example.audit.entry","source":"audit-service","id":"...","subject":"audit-log","time":"...","data":{"action":"user.login"}}

Was this page helpful?

On this page