Client SDKsElixirReference
Queues
Queue message structs, policies, and pull-based messaging — KubeMQ Elixir SDK reference.
Queues provide pull-based messaging with acknowledgement, visibility timeout, dead-letter routing, and transactional batch operations.
Structs
KubeMQ.QueueMessage
Used to send messages to a queue.
| Field | Type | Description |
|---|---|---|
id | String.t() | Auto-generated message ID |
channel | String.t() | Target queue channel |
metadata | String.t() | Optional metadata |
body | String.t() | binary() | Message payload |
client_id | String.t() | Sender client ID |
tags | map() | Optional key-value tags |
policy | QueuePolicy.t() | Delivery policy (delay, expiration, DLQ) |
attributes | QueueAttributes.t() | Server-assigned attributes (read-only) |
KubeMQ.QueuePolicy
Controls message delivery behavior.
| Field | Type | Description |
|---|---|---|
expiration_seconds | integer() | Message TTL in seconds |
delay_seconds | integer() | Defer delivery by N seconds |
max_receive_count | integer() | Max receives before DLQ routing |
max_receive_queue | String.t() | Dead-letter queue channel name |
KubeMQ.QueueAttributes
Server-assigned metadata on received messages.
| Field | Type | Description |
|---|---|---|
timestamp | integer() | Send timestamp |
sequence | integer() | Sequence number |
receive_count | integer() | Number of times received |
re_routed | boolean() | Whether message was re-routed |
re_routed_from_queue | String.t() | Original queue if re-routed |
KubeMQ.PollResponse
Returned by poll_queue/2 for transactional control.
| Function | Spec | Description |
|---|---|---|
ack_all/1 | (PollResponse.t()) :: {:ok, term()} | {:error, Error.t()} | Acknowledge all messages |
nack_all/1 | (PollResponse.t()) :: {:ok, term()} | {:error, Error.t()} | Reject all messages (re-queue) |
requeue_all/2 | (PollResponse.t(), channel) :: {:ok, term()} | {:error, Error.t()} | Move all to another queue |
ack_range/2 | (PollResponse.t(), sequences) :: :ok | {:error, Error.t()} | Ack specific sequences |
nack_range/2 | (PollResponse.t(), sequences) :: :ok | {:error, Error.t()} | Nack specific sequences |
Simple Send/Receive
msg = KubeMQ.QueueMessage.new(channel: "tasks", body: "process order")
{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
{:ok, recv} = KubeMQ.Client.receive_queue_messages(client, "tasks",
max_messages: 10, wait_timeout: 5_000)Stream API (Poll)
{:ok, poll} = KubeMQ.Client.poll_queue(client,
channel: "tasks", max_items: 10, wait_timeout: 5_000)
Enum.each(poll.messages, fn msg -> IO.puts(msg.body) end)
{:ok, _} = KubeMQ.PollResponse.ack_all(poll)Upstream Stream
{:ok, handle} = KubeMQ.Client.queue_upstream(client)
{:ok, results} = KubeMQ.QueueUpstreamHandle.send(handle, messages)
KubeMQ.QueueUpstreamHandle.close(handle)See Also
Was this page helpful?