Device Command & Control
Send commands to IoT devices and query their status using KubeMQ RPC.
Scenario
A central controller sends commands to field devices (restart, configure, update firmware) and queries their telemetry data (temperature, battery level, connectivity status). Each device runs an agent that subscribes to its own command/query channel. A status dashboard uses cached queries to display device status without hitting each device on every request.
Architecture
A controller issues commands and queries through KubeMQ; each device agent owns its command and query channels, while the dashboard reads device status through cached queries.
Implementation
Device Agent (Responder)
Each device agent subscribes to commands and queries on its own channel (e.g., device.sensor-001).
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"os"
"time"
"github.com/kubemq-io/kubemq-go/v2"
)
func main() {
deviceId := os.Getenv("DEVICE_ID")
channel := "device." + deviceId
ctx := context.Background()
client, err := kubemq.NewClient(ctx,
kubemq.WithAddress("localhost", 50000),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
client.SubscribeToCommands(ctx, channel, "",
kubemq.WithOnCommandReceive(func(cmd *kubemq.CommandReceive) {
action := string(cmd.Body)
fmt.Printf("[%s] Executing: %s\n", deviceId, action)
switch action {
case "restart":
fmt.Printf("[%s] Restarting...\n", deviceId)
case "configure":
fmt.Printf("[%s] Applying config from metadata\n", deviceId)
}
resp := kubemq.NewCommandReply().
SetRequestId(cmd.Id).
SetResponseTo(cmd.ResponseTo).
SetExecutedAt(time.Now())
client.SendCommandResponse(ctx, resp)
}),
kubemq.WithOnError(func(err error) { log.Println("CMD error:", err) }),
)
client.SubscribeToQueries(ctx, channel, "",
kubemq.WithOnQueryReceive(func(query *kubemq.QueryReceive) {
fmt.Printf("[%s] Status query received\n", deviceId)
status := map[string]interface{}{
"deviceId": deviceId, "temperature": 42.5,
"battery": 87, "online": true,
}
data, _ := json.Marshal(status)
resp := kubemq.NewQueryReply().
SetRequestId(query.Id).
SetResponseTo(query.ResponseTo).
SetBody(data).
SetExecutedAt(time.Now())
client.SendQueryResponse(ctx, resp)
}),
kubemq.WithOnError(func(err error) { log.Println("QRY error:", err) }),
)
fmt.Printf("Device agent %s ready\n", deviceId)
<-ctx.Done()
}import os
import json
import time
from kubemq.cq import (
Client as CQClient,
CommandsSubscription, CommandReceived, CommandResponse,
QueriesSubscription, QueryReceived, QueryResponse,
CancellationToken,
)
device_id = os.getenv("DEVICE_ID", "sensor-001")
channel = f"device.{device_id}"
client = CQClient(address="localhost:50000")
cancel = CancellationToken()
def on_command(request: CommandReceived) -> None:
action = request.body.decode("utf-8")
print(f"[{device_id}] Executing: {action}")
client.send_response_message(
CommandResponse(command_received=request, is_executed=True))
def on_query(request: QueryReceived) -> None:
print(f"[{device_id}] Status query received")
status = {"deviceId": device_id, "temperature": 42.5,
"battery": 87, "online": True}
client.send_response_message(
QueryResponse(query_received=request, is_executed=True,
body=json.dumps(status).encode()))
client.subscribe_to_commands(CommandsSubscription(
channel=channel, on_receive_command_callback=on_command,
on_error_callback=lambda e: print(f"Error: {e}")), cancel=cancel)
client.subscribe_to_queries(QueriesSubscription(
channel=channel, on_receive_query_callback=on_query,
on_error_callback=lambda e: print(f"Error: {e}")), cancel=cancel)
print(f"Device agent {device_id} ready")
time.sleep(3600)const { KubeMQClient } = require("kubemq-js");
const deviceId = process.env.DEVICE_ID || "sensor-001";
const channel = `device.${deviceId}`;
const client = new KubeMQClient({ address: "localhost:50000" });
client.subscribeToCommands({
channel,
onCommand: (cmd) => {
const action = Buffer.from(cmd.body).toString();
console.log(`[${deviceId}] Executing: ${action}`);
client.sendCommandResponse({ requestId: cmd.id, isExecuted: true });
},
onError: (err) => console.error("CMD error:", err.message),
});
client.subscribeToQueries({
channel,
onQuery: (query) => {
console.log(`[${deviceId}] Status query received`);
client.sendQueryResponse({
requestId: query.id,
isExecuted: true,
body: Buffer.from(JSON.stringify({
deviceId, temperature: 42.5, battery: 87, online: true,
})),
});
},
onError: (err) => console.error("QRY error:", err.message),
});
console.log(`Device agent ${deviceId} ready`);String deviceId = System.getenv("DEVICE_ID");
String channel = "device." + deviceId;
CQClient client = CQClient.builder()
.address("localhost:50000").clientId(deviceId).build();
client.subscribeToCommands(CommandsSubscription.builder()
.channel(channel)
.onReceiveCommandCallback(cmd -> {
System.out.printf("[%s] Executing: %s%n", deviceId, new String(cmd.getBody()));
return CommandResponseMessage.builder()
.requestId(cmd.getId()).isExecuted(true).build();
})
.onErrorCallback(err -> System.err.println("Error: " + err.getMessage()))
.build());
client.subscribeToQueries(QueriesSubscription.builder()
.channel(channel)
.onReceiveQueryCallback(query -> {
System.out.printf("[%s] Status query received%n", deviceId);
String status = String.format(
"{\"deviceId\":\"%s\",\"temperature\":42.5,\"battery\":87,\"online\":true}",
deviceId);
return QueryResponseMessage.builder()
.requestId(query.getId()).isExecuted(true)
.body(status.getBytes()).build();
})
.onErrorCallback(err -> System.err.println("Error: " + err.getMessage()))
.build());
System.out.printf("Device agent %s ready%n", deviceId);
Thread.sleep(3600000);var deviceId = Environment.GetEnvironmentVariable("DEVICE_ID") ?? "sensor-001";
var channel = $"device.{deviceId}";
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();
var cmdTask = Task.Run(async () =>
{
await foreach (var cmd in client.SubscribeToCommandsAsync(
new CommandsSubscription { Channel = channel }))
{
Console.WriteLine($"[{deviceId}] Executing: {Encoding.UTF8.GetString(cmd.Body.Span)}");
await client.SendCommandResponseAsync(new CommandResponse
{ RequestId = cmd.Id, IsExecuted = true });
}
});
var queryTask = Task.Run(async () =>
{
await foreach (var query in client.SubscribeToQueriesAsync(
new QueriesSubscription { Channel = channel }))
{
Console.WriteLine($"[{deviceId}] Status query received");
await client.SendQueryResponseAsync(new QueryResponse
{
RequestId = query.Id, IsExecuted = true,
Body = Encoding.UTF8.GetBytes(
$"{{\"deviceId\":\"{deviceId}\",\"temperature\":42.5,\"battery\":87,\"online\":true}}")
});
}
});
Console.WriteLine($"Device agent {deviceId} ready");
await Task.WhenAll(cmdTask, queryTask);val deviceId = System.getenv("DEVICE_ID") ?: "sensor-001"
val channel = "device.$deviceId"
val client = CQClient("localhost:50000")
client.subscribeToCommands(channel = channel,
onCommand = { cmd ->
println("[$deviceId] Executing: ${String(cmd.body)}")
client.sendCommandResponse(requestId = cmd.id, isExecuted = true)
},
onError = { err -> System.err.println("Error: ${err.message}") })
client.subscribeToQueries(channel = channel,
onQuery = { query ->
println("[$deviceId] Status query received")
val status = """{"deviceId":"$deviceId","temperature":42.5,"battery":87,"online":true}"""
client.sendQueryResponse(requestId = query.id, isExecuted = true,
body = status.toByteArray())
},
onError = { err -> System.err.println("Error: ${err.message}") })
println("Device agent $deviceId ready")
Thread.sleep(3600000)std::string deviceId = std::getenv("DEVICE_ID") ? std::getenv("DEVICE_ID") : "sensor-001";
std::string channel = "device." + deviceId;
auto client = kubemq::CQClient("localhost:50000");
client.subscribeToCommands(channel, "",
[&](const kubemq::CommandReceive& cmd) {
std::cout << "[" << deviceId << "] Executing: " << cmd.body << std::endl;
client.sendCommandResponse(cmd.id, true);
},
[](const std::string& err) { std::cerr << "Error: " << err << std::endl; });
client.subscribeToQueries(channel, "",
[&](const kubemq::QueryReceive& query) {
std::cout << "[" << deviceId << "] Status query received" << std::endl;
std::string status = R"({"deviceId":")" + deviceId +
R"(","temperature":42.5,"battery":87,"online":true})";
client.sendQueryResponse(query.id, true, status);
},
[](const std::string& err) { std::cerr << "Error: " << err << std::endl; });
std::cout << "Device agent " << deviceId << " ready" << std::endl;
std::this_thread::sleep_for(std::chrono::hours(1));use kubemq::prelude::*;
use kubemq::{CommandReplyBuilder, QueryReplyBuilder};
use std::env;
use std::time::Duration;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let device_id = env::var("DEVICE_ID").unwrap_or_else(|_| "sensor-001".into());
let channel = format!("device.{}", device_id);
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.build()
.await?;
let cmd_client = client.clone();
let cmd_device = device_id.clone();
let cmd_sub = client
.subscribe_to_commands(
&channel,
"",
move |cmd| {
let c = cmd_client.clone();
let id = cmd_device.clone();
Box::pin(async move {
println!("[{}] Executing: {}", id, String::from_utf8_lossy(&cmd.body));
let reply = CommandReplyBuilder::new()
.request_id(&cmd.id)
.response_to(&cmd.response_to)
.build();
tokio::spawn(async move {
let _ = c.send_command_response(reply).await;
});
})
},
None,
)
.await?;
let qry_client = client.clone();
let qry_device = device_id.clone();
let qry_sub = client
.subscribe_to_queries(
&channel,
"",
move |query| {
let c = qry_client.clone();
let id = qry_device.clone();
Box::pin(async move {
println!("[{}] Status query received", id);
let status = format!(
r#"{{"deviceId":"{}","temperature":42.5,"battery":87,"online":true}}"#,
id
);
let reply = QueryReplyBuilder::new()
.request_id(&query.id)
.response_to(&query.response_to)
.body(status.into_bytes())
.build();
tokio::spawn(async move {
let _ = c.send_query_response(reply).await;
});
})
},
None,
)
.await?;
println!("Device agent {} ready", device_id);
tokio::signal::ctrl_c().await.ok();
cmd_sub.unsubscribe().await;
qry_sub.unsubscribe().await;
client.close().await?;
Ok(())
}require 'kubemq'
require 'json'
device_id = ENV.fetch('DEVICE_ID', 'sensor-001')
channel = "device.#{device_id}"
client = KubeMQ::CQClient.new(address: 'localhost:50000', client_id: device_id)
cancel = KubeMQ::CancellationToken.new
cmd_sub = KubeMQ::CQ::CommandsSubscription.new(channel: channel)
client.subscribe_to_commands(cmd_sub, cancellation_token: cancel,
on_error: ->(e) { puts "CMD error: #{e.message}" }) do |cmd|
puts "[#{device_id}] Executing: #{cmd.body}"
client.send_response(KubeMQ::CQ::CommandResponseMessage.new(
request_id: cmd.id, reply_channel: cmd.reply_channel, executed: true))
end
qry_sub = KubeMQ::CQ::QueriesSubscription.new(channel: channel)
client.subscribe_to_queries(qry_sub, cancellation_token: cancel,
on_error: ->(e) { puts "QRY error: #{e.message}" }) do |query|
puts "[#{device_id}] Status query received"
status = { deviceId: device_id, temperature: 42.5, battery: 87, online: true }
client.send_response(KubeMQ::CQ::QueryResponseMessage.new(
request_id: query.id, reply_channel: query.reply_channel,
executed: true, body: status.to_json))
end
puts "Device agent #{device_id} ready"
cancel.waitdevice_id = System.get_env("DEVICE_ID") || "sensor-001"
channel = "device.#{device_id}"
{:ok, client} =
KubeMQ.Client.start_link(address: "localhost:50000", client_id: device_id)
{:ok, _cmd_sub} =
KubeMQ.Client.subscribe_to_commands(client, channel,
on_command: fn cmd ->
IO.puts("[#{device_id}] Executing: #{cmd.body}")
KubeMQ.CommandReply.new(
request_id: cmd.id,
response_to: cmd.reply_channel,
executed: true
)
end,
on_error: fn err -> IO.puts("CMD error: #{err.message}") end
)
{:ok, _qry_sub} =
KubeMQ.Client.subscribe_to_queries(client, channel,
on_query: fn query ->
IO.puts("[#{device_id}] Status query received")
status = ~s({"deviceId":"#{device_id}","temperature":42.5,"battery":87,"online":true})
KubeMQ.QueryReply.new(
request_id: query.id,
response_to: query.reply_channel,
executed: true,
body: status
)
end,
on_error: fn err -> IO.puts("QRY error: #{err.message}") end
)
IO.puts("Device agent #{device_id} ready")
Process.sleep(:infinity)Controller (Sender)
The controller sends commands to specific devices and queries their status.
resp, _ := client.SendCommand(ctx, kubemq.NewCommand().
SetChannel("device.sensor-001").
SetBody([]byte("restart")).
SetTimeout(15 * time.Second))
log.Printf("Restart sensor-001: Executed=%v", resp.Executed)
status, _ := client.SendQuery(ctx, kubemq.NewQuery().
SetChannel("device.sensor-001").
SetBody([]byte("get-status")).
SetTimeout(10 * time.Second))
log.Printf("Status: %s", status.Body)resp = client.send_command(CommandMessage(
channel="device.sensor-001", body=b"restart", timeout_in_seconds=15))
print(f"Restart sensor-001: Executed={resp.is_executed}")
status = client.send_query(QueryMessage(
channel="device.sensor-001", body=b"get-status", timeout_in_seconds=10))
print(f"Status: {status.body.decode('utf-8')}")const resp = await client.sendCommand({
channel: "device.sensor-001", body: Buffer.from("restart"), timeoutInSeconds: 15,
});
console.log("Restart sensor-001: Executed=", resp.isExecuted);
const status = await client.sendQuery({
channel: "device.sensor-001", body: Buffer.from("get-status"), timeoutInSeconds: 10,
});
console.log("Status:", Buffer.from(status.body).toString());var resp = client.sendCommandRequest(CommandMessage.builder()
.channel("device.sensor-001").body("restart".getBytes()).timeout(15000).build());
System.out.println("Restart sensor-001: Executed=" + resp.isExecuted());
var status = client.sendQueryRequest(QueryMessage.builder()
.channel("device.sensor-001").body("get-status".getBytes()).timeout(10000).build());
System.out.println("Status: " + new String(status.getBody()));var resp = await client.SendCommandAsync(new CommandMessage
{
Channel = "device.sensor-001", Body = Encoding.UTF8.GetBytes("restart"),
Timeout = TimeSpan.FromSeconds(15)
});
Console.WriteLine($"Restart sensor-001: Executed={resp.IsExecuted}");
var status = await client.SendQueryAsync(new QueryMessage
{
Channel = "device.sensor-001", Body = Encoding.UTF8.GetBytes("get-status"),
Timeout = TimeSpan.FromSeconds(10)
});
Console.WriteLine($"Status: {Encoding.UTF8.GetString(status.Body.Span)}");val resp = client.sendCommand(CommandMessage(
channel = "device.sensor-001", body = "restart".toByteArray(), timeout = 15000))
println("Restart sensor-001: Executed=${resp.isExecuted}")
val status = client.sendQuery(QueryMessage(
channel = "device.sensor-001", body = "get-status".toByteArray(), timeout = 10000))
println("Status: ${String(status.body)}")kubemq::CommandMessage cmd;
cmd.channel = "device.sensor-001";
cmd.body = "restart";
cmd.timeout = 15000;
auto resp = client.sendCommand(cmd);
std::cout << "Restart: Executed=" << resp.isExecuted << std::endl;
kubemq::QueryMessage query;
query.channel = "device.sensor-001";
query.body = "get-status";
query.timeout = 10000;
auto status = client.sendQuery(query);
std::cout << "Status: " << status.body << std::endl;use kubemq::prelude::*;
use kubemq::{CommandBuilder, QueryBuilder};
use std::time::Duration;
let command = CommandBuilder::new()
.channel("device.sensor-001")
.body(b"restart".to_vec())
.timeout(Duration::from_secs(15))
.build();
let resp = client.send_command(command).await?;
println!("Restart sensor-001: executed={}", resp.executed);
let query = QueryBuilder::new()
.channel("device.sensor-001")
.body(b"get-status".to_vec())
.timeout(Duration::from_secs(10))
.build();
let status = client.send_query(query).await?;
println!("Status: {}", String::from_utf8_lossy(&status.body));require 'kubemq'
client = KubeMQ::CQClient.new(address: 'localhost:50000', client_id: 'controller')
resp = client.send_command(KubeMQ::CQ::CommandMessage.new(
channel: 'device.sensor-001', timeout: 15, body: 'restart'))
puts "Restart sensor-001: executed=#{resp.executed}"
status = client.send_query(KubeMQ::CQ::QueryMessage.new(
channel: 'device.sensor-001', timeout: 10, body: 'get-status'))
puts "Status: #{status.body}"command =
KubeMQ.Command.new(channel: "device.sensor-001", body: "restart", timeout: 15_000)
{:ok, resp} = KubeMQ.Client.send_command(client, command)
IO.puts("Restart sensor-001: executed=#{resp.executed}")
query =
KubeMQ.Query.new(channel: "device.sensor-001", body: "get-status", timeout: 10_000)
{:ok, status} = KubeMQ.Client.send_query(client, query)
IO.puts("Status: #{status.body}")Status Dashboard (Cached Query)
The dashboard queries device status with caching to avoid hitting each device on every page load.
status, _ := client.SendQuery(ctx, kubemq.NewQuery().
SetChannel("device.sensor-001").
SetBody([]byte("get-status")).
SetTimeout(10 * time.Second).
SetCacheKey("device-status-sensor-001").
SetCacheTTL(10 * time.Second))
log.Printf("Status (cache hit: %v): %s", status.CacheHit, status.Body)status = client.send_query(QueryMessage(
channel="device.sensor-001", body=b"get-status",
timeout_in_seconds=10,
cache_key="device-status-sensor-001", cache_ttl_in_seconds=10))
print(f"Status (cache hit: {status.cache_hit}): {status.body.decode('utf-8')}")const status = await client.sendQuery({
channel: "device.sensor-001", body: Buffer.from("get-status"),
timeoutInSeconds: 10,
cacheKey: "device-status-sensor-001", cacheTTL: 10000,
});
console.log(`Status (cache hit: ${status.cacheHit}):`,
Buffer.from(status.body).toString());var status = client.sendQueryRequest(QueryMessage.builder()
.channel("device.sensor-001").body("get-status".getBytes())
.timeout(10000).cacheKey("device-status-sensor-001").cacheTTL(10000).build());
System.out.printf("Status (cache hit: %s): %s%n",
status.isCacheHit(), new String(status.getBody()));var status = await client.SendQueryAsync(new QueryMessage
{
Channel = "device.sensor-001", Body = Encoding.UTF8.GetBytes("get-status"),
Timeout = TimeSpan.FromSeconds(10),
CacheKey = "device-status-sensor-001", CacheTTL = TimeSpan.FromSeconds(10)
});
Console.WriteLine($"Status (cache hit: {status.CacheHit}): "
+ Encoding.UTF8.GetString(status.Body.Span));val status = client.sendQuery(QueryMessage(
channel = "device.sensor-001", body = "get-status".toByteArray(),
timeout = 10000, cacheKey = "device-status-sensor-001", cacheTTL = 10000))
println("Status (cache hit: ${status.cacheHit}): ${String(status.body)}")kubemq::QueryMessage query;
query.channel = "device.sensor-001";
query.body = "get-status";
query.timeout = 10000;
query.cacheKey = "device-status-sensor-001";
query.cacheTTL = 10000;
auto status = client.sendQuery(query);
std::cout << "Status (cache hit: " << status.cacheHit << "): "
<< status.body << std::endl;use kubemq::prelude::*;
use kubemq::QueryBuilder;
use std::time::Duration;
let query = QueryBuilder::new()
.channel("device.sensor-001")
.body(b"get-status".to_vec())
.timeout(Duration::from_secs(10))
.cache_key("device-status-sensor-001")
.cache_ttl(Duration::from_secs(10))
.build();
let status = client.send_query(query).await?;
println!(
"Status (cache hit: {}): {}",
status.cache_hit,
String::from_utf8_lossy(&status.body)
);require 'kubemq'
client = KubeMQ::CQClient.new(address: 'localhost:50000', client_id: 'dashboard')
status = client.send_query(KubeMQ::CQ::QueryMessage.new(
channel: 'device.sensor-001', timeout: 10, body: 'get-status',
cache_key: 'device-status-sensor-001', cache_ttl: 10))
puts "Status (cache hit: #{status.cache_hit}): #{status.body}"query =
KubeMQ.Query.new(
channel: "device.sensor-001",
body: "get-status",
timeout: 10_000,
cache_key: "device-status-sensor-001",
cache_ttl: 10_000
)
{:ok, status} = KubeMQ.Client.send_query(client, query)
IO.puts("Status (cache hit: #{status.cache_hit}): #{status.body}")Production Considerations
Was this page helpful?