Events Reference
Complete reference for KubeMQ Events message structure, validation rules, configuration, and error codes.
Message Structure
Event Message (Send)
Prop
Type
Event Receive (Subscribe)
Prop
Type
Unlike Events Store, plain Events do not include Timestamp or Sequence fields because there is no persistence layer.
Send Result
Prop
Type
Subscription Options
Prop
Type
Consumer Groups
When multiple subscribers specify the same group value on the same channel:
- Each event is delivered to exactly one member of the group (round-robin)
- When
groupis empty, every subscriber receives every event (standard fan-out) - Groups are independent per channel
- There is no limit on the number of group members
See Consumer Groups and Scale Subscribers for examples.
Channel Naming Rules
| Rule | Constraint | Error Code |
|---|---|---|
| Required | Cannot be empty | 102 |
| No trailing dot | Cannot end with . | 119 |
| No whitespace | Cannot contain spaces | 108 |
| No wildcards (publish) | Cannot contain * or > when publishing | 107 |
| Max length | 256 characters | — |
Valid channel regex (publish): ^[^\s*>]+[^.]$
Wildcards (* and >) are allowed in subscription channel patterns only. See Wildcard Subscriptions.
Channel Naming Conventions
Use dot-separated hierarchical names for best results with wildcard subscriptions:
{domain}.{entity}.{action}
orders.created
orders.us-east.shipped
payments.completed
inventory.reservedRouting (Multicast)
Events support multicast publishing through special channel syntax:
| Character | Purpose | Example |
|---|---|---|
; | Separate channels of the same type | orders;notifications |
: | Specify target pattern type | events:orders;events_store:audit |
Pattern Type Prefixes
| Prefix | Target Pattern |
|---|---|
events: | Events (fire-and-forget) |
events_store: | Events Store (persistent) |
queues: | Queues (guaranteed delivery) |
Routed messages are automatically tagged with X-KUBEMQ-ROUTED=true. See Multicast Events.
Transport Protocols
gRPC
- Publish:
SendEvent(pb.Event)— unary RPC - Publish Stream:
SendEventsStream()— bidirectional streaming for high-throughput publishing - Subscribe:
SubscribeToEvents(pb.Subscribe)— server-streaming RPC - Default port:
50000 - Max message size: ~1 GB (configurable)
REST
- Publish:
POST /sendwith JSON body - Subscribe: WebSocket upgrade for streaming delivery
- Default port:
9090 - Max body size: 100 MB (configurable)
WebSocket
- Publish and subscribe over persistent WebSocket connections
- JSON-encoded messages
- Read limit: 1 MB
Configuration
Server Configuration
| Setting | Default | Description |
|---|---|---|
Connectors.Grpc.BodyLimit | ~1 GB | Max gRPC message size |
Connectors.Rest.BodyLimit | 100 MB | Max REST body size |
Client Configuration
Prop
Type
Slow Consumer Handling
When a subscriber's receive buffer is full, KubeMQ waits up to the write deadline (default 2 seconds) for the buffer to clear. If the deadline expires:
- The event is dropped for that subscriber
- A warning is logged server-side with the channel, event ID, and metadata
- Other subscribers are not affected
See Handle Slow Consumers for mitigation strategies.
Delivery Semantics
| Aspect | Behavior |
|---|---|
| Delivery guarantee | At-most-once |
| Persistence | None — events flow through memory only |
| Ordering | Events are delivered in publish order per channel |
| Duplicates | No duplicates (single delivery attempt) |
| Acknowledgment | None — fire-and-forget |
| Retry | None — failed deliveries are not retried |
Events vs Events Store
| Feature | Events | Events Store |
|---|---|---|
| Persistence | No | Yes (disk-backed) |
| Replay | No | Yes (from offset, time, sequence) |
| Delivery guarantee | At-most-once | At-least-once |
| Wildcard subscriptions | Yes (*, >) | No |
| Consumer groups | Yes (round-robin) | Yes (durable) |
| Latency | Lowest | Slightly higher (disk write) |
| Use cases | Real-time notifications, streaming | Audit trails, event sourcing |
Error Codes
| Code | Error | Description |
|---|---|---|
| 101 | Invalid ClientID | ClientID is empty |
| 102 | Invalid Channel | Channel is empty |
| 107 | Invalid Wildcards | Channel contains * or > (publish only) |
| 108 | Invalid Whitespace | Channel contains spaces |
| 110 | Invalid Message | Both Body and Metadata are empty |
| 119 | Invalid Channel Separator | Channel ends with . |
Runtime Errors
| Error | Cause | Resolution |
|---|---|---|
ErrShutdownMode | Server is shutting down | Reconnect after server restart |
ErrConnectionNoAvailable | broker connection is down | SDK auto-reconnects; check server health |
| Authorization denied | Casbin policy rejected the operation | Verify client permissions |
SDK Quick Reference
// Publish
client.SendEvent(ctx, kubemq.NewEvent().
SetChannel("ch").SetBody([]byte("data")))
// Subscribe
client.SubscribeToEvents(ctx, "ch", "group",
kubemq.WithOnEvent(handler),
kubemq.WithOnError(errHandler))
// Stream Publish
streamCh := make(chan *kubemq.Event, 100)
resultCh := make(chan *kubemq.EventSendResult, 100)
go client.StreamEvents(ctx, streamCh, resultCh)# Publish
client.send_event(EventMessage(channel="ch", body=b"data"))
# Subscribe
client.subscribe_to_events(
EventsSubscription(channel="ch", group="group",
on_receive_event_callback=handler,
on_error_callback=err_handler),
cancel=CancellationToken())
# Stream Publish
stream = client.open_events_stream()
stream.send(EventMessage(channel="ch", body=b"data"))// Publish
await client.sendEvent({ channel: "ch", body: Buffer.from("data") });
// Subscribe
client.subscribeToEvents({
channel: "ch", group: "group",
onEvent: handler, onError: errHandler
});
// Stream Publish
const stream = client.createEventStream();
stream.send({ channel: "ch", body: Buffer.from("data") });// Publish
client.sendEventsMessage(EventMessage.builder()
.channel("ch").body("data".getBytes()).build());
// Subscribe
client.subscribeToEvents(EventsSubscription.builder()
.channel("ch").group("group")
.onReceiveEventCallback(handler)
.onErrorCallback(errHandler).build());
// Stream Publish
EventsStream stream = client.openEventsStream();
stream.send(EventMessage.builder()
.channel("ch").body("data".getBytes()).build());// Publish
await client.SendEventAsync(new EventMessage {
Channel = "ch", Body = Encoding.UTF8.GetBytes("data") });
// Subscribe
await foreach (var msg in client.SubscribeToEventsAsync(
new EventsSubscription { Channel = "ch", Group = "group" })) { }
// Stream Publish
var stream = client.OpenEventsStream();
await stream.SendAsync(new EventMessage {
Channel = "ch", Body = Encoding.UTF8.GetBytes("data") });// Publish
client.sendEvent(EventMessage(
channel = "ch", body = "data".toByteArray()))
// Subscribe
client.subscribeToEvents(
channel = "ch", group = "group",
onEvent = handler, onError = errHandler)
// Stream Publish
val stream = client.openEventsStream()
stream.send(EventMessage(channel = "ch", body = "data".toByteArray()))// Publish
kubemq::EventMessage event;
event.channel = "ch";
event.body = "data";
client.sendEvent(event);
// Subscribe
client.subscribeToEvents("ch", "group", handler, errHandler);
// Stream Publish
auto stream = client.openEventsStream();
stream.send(event);// Publish
let event = EventBuilder::new()
.channel("ch").body(b"data".to_vec()).build();
client.send_event(event).await?;
// Subscribe
let sub = client.subscribe_to_events("ch", "group",
|event| Box::pin(async move { handle(event).await }),
None).await?;
// Stream Publish
let mut stream = client.send_event_stream().await?;
stream.send(EventBuilder::new()
.channel("ch").body(b"data".to_vec()).build()).await?;# Publish
client.send_event(KubeMQ::PubSub::EventMessage.new(
channel: "ch", body: "data"))
# Subscribe
sub = KubeMQ::PubSub::EventsSubscription.new(channel: "ch", group: "group")
client.subscribe_to_events(sub, cancellation_token: cancel,
on_error: ->(e) { handle_error(e) }) { |event| handle(event) }
# Stream Publish
sender = client.create_events_sender
sender.publish(KubeMQ::PubSub::EventMessage.new(
channel: "ch", body: "data"))# Publish
event = KubeMQ.Event.new(channel: "ch", body: "data")
KubeMQ.Client.send_event(client, event)
# Subscribe
{:ok, sub} = KubeMQ.Client.subscribe_to_events(client, "ch",
group: "group", on_event: fn event -> handle(event) end)
# Stream Publish
{:ok, handle} = KubeMQ.Client.send_event_stream(client)
KubeMQ.EventStreamHandle.send(handle,
KubeMQ.Event.new(channel: "ch", body: "data"))Related
- Getting Started with Events
- Events Store Reference for the persistent variant
- Queues Reference for guaranteed delivery
Was this page helpful?