Send Query
Send a KubeMQ query and receive a data response from a handler using the Elixir SDK.
Overview
This tutorial builds the RPC half of KubeMQ's request/reply patterns: a query, where the caller blocks for a handler's data payload instead of just a completion status. Reach for it whenever a caller needs an answer — fetching a record, running a lookup, or asking another service to compute a value on demand. You'll run a handler and a sender in the same process to see the full round trip.
KubeMQ.Client.subscribe_to_queries(client, channel, on_query: fn query -> ... end) registers a handler; the callback returns a KubeMQ.QueryReply.new(request_id: query.id, response_to: query.reply_channel, ...) — copied from the incoming query — plus executed: true and a body. The sender builds a KubeMQ.Query.new(...) and calls KubeMQ.Client.send_query(client, query), blocking until {:ok, response} arrives; KubeMQ routes the reply back to the caller waiting on it.
Gotchas: the timeout must cover however long the handler takes to run — a slow handler returns {:error, err} even though the handler eventually succeeds. No handler subscribed yet also times out rather than erroring immediately. body is a raw payload — metadata is a convention for content type, but KubeMQ never parses either for you.
Prerequisites
- KubeMQ server running on
localhost:50000 - Elixir SDK installed (
{:kubemq, "~> 1.0"}in mix.exs)
Code
channel = "elixir-rpc.query-send"
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "elixir-query-sender")
{:ok, sub} =
KubeMQ.Client.subscribe_to_queries(client, channel,
on_query: fn query ->
IO.puts("[Handler] Query received: #{query.body}")
KubeMQ.QueryReply.new(
request_id: query.id,
response_to: query.reply_channel,
executed: true,
body: ~s({"name": "Widget", "price": 29.99}),
metadata: "application/json"
)
end
)
Process.sleep(500)
query = KubeMQ.Query.new(
channel: channel,
body: "sku-12345",
timeout: 10_000
)
case KubeMQ.Client.send_query(client, query) do
{:ok, response} ->
IO.puts("Query executed: #{response.executed}")
IO.puts("Response body: #{response.body}")
IO.puts("Response metadata: #{response.metadata}")
{:error, err} ->
IO.puts("Query failed: #{err.message}")
end
KubeMQ.Subscription.cancel(sub)
KubeMQ.Client.close(client)How It Works
- The
on_querycallback returns a%QueryReply{}with abodyfield for response data - Unlike commands, queries return data to the caller via the
response.bodyfield - The
metadatafield can be used to indicate content type
Related
Was this page helpful?