KubeMQ
LearnRPC

Commands & Queries Reference

Complete reference for KubeMQ RPC — request/response structure, caching, timeouts, and error codes.

Request Model

Commands and queries share the same request structure. The RequestTypeData field determines whether the request is treated as a command or query.

Prop

Type

Response Model

Prop

Type

Command vs Query Response Differences

Response FieldCommandQuery
BodyAlways nil (stripped)Preserved
MetadataAlways "" (stripped)Preserved
CacheHitAlways false (stripped)Preserved
ReplyChannelAlways "" (stripped)Always "" (stripped)
ExecutedPreservedPreserved
ErrorPreservedPreserved

Validation Rules

Channel Name

RuleConstraintError Code
RequiredCannot be empty102
No trailing dotCannot end with .119
No whitespaceCannot contain spaces108
No wildcardsCannot contain * or >107

Valid channel name regex: ^[^\s*>]+[^.]$

Client ID

RuleConstraintError Code
RequiredCannot be empty101
AlphanumericMust match ^[a-zA-Z0-9_-]+$

Request Content

At least one of Body or Metadata must be provided. If both are empty, the request is rejected with error code 115.

Timeout

RuleConstraintError Code
RequiredMust be greater than zero109

Cache (Queries Only)

RuleConstraintError Code
CacheTTL requiredIf CacheKey is set, CacheTTL must be > 0116

Subscription Model

Prop

Type

Consumer Groups

When multiple responders specify the same group value on the same channel:

  • Each request is delivered to exactly one member of the group (round-robin)
  • When group is empty, every responder receives every request (fan-out)
  • Groups are independent per channel
  • There is no limit on the number of group members

See Load Balancing for examples.

Caching Configuration

Transport Protocols

gRPC

MethodTypeDescription
SendRequest(Request) → ResponseUnarySend command or query
SendResponse(Response) → EmptyUnarySend response back to requester
SubscribeToRequests(Subscribe) → stream RequestServer streamSubscribe to commands or queries

Default port: 50000

REST

MethodPathDescription
POST/send/requestSend command or query
POST/send/responseSend RPC response
GET/subscribe/requestsWebSocket: subscribe to commands or queries

Default port: 9090

REST subscription query parameters:

ParameterDescriptionExample
client_idClient identifiermy-responder
channelChannel nameorders.process
groupLoad balancing groupworkers
subscribe_typecommands or queriescommands

Internal Channel Mapping

PatternChannel PrefixExample
Commands_COMMANDS_._COMMANDS_.orders.process
Queries_QUERIES_._QUERIES_.inventory.lookup

Middleware Chain

Command Sender

Request → Logging → Monitor → Metrics → broker request

Query Sender

Request → Logging → Monitor → Cache → Metrics → broker request

Receiver (Commands and Queries)

broker queue-subscribe → Logging → reqCh delivery

Error Codes

Input Validation

CodeErrorDescription
101Invalid ClientIDClientID is empty
102Invalid ChannelChannel is empty
107Invalid ChannelChannel contains wildcards (* or >)
108Invalid ChannelChannel contains whitespace
109Invalid TimeoutTimeout is zero or negative
115Invalid RequestBoth body and metadata are empty
116Invalid CacheTTLCacheKey is set but CacheTTL is zero or negative
117Invalid RequestIDRequestID is empty (on response)
119Invalid ChannelChannel ends with .

Runtime Errors

CodeErrorDescription
206Invalid Request TypeRequest type is not Command or Query
207Invalid Subscribe TypeSubscribe type is not Commands or Queries
301Request TimeoutNo reply received before timeout expired
302Connection Unavailablebroker connection is down
303Invalid Response FormatReply data cannot be unmarshaled
409Shutdown ModeServer is shutting down, all operations rejected
412Access DeniedAuthorization denied for the resource

Delivery Semantics

AspectBehavior
Delivery guaranteeAt-most-once (request is sent once, no automatic retry)
OrderingRequests are independent (no ordering guarantee between requests)
TimeoutSender blocks until response arrives or timeout expires
AcknowledgmentImplicit — response is the acknowledgment

SDK Quick Reference

// Send Command
client.SendCommand(ctx, kubemq.NewCommand().
    SetChannel("ch").SetBody([]byte("data")).SetTimeout(10*time.Second))

// Send Query
client.SendQuery(ctx, kubemq.NewQuery().
    SetChannel("ch").SetBody([]byte("data")).SetTimeout(10*time.Second))

// Send Query with Cache
client.SendQuery(ctx, kubemq.NewQuery().
    SetChannel("ch").SetBody([]byte("data")).SetTimeout(10*time.Second).
    SetCacheKey("key").SetCacheTTL(60*time.Second))

// Subscribe to Commands
client.SubscribeToCommands(ctx, "ch", "group",
    kubemq.WithOnCommandReceive(handler),
    kubemq.WithOnError(errHandler))

// Subscribe to Queries
client.SubscribeToQueries(ctx, "ch", "group",
    kubemq.WithOnQueryReceive(handler),
    kubemq.WithOnError(errHandler))
# Send Command
client.send_command(CommandMessage(
    channel="ch", body=b"data", timeout_in_seconds=10))

# Send Query
client.send_query(QueryMessage(
    channel="ch", body=b"data", timeout_in_seconds=10))

# Send Query with Cache
client.send_query(QueryMessage(
    channel="ch", body=b"data", timeout_in_seconds=10,
    cache_key="key", cache_ttl_in_seconds=60))

# Subscribe to Commands
client.subscribe_to_commands(CommandsSubscription(
    channel="ch", group="group",
    on_receive_command_callback=handler,
    on_error_callback=err_handler), cancel=cancel)

# Subscribe to Queries
client.subscribe_to_queries(QueriesSubscription(
    channel="ch", group="group",
    on_receive_query_callback=handler,
    on_error_callback=err_handler), cancel=cancel)
// Send Command
await client.sendCommand({
  channel: "ch", body: Buffer.from("data"), timeoutInSeconds: 10 });

// Send Query
await client.sendQuery({
  channel: "ch", body: Buffer.from("data"), timeoutInSeconds: 10 });

// Send Query with Cache
await client.sendQuery({
  channel: "ch", body: Buffer.from("data"), timeoutInSeconds: 10,
  cacheKey: "key", cacheTTL: 60000 });

// Subscribe to Commands
client.subscribeToCommands({
  channel: "ch", group: "group",
  onCommand: handler, onError: errHandler });

// Subscribe to Queries
client.subscribeToQueries({
  channel: "ch", group: "group",
  onQuery: handler, onError: errHandler });
// Send Command
client.sendCommandRequest(CommandMessage.builder()
    .channel("ch").body("data".getBytes()).timeout(10000).build());

// Send Query
client.sendQueryRequest(QueryMessage.builder()
    .channel("ch").body("data".getBytes()).timeout(10000).build());

// Send Query with Cache
client.sendQueryRequest(QueryMessage.builder()
    .channel("ch").body("data".getBytes()).timeout(10000)
    .cacheKey("key").cacheTTL(60000).build());

// Subscribe to Commands
client.subscribeToCommands(CommandsSubscription.builder()
    .channel("ch").group("group")
    .onReceiveCommandCallback(handler)
    .onErrorCallback(errHandler).build());

// Subscribe to Queries
client.subscribeToQueries(QueriesSubscription.builder()
    .channel("ch").group("group")
    .onReceiveQueryCallback(handler)
    .onErrorCallback(errHandler).build());
// Send Command
await client.SendCommandAsync(new CommandMessage {
    Channel = "ch", Body = Encoding.UTF8.GetBytes("data"),
    Timeout = TimeSpan.FromSeconds(10) });

// Send Query
await client.SendQueryAsync(new QueryMessage {
    Channel = "ch", Body = Encoding.UTF8.GetBytes("data"),
    Timeout = TimeSpan.FromSeconds(10) });

// Send Query with Cache
await client.SendQueryAsync(new QueryMessage {
    Channel = "ch", Body = Encoding.UTF8.GetBytes("data"),
    Timeout = TimeSpan.FromSeconds(10),
    CacheKey = "key", CacheTTL = TimeSpan.FromSeconds(60) });

// Subscribe to Commands
await foreach (var cmd in client.SubscribeToCommandsAsync(
    new CommandsSubscription { Channel = "ch", Group = "group" })) { }

// Subscribe to Queries
await foreach (var q in client.SubscribeToQueriesAsync(
    new QueriesSubscription { Channel = "ch", Group = "group" })) { }
// Send Command
client.sendCommand(CommandMessage(
    channel = "ch", body = "data".toByteArray(), timeout = 10000))

// Send Query
client.sendQuery(QueryMessage(
    channel = "ch", body = "data".toByteArray(), timeout = 10000))

// Send Query with Cache
client.sendQuery(QueryMessage(
    channel = "ch", body = "data".toByteArray(), timeout = 10000,
    cacheKey = "key", cacheTTL = 60000))

// Subscribe to Commands
client.subscribeToCommands(
    channel = "ch", group = "group",
    onCommand = handler, onError = errHandler)

// Subscribe to Queries
client.subscribeToQueries(
    channel = "ch", group = "group",
    onQuery = handler, onError = errHandler)
// Send Command
kubemq::CommandMessage cmd;
cmd.channel = "ch"; cmd.body = "data"; cmd.timeout = 10000;
client.sendCommand(cmd);

// Send Query
kubemq::QueryMessage query;
query.channel = "ch"; query.body = "data"; query.timeout = 10000;
client.sendQuery(query);

// Send Query with Cache
query.cacheKey = "key"; query.cacheTTL = 60000;
client.sendQuery(query);

// Subscribe to Commands
client.subscribeToCommands("ch", "group", handler, errHandler);

// Subscribe to Queries
client.subscribeToQueries("ch", "group", handler, errHandler);
// Send Command
let command = CommandBuilder::new()
    .channel("ch").body(b"data".to_vec())
    .timeout(Duration::from_secs(10)).build();
client.send_command(command).await?;

// Send Query
let query = QueryBuilder::new()
    .channel("ch").body(b"data".to_vec())
    .timeout(Duration::from_secs(10)).build();
client.send_query(query).await?;

// Send Query with Cache
let query = QueryBuilder::new()
    .channel("ch").body(b"data".to_vec())
    .timeout(Duration::from_secs(10))
    .cache_key("key").cache_ttl(Duration::from_secs(60)).build();
client.send_query(query).await?;

// Subscribe to Commands
client.subscribe_to_commands("ch", "group", handler, None).await?;

// Subscribe to Queries
client.subscribe_to_queries("ch", "group", handler, None).await?;
# Send Command
msg = KubeMQ::CQ::CommandMessage.new(
  channel: "ch", body: "data", timeout: 10)
client.send_command(msg)

# Send Query
msg = KubeMQ::CQ::QueryMessage.new(
  channel: "ch", body: "data", timeout: 10)
client.send_query(msg)

# Send Query with Cache
msg = KubeMQ::CQ::QueryMessage.new(
  channel: "ch", body: "data", timeout: 10,
  cache_key: "key", cache_ttl: 60)
client.send_query(msg)

# Subscribe to Commands
sub = KubeMQ::CQ::CommandsSubscription.new(channel: "ch", group: "group")
client.subscribe_to_commands(sub, cancellation_token: cancel,
  on_error: err_handler) { |cmd| handler.call(cmd) }

# Subscribe to Queries
sub = KubeMQ::CQ::QueriesSubscription.new(channel: "ch", group: "group")
client.subscribe_to_queries(sub, cancellation_token: cancel,
  on_error: err_handler) { |query| handler.call(query) }
# Send Command
command = KubeMQ.Command.new(
  channel: "ch", body: "data", timeout: 10_000)
KubeMQ.Client.send_command(client, command)

# Send Query
query = KubeMQ.Query.new(
  channel: "ch", body: "data", timeout: 10_000)
KubeMQ.Client.send_query(client, query)

# Send Query with Cache
query = KubeMQ.Query.new(
  channel: "ch", body: "data", timeout: 10_000,
  cache_key: "key", cache_ttl: 60_000)
KubeMQ.Client.send_query(client, query)

# Subscribe to Commands
KubeMQ.Client.subscribe_to_commands(client, "ch",
  group: "group", on_command: handler, on_error: err_handler)

# Subscribe to Queries
KubeMQ.Client.subscribe_to_queries(client, "ch",
  group: "group", on_query: handler, on_error: err_handler)

Was this page helpful?

On this page