Queues
Durable queues, polling, and streaming — KubeMQ Go SDK reference.
Queues provide at-least-once delivery with explicit or automatic acknowledgement. The Go SDK supports unary send/receive, batch send, bidirectional upstream streaming for high throughput, and a downstream receiver for long-lived poll loops.
Simple Queues
SendQueueMessage
func (c *Client) SendQueueMessage(ctx context.Context, msg *QueueMessage) (*QueueSendResult, error)Sends a single message with optional policy (delay, expiration, dead-letter, max receive count).
Parameters:
| Name | Type | Required | Description |
|---|---|---|---|
ctx | context.Context | Yes | Send deadline |
msg | *QueueMessage | Yes | Channel, body/metadata/tags, optional Policy |
Returns: *QueueSendResult — server message ID, timing fields, and error flags.
SendQueueMessages
func (c *Client) SendQueueMessages(ctx context.Context, msgs []*QueueMessage) ([]*QueueSendResult, error)Batch send with per-message validation.
PollQueue
func (c *Client) PollQueue(ctx context.Context, req *PollRequest) (*PollResponse, error)One-shot poll helper that forces AutoAck, creates a temporary downstream receiver, polls once, and closes the receiver.
Parameters:
| Name | Type | Required | Description |
|---|---|---|---|
ctx | context.Context | Yes | Poll deadline |
req | *PollRequest | Yes | Channel, max items, wait time, visibility, auto-ack |
Returns: *PollResponse — transaction ID and []*QueueDownstreamMessage for manual ack/nack when not auto-acking via full receiver.
AckAllQueueMessages
func (c *Client) AckAllQueueMessages(ctx context.Context, req *AckAllQueueMessagesRequest) (*AckAllQueueMessagesResponse, error)Acknowledges all visible messages on a channel (admin-style operation).
Queue Streaming
QueueUpstream
func (c *Client) QueueUpstream(ctx context.Context) (*QueueUpstreamHandle, error)Opens a bidirectional stream for high-throughput publishing. Use Send with a request ID and batch of *QueueMessage, read Results for per-batch outcomes, and Close when finished.
Returns: *QueueUpstreamHandle — Send, Results channel, Done, Close.
NewQueueDownstreamReceiver
func (c *Client) NewQueueDownstreamReceiver(ctx context.Context) (*QueueDownstreamReceiver, error)Creates a persistent downstream client for repeated Poll calls over a single gRPC stream with automatic reconnection.
Poll (receiver)
func (r *QueueDownstreamReceiver) Poll(ctx context.Context, req *PollRequest) (*PollResponse, error)Polls messages according to PollRequest (auto-ack vs manual settle). Use PollResponse helpers to ack/nack/requeue batches when AutoAck is false.
Errors
func (r *QueueDownstreamReceiver) Errors() <-chan errorBackground error channel for stream issues; consume alongside polling loops.
QueueMessage policy
| Policy field | Purpose |
|---|---|
ExpirationSeconds | Drop undelivered messages after TTL |
DelaySeconds | Defer initial delivery |
MaxReceiveCount / MaxReceiveQueue | Dead-letter and retry limits |
Quick Usage
msg := client.NewQueueMessage()
msg.SetChannel("jobs")
msg.SetBody([]byte("work"))
res, err := client.SendQueueMessage(ctx, msg)
if err != nil {
return err
}
_ = res.MessageIDSee Also
Was this page helpful?