KubeMQ
LearnEvents

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 group is 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

RuleConstraintError Code
RequiredCannot be empty102
No trailing dotCannot end with .119
No whitespaceCannot contain spaces108
No wildcards (publish)Cannot contain * or > when publishing107
Max length256 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.reserved

Routing (Multicast)

Events support multicast publishing through special channel syntax:

CharacterPurposeExample
;Separate channels of the same typeorders;notifications
:Specify target pattern typeevents:orders;events_store:audit

Pattern Type Prefixes

PrefixTarget 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 /send with 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

SettingDefaultDescription
Connectors.Grpc.BodyLimit~1 GBMax gRPC message size
Connectors.Rest.BodyLimit100 MBMax 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

AspectBehavior
Delivery guaranteeAt-most-once
PersistenceNone — events flow through memory only
OrderingEvents are delivered in publish order per channel
DuplicatesNo duplicates (single delivery attempt)
AcknowledgmentNone — fire-and-forget
RetryNone — failed deliveries are not retried

Events vs Events Store

FeatureEventsEvents Store
PersistenceNoYes (disk-backed)
ReplayNoYes (from offset, time, sequence)
Delivery guaranteeAt-most-onceAt-least-once
Wildcard subscriptionsYes (*, >)No
Consumer groupsYes (round-robin)Yes (durable)
LatencyLowestSlightly higher (disk write)
Use casesReal-time notifications, streamingAudit trails, event sourcing

Error Codes

CodeErrorDescription
101Invalid ClientIDClientID is empty
102Invalid ChannelChannel is empty
107Invalid WildcardsChannel contains * or > (publish only)
108Invalid WhitespaceChannel contains spaces
110Invalid MessageBoth Body and Metadata are empty
119Invalid Channel SeparatorChannel ends with .

Runtime Errors

ErrorCauseResolution
ErrShutdownModeServer is shutting downReconnect after server restart
ErrConnectionNoAvailablebroker connection is downSDK auto-reconnects; check server health
Authorization deniedCasbin policy rejected the operationVerify 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"))

Was this page helpful?

On this page