Handle Command
Register a handler for incoming KubeMQ Commands and send replies using the Go SDK.
Overview
A command handler is the receiving side of KubeMQ's Commands pattern — the code that actually does the work a caller is blocked waiting on. Instead of building your own request-routing layer on top of a queue, you register a handler once with SubscribeToCommands and KubeMQ delivers every matching command on that channel to it as a long-lived, server-streamed subscription, turning the channel into a synchronous RPC endpoint.
Handling happens inside the WithOnCommandReceive callback: you read the command's Id, Body, and Metadata, run your business logic, then build a reply with NewCommandReply().SetRequestId(cmd.Id).SetResponseTo(cmd.ResponseTo) and send it with SendCommandResponse. Copying RequestId and ResponseTo from the received command is what lets the broker correlate the reply back to the exact caller blocked on SendCommand — nothing else identifies which request the response belongs to.
Gotchas: the reply must be sent before the caller's SetTimeout deadline or the caller sees a timeout even if you eventually respond; the callback runs on a shared delivery path, so slow or blocking business logic head-of-line blocks the next command; and an uncaught panic inside the callback can take down the subscription without ever notifying the caller.
Prerequisites
- KubeMQ server running on
localhost:50000 - Go SDK installed (
go get github.com/kubemq-io/kubemq-go/v2)
Code
// Example: commands/handle-command
//
// Demonstrates subscribing to commands and handling them with business logic.
// The handler processes incoming commands and sends back responses.
//
// Channel: go-commands.handle-command
// Client ID: go-commands-handle-command-client
//
// Run with a KubeMQ server on localhost:50000
// (see https://docs.kubemq.io/deploy).
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/kubemq-io/kubemq-go/v2"
)
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
client, err := kubemq.NewClient(ctx,
kubemq.WithAddress("localhost", 50000),
kubemq.WithClientId("go-commands-handle-command-client"),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
channel := "go-commands.handle-command"
done := make(chan struct{})
// Register a command handler that processes incoming commands.
sub, err := client.SubscribeToCommands(ctx, channel, "",
kubemq.WithOnCommandReceive(func(cmd *kubemq.CommandReceive) {
fmt.Printf("Handling command: id=%s body=%s metadata=%s\n",
cmd.Id, cmd.Body, cmd.Metadata)
// Process the command (business logic goes here).
// Then send back a response.
resp := kubemq.NewCommandReply().
SetRequestId(cmd.Id).
SetResponseTo(cmd.ResponseTo).
SetExecutedAt(time.Now())
if err := client.SendCommandResponse(ctx, resp); err != nil {
log.Printf("Failed to send response: %v", err)
}
close(done)
}),
kubemq.WithOnError(func(err error) {
log.Println("Handler error:", err)
}),
)
if err != nil {
log.Fatal(err)
}
defer sub.Unsubscribe()
time.Sleep(300 * time.Millisecond)
// Send a command to trigger the handler.
cmdResp, err := client.SendCommand(ctx, kubemq.NewCommand().
SetChannel(channel).
SetBody([]byte("process-order")).
SetMetadata("order-123").
SetTimeout(10*time.Second))
if err != nil {
log.Fatal(err)
}
fmt.Printf("Response: executed=%v\n", cmdResp.Executed)
<-done
}
How It Works
client.SubscribeToCommands(ctx, channel, "", opts...)opens a long-lived gRPC subscription that callskubemq.WithOnCommandReceivefor every incoming command ongo-commands.handle-command.- Inside the callback,
kubemq.NewCommandReply().SetRequestId(cmd.Id).SetResponseTo(cmd.ResponseTo).SetExecutedAt(time.Now())constructs the response;SetRequestIdandSetResponseToare mandatory for the broker to route the reply to the correct caller. client.SendCommandResponse(ctx, resp)is called from within the callback goroutine and sends the reply back to the broker; the caller'sSendCommandunblocks when the reply arrives.defer sub.Unsubscribe()cancels the subscription cleanly; any in-flight command deliveries that arrive afterUnsubscribeare dropped by the broker.
Related
Was this page helpful?