KubeMQ
Client SDKsGoReference

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:

NameTypeRequiredDescription
ctxcontext.ContextYesSend deadline
msg*QueueMessageYesChannel, 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:

NameTypeRequiredDescription
ctxcontext.ContextYesPoll deadline
req*PollRequestYesChannel, 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: *QueueUpstreamHandleSend, 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 error

Background error channel for stream issues; consume alongside polling loops.

QueueMessage policy

Policy fieldPurpose
ExpirationSecondsDrop undelivered messages after TTL
DelaySecondsDefer initial delivery
MaxReceiveCount / MaxReceiveQueueDead-letter and retry limits

Quick Usage

queues.go
msg := client.NewQueueMessage()
msg.SetChannel("jobs")
msg.SetBody([]byte("work"))

res, err := client.SendQueueMessage(ctx, msg)
if err != nil {
    return err
}
_ = res.MessageID

See Also

Was this page helpful?

On this page