Handle Command
Register a handler for incoming KubeMQ commands using the C# SDK RPC client.
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 by iterating SubscribeToCommandsAsync, and KubeMQ delivers every matching command on that channel to it as a long-lived, server-streamed IAsyncEnumerable<CommandReceived>, turning the channel into a synchronous RPC endpoint.
Handling happens inside the await foreach loop: each iteration yields a command with RequestId, ReplyChannel, and Body; you run your business logic, then send a reply with SendCommandResponseAsync(new CommandResponse { RequestId = cmd.RequestId, ReplyChannel = cmd.ReplyChannel, Executed = true }). Copying RequestId and ReplyChannel from the received command is what lets the broker correlate the reply back to the exact caller blocked on the send — nothing else identifies which request the response belongs to.
Gotchas: the reply must be sent before the caller's timeout elapses or the caller sees a timeout even if you eventually respond; set Executed = false and populate Error to signal a business-logic failure rather than a successful run; and the await foreach loop processes commands one at a time, so slow business logic head-of-line blocks the next command.
Prerequisites
- KubeMQ server running on
localhost:50000 - C# SDK installed (
dotnet add package KubeMQ.SDK.CSharp)
Code
// KubeMQ .NET SDK — Commands: Handle Command
//
// This example demonstrates subscribing to incoming commands and responding.
// Run this before Commands.SendCommand to handle the request.
//
// Prerequisites:
// - KubeMQ server running on localhost:50000
// - dotnet run
using KubeMQ.Sdk.Client;
using KubeMQ.Sdk.Commands;
using System.Text;
await using var client = new KubeMQClient(new KubeMQClientOptions
{
ClientId = "csharp-commands-handle-command-client",
});
await client.ConnectAsync();
Console.WriteLine("Waiting for commands on 'csharp-demo.commands'...");
var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, e) =>
{
e.Cancel = true;
cts.Cancel();
};
await foreach (var cmd in client.SubscribeToCommandsAsync(
new CommandsSubscription { Channel = "csharp-commands.handle-command" }, cts.Token))
{
var body = Encoding.UTF8.GetString(cmd.Body.Span);
Console.WriteLine($"Received command: {body}");
// Process the command and respond
await client.SendCommandResponseAsync(new CommandResponse
{
RequestId = cmd.RequestId,
ReplyChannel = cmd.ReplyChannel!,
Executed = true,
});
Console.WriteLine(" -> Responded: executed=true");
}
Console.WriteLine("Done.");
How It Works
SubscribeToCommandsAsyncreturns anIAsyncEnumerable<CommandReceived>. Theawait foreachloop drives the gRPC stream; each iteration yields one incoming command.cmd.RequestIdandcmd.ReplyChannelare populated by the broker — both must be copied into the response for the broker to correlate and route the reply back to the sender.Executed = truesignals successful processing. Set it tofalseand populate the optionalErrorstring to return a failure response.- The
CancellationTokenSourcewired toConsole.CancelKeyPressgives this long-running handler a clean shutdown path on Ctrl+C.
Related
Was this page helpful?