Send Command
Send a KubeMQ Command and wait for the handler's execution result using the Go SDK.
Overview
A command is KubeMQ's fire-and-confirm RPC pattern: reach for it when you need to know an action actually ran on the other end — "restart the service" — but don't need data back, just a yes/no on execution. It sits between one-way pub/sub, which gives no confirmation, and a query, which returns a result payload. Commands turn "I hope that worked" into a definite outcome your caller can branch on.
This sample builds that lesson: a handler subscribes with SubscribeToCommands, and the sender calls client.SendCommand, which blocks until a reply arrives or the timeout elapses. The reply is built with kubemq.NewCommandReply() and correlated back via SetRequestId(cmd.Id) and SetResponseTo(cmd.ResponseTo) — that correlation is what lets the broker route the response to the exact caller waiting on it, even with many senders sharing one channel.
Gotchas: if no handler is subscribed (or it's still starting up), SendCommand blocks for the full timeout before failing — there's no fast "nobody's listening" error. A handler that forgets to set SetResponseTo/SetRequestId on the reply leaves the caller hanging until timeout. And a command's reply carries no business data — if you need the handler to return a value, use a query instead.
Prerequisites
- KubeMQ server running on
localhost:50000 - Go SDK installed (
go get github.com/kubemq-io/kubemq-go/v2)
Code
// Example: commands/send-command
//
// Demonstrates sending a command (RPC-style request) and receiving a response.
// A handler subscribes to the command channel, processes the command,
// and sends back an execution response.
//
// Channel: go-commands.send-command
// Client ID: go-commands-send-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-send-command-client"),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
channel := "go-commands.send-command"
done := make(chan struct{})
// Subscribe to handle incoming commands.
sub, err := client.SubscribeToCommands(ctx, channel, "",
kubemq.WithOnCommandReceive(func(cmd *kubemq.CommandReceive) {
fmt.Printf("Command received: channel=%s body=%s\n", cmd.Channel, cmd.Body)
// Send a response indicating successful execution.
resp := kubemq.NewCommandReply().
SetRequestId(cmd.Id).
SetResponseTo(cmd.ResponseTo).
SetBody([]byte("executed")).
SetExecutedAt(time.Now())
_ = client.SendCommandResponse(ctx, resp)
close(done)
}),
kubemq.WithOnError(func(err error) {
log.Println("Command subscription error:", err)
}),
)
if err != nil {
log.Fatal(err)
}
defer sub.Unsubscribe()
time.Sleep(300 * time.Millisecond) // Allow subscription to establish.
// Send a command and wait for the response.
cmdResp, err := client.SendCommand(ctx, kubemq.NewCommand().
SetChannel(channel).
SetBody([]byte("do-something")).
SetTimeout(10*time.Second))
if err != nil {
log.Fatal(err)
}
fmt.Printf("Command response: executed=%v\n", cmdResp.Executed)
<-done
}
// Expected output:
// Command received: channel=go-commands.send-command body=do-something
// Command response: executed=true
How It Works
- A
SubscribeToCommandscall registers a command handler viakubemq.WithOnCommandReceive; the handler constructs akubemq.CommandReplysettingSetRequestId(cmd.Id)andSetResponseTo(cmd.ResponseTo)to route the response back to the correct caller. - A 300 ms sleep lets the subscription establish on the broker before the sender issues the command — in production, prefer
WithWaitForReadyor an explicitPing. client.SendCommand(ctx, kubemq.NewCommand().SetTimeout(10*time.Second))blocks until a handler responds or the timeout fires; the response is available incmdResp.Executed.client.SendCommandResponse(ctx, resp)sends the reply asynchronously back to the broker, which forwards it to the waitingSendCommandcaller.
Related
Was this page helpful?