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 Field | Command | Query |
|---|---|---|
Body | Always nil (stripped) | Preserved |
Metadata | Always "" (stripped) | Preserved |
CacheHit | Always false (stripped) | Preserved |
ReplyChannel | Always "" (stripped) | Always "" (stripped) |
Executed | Preserved | Preserved |
Error | Preserved | Preserved |
Validation Rules
Channel Name
| Rule | Constraint | Error Code |
|---|---|---|
| Required | Cannot be empty | 102 |
| No trailing dot | Cannot end with . | 119 |
| No whitespace | Cannot contain spaces | 108 |
| No wildcards | Cannot contain * or > | 107 |
Valid channel name regex: ^[^\s*>]+[^.]$
Client ID
| Rule | Constraint | Error Code |
|---|---|---|
| Required | Cannot be empty | 101 |
| Alphanumeric | Must 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
| Rule | Constraint | Error Code |
|---|---|---|
| Required | Must be greater than zero | 109 |
Cache (Queries Only)
| Rule | Constraint | Error Code |
|---|---|---|
| CacheTTL required | If CacheKey is set, CacheTTL must be > 0 | 116 |
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
groupis 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
| Method | Type | Description |
|---|---|---|
SendRequest(Request) → Response | Unary | Send command or query |
SendResponse(Response) → Empty | Unary | Send response back to requester |
SubscribeToRequests(Subscribe) → stream Request | Server stream | Subscribe to commands or queries |
Default port: 50000
REST
| Method | Path | Description |
|---|---|---|
POST | /send/request | Send command or query |
POST | /send/response | Send RPC response |
GET | /subscribe/requests | WebSocket: subscribe to commands or queries |
Default port: 9090
REST subscription query parameters:
| Parameter | Description | Example |
|---|---|---|
client_id | Client identifier | my-responder |
channel | Channel name | orders.process |
group | Load balancing group | workers |
subscribe_type | commands or queries | commands |
Internal Channel Mapping
| Pattern | Channel Prefix | Example |
|---|---|---|
| Commands | _COMMANDS_. | _COMMANDS_.orders.process |
| Queries | _QUERIES_. | _QUERIES_.inventory.lookup |
Middleware Chain
Command Sender
Request → Logging → Monitor → Metrics → broker requestQuery Sender
Request → Logging → Monitor → Cache → Metrics → broker requestReceiver (Commands and Queries)
broker queue-subscribe → Logging → reqCh deliveryError Codes
Input Validation
| Code | Error | Description |
|---|---|---|
| 101 | Invalid ClientID | ClientID is empty |
| 102 | Invalid Channel | Channel is empty |
| 107 | Invalid Channel | Channel contains wildcards (* or >) |
| 108 | Invalid Channel | Channel contains whitespace |
| 109 | Invalid Timeout | Timeout is zero or negative |
| 115 | Invalid Request | Both body and metadata are empty |
| 116 | Invalid CacheTTL | CacheKey is set but CacheTTL is zero or negative |
| 117 | Invalid RequestID | RequestID is empty (on response) |
| 119 | Invalid Channel | Channel ends with . |
Runtime Errors
| Code | Error | Description |
|---|---|---|
| 206 | Invalid Request Type | Request type is not Command or Query |
| 207 | Invalid Subscribe Type | Subscribe type is not Commands or Queries |
| 301 | Request Timeout | No reply received before timeout expired |
| 302 | Connection Unavailable | broker connection is down |
| 303 | Invalid Response Format | Reply data cannot be unmarshaled |
| 409 | Shutdown Mode | Server is shutting down, all operations rejected |
| 412 | Access Denied | Authorization denied for the resource |
Delivery Semantics
| Aspect | Behavior |
|---|---|
| Delivery guarantee | At-most-once (request is sent once, no automatic retry) |
| Ordering | Requests are independent (no ordering guarantee between requests) |
| Timeout | Sender blocks until response arrives or timeout expires |
| Acknowledgment | Implicit — 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)Related
- Getting Started — send your first command and query
- Configure Timeouts — per-request timeouts and retries
- Load Balancing — distribute requests across responders
- Events Reference — for the fire-and-forget pattern
- Queues Reference — for guaranteed delivery
Was this page helpful?