KubeMQ
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.

FieldTypeDescription
idString.t()Auto-generated message ID
channelString.t()Target queue channel
metadataString.t()Optional metadata
bodyString.t() | binary()Message payload
client_idString.t()Sender client ID
tagsmap()Optional key-value tags
policyQueuePolicy.t()Delivery policy (delay, expiration, DLQ)
attributesQueueAttributes.t()Server-assigned attributes (read-only)

KubeMQ.QueuePolicy

Controls message delivery behavior.

FieldTypeDescription
expiration_secondsinteger()Message TTL in seconds
delay_secondsinteger()Defer delivery by N seconds
max_receive_countinteger()Max receives before DLQ routing
max_receive_queueString.t()Dead-letter queue channel name

KubeMQ.QueueAttributes

Server-assigned metadata on received messages.

FieldTypeDescription
timestampinteger()Send timestamp
sequenceinteger()Sequence number
receive_countinteger()Number of times received
re_routedboolean()Whether message was re-routed
re_routed_from_queueString.t()Original queue if re-routed

KubeMQ.PollResponse

Returned by poll_queue/2 for transactional control.

FunctionSpecDescription
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?

On this page