Skip to content

Go client

Complete surface of the Go client package: ClientConfig with defaults, every builder method and its signature, the Admin methods that reach a route that exists, and the streaming SDK.

Updated View as Markdown

Import path github.com/smartpricing/queen/clients/client-go, package name queen. Requires Go 1.24 or newer. Every network method takes a context.Context as its first argument. Version 1.0.0, aligned with the broker.

import queen "github.com/smartpricing/queen/clients/client-go"

The streaming SDK is a separate package under .../client-go/streams.

Constructor

func New(config interface{}) (*Queen, error)

config accepts a string (single URL), a []string (multiple URLs), a ClientConfig, or a *ClientConfig. Anything else returns invalid config type: expected string, []string, or ClientConfig.

ClientConfig

Field Type Default applied Effect
URLs []string none Multiple brokers. More than one entry activates the load balancer.
URL string none Single broker, converted to URLs.
TimeoutMillis int 30000 Per-request deadline.
RetryAttempts int 3 Retries after the first attempt, so the default is 4 attempts in total. A negative value means exactly one attempt.
RetryDelayMillis int 1000 First retry delay; doubles per attempt.
LoadBalancingStrategy string "affinity" "affinity", "round-robin" or "session".
AffinityHashRing int 128 Virtual nodes per backend on the hash ring.
EnableFailover bool false in practice (see below) Whether a failing backend is marked unhealthy.
HealthRetryAfterMillis int 5000 How long a backend stays marked unhealthy.
BearerToken string "" Sent as Authorization: Bearer <token>.
Headers map[string]string nil Extra headers on every request.
MaxIdleConnsPerHost int 256 when <= 0 Go’s own default is 2, which cripples a single-host high-concurrency producer pool.
MaxConnsPerHost int 0 (unlimited) Hard cap on total connections per host.
Retry429 *Retry429Config nil HTTP 429 backoff policy.

The exported ClientDefaults, QueueDefaults, ConsumeDefaults, PopDefaults and BufferDefaults variables hold these values, so an application can read one instead of restating it. Unlike the other clients, the load balancer is constructed unconditionally, even for a single URL, so LoadBalancingStrategy and the health settings are always in play.

Retry429Config

Separate from RetryAttempts, which covers 5xx and network faults. A zero field means “use the kind-based default”.

Field Default Effect
MaxAttempts 10 for ordinary requests; unbounded for a long-poll pop Total attempts including the first. Setting it applies the same bound to both kinds.
BaseMs 500 First delay absent a Retry-After header. Doubles per attempt.
CapMs 30000 Ceiling for the exponential delay.

A Retry-After header (seconds) wins over the computed delay, with ±20% jitter either way. A 429 is retried against the same backend and never triggers failover.

*Queen

Method Signature
Queue Queue(name string) *QueueBuilder
Admin Admin() *Admin. Lazily created, one per client
Transaction Transaction() *TransactionBuilder
Ack Ack(ctx, messages interface{}, success bool, opts AckOptions) ([]AckResponse, error)
Renew Renew(ctx, messages interface{}) ([]RenewResponse, error)
FlushAllBuffers FlushAllBuffers(ctx) error
GetBufferStats GetBufferStats() BufferStats
DeleteConsumerGroup DeleteConsumerGroup(ctx, consumerGroup string, deleteMetadata bool) error
UpdateConsumerGroupTimestamp UpdateConsumerGroupTimestamp(ctx, consumerGroup string, timestamp time.Time) error
Close Close(ctx) error. Flushes buffers, then closes the HTTP client
GetHttpClient GetHttpClient() *HttpClient
GetBufferManager GetBufferManager() *BufferManager

Ack

messages accepts *Message, Message, []*Message or []Message; anything else returns an error. success maps to AckStatusCompleted ("completed") or AckStatusFailed ("failed"). AckOptions carries ConsumerGroup and Error.

One message uses POST /api/v1/ack, more than one uses POST /api/v1/ack/batch. Both return one AckResponse{Success, Error} per message in request order. ValidateMessages runs first, so a message missing TransactionID or PartitionID fails before any request is sent.

Renew

Accepts a lease-id string, []string, or any of the message forms. Lease ids are deduplicated first: with multi-partition pop every message in a batch shares one LeaseID and one extend call renews every claimed partition, so passing the whole slice issues one HTTP call. RenewResponse carries LeaseID, Success, NewExpiresAt time.Time and Error.

*QueueBuilder

client.Queue(name). Configuration methods return *QueueBuilder; terminal methods are Push, Pop, Consume, ConsumeBatch, Create, Delete, DLQ and FlushBuffer.

Addressing

Method Default Notes
Name() string none The queue name (used by the streaming SDK).
Partition(name) DefaultPartition ("Default") The ordered lane. Anything else switches pop to the partition-scoped route.
Namespace(name) "" Grouping label at create time; also a pop filter.
Task(name) "" Second grouping label; same dual role.
Group(name) "" Consumer group. Absent, the broker uses QueueModeConsumerGroup ("__QUEUE_MODE__").

Queue lifecycle

Method Wire call Returns
Config(config QueueConfig) none *QueueBuilder
Create() POST /api/v1/configure *OperationBuilder
Delete() DELETE /api/v1/resources/queues/:queue *OperationBuilder

*OperationBuilder has one method: Execute(ctx) (map[string]interface{}, error). Create validates the queue name first and adds "configured": true to the returned map for parity with the other clients.

QueueConfig holds the fields Create can send, with the value in QueueDefaults:

Field Wire key QueueDefaults Unit
LeaseTime leaseTime 300 seconds
RetryLimit retryLimit 3 attempts
DelayedProcessing delayedProcessing 0 seconds
WindowBuffer windowBuffer 0 seconds
RetentionSeconds retentionSeconds 0 seconds
CompletedRetentionSeconds completedRetentionSeconds 0 seconds
RetentionEnabled retentionEnabled false none
DeadLetterQueue deadLetterQueue false none
DlqAfterMaxRetries dlqAfterMaxRetries false none
EncryptionEnabled encryptionEnabled false none

Producing

Method Notes
Buffer(config BufferConfig) {MessageCount, TimeMillis}; BufferDefaults is 100 / 1000. Turns Push into an enqueue.
Push(payload interface{}) *PushBuilder payload is one value or a slice.

*PushBuilder:

Method Notes
TransactionID(id string) Explicit id. Applied to the first item only.
TraceID(id string) Trace id, forwarded only when it is a valid UUID.
Execute(ctx) ([]PushResponse, error) Sends POST /api/v1/push, or enqueues when Buffer is set.

PushResponse carries Status ("queued", "duplicate", "failed"), TransactionID and Error. The client mints a UUIDv7 transactionId for items that have none.

clients/client-go/tests/push_test.gogo
responses, err := client.Queue(queueName).Push(payload).Execute(ctx)
if err != nil {
	t.Fatalf("Failed to push message: %v", err)
}

Consuming

Method Pop default Consume default Notes
Batch(size) 1 1 With Partitions(n), this is the total budget across all claimed partitions.
Partitions(n) 1 1 Claim up to N partitions per call; all share one LeaseID. Sent only when > 1.
Wait(enabled) false true Long-poll. Tri-state internally: unset falls back to the per-method default.
TimeoutMillis(ms) 30000 30000 Server-side deadline. Pop adds 5 s of client slack when waiting.
AutoAck(enabled) not sent true The two meanings differ. See the warning.
SubscriptionMode(mode) "" "" SubscriptionModeAll, SubscriptionModeNew, SubscriptionModeNewOnly.
SubscriptionFrom(from) "" "" SubscriptionFromNow ("now") or an ISO timestamp.
Concurrency(count) n/a 1 Consume only. Worker goroutines.
Limit(count) n/a 0 (unlimited) Consume only.
IdleMillis(ms) n/a 0 (no timeout) Consume only.
RenewLease(enabled, intervalMillis) n/a false Consume only.
Each() n/a batch mode Consume only. One message per handler call.

Pop

func (qb *QueueBuilder) Pop(ctx context.Context) ([]*Message, error)

Returns queue name, namespace, or task is required when none of the three is set. Route selection: queue plus non-Default partition uses /api/v1/pop/queue/:queue/partition/:partition; queue alone uses /api/v1/pop/queue/:queue; namespace or task alone uses /api/v1/pop.

Consume and ConsumeBatch

func (qb *QueueBuilder) Consume(ctx, handler MessageHandler) *ConsumeBuilder
func (qb *QueueBuilder) ConsumeBatch(ctx, handler BatchMessageHandler) *ConsumeBuilder

MessageHandler is func(ctx context.Context, msg *Message) error; BatchMessageHandler is func(ctx context.Context, msgs []*Message) error. *ConsumeBuilder exposes Execute(ctx) error and Start(ctx) error (Start is an alias). Both block until every worker stops; cancel the context to stop them.

Worker loop behaviour:

  • A handler error with AutoAck on nacks the message and the worker continues.
  • A 429 that escapes the HTTP layer’s own retry backs off and continues.
  • A 403 stops the worker. Use (*HTTPError).IsClusterSuspended() to distinguish the terminal case.

(*Message).Trace

func (m *Message) Trace(ctx context.Context, config TraceConfig) (*TraceResponse, error)

TraceConfig carries TraceName string, TraceNames []string, EventType string and Data map[string]interface{}. Posts to /api/v1/traces.

Buffering

FlushBuffer(ctx) error flushes this builder’s queue/partition buffer. *BufferManager exposes Add, Flush(ctx, key), FlushAll(ctx), GetStats() BufferStats, GetBuffer(key) and Clear(). BufferStats carries ActiveBuffers, TotalBufferedMessages, OldestBufferAge and FlushesPerformed.

Dead-letter queue

DLQ(consumerGroup string) *DLQBuilder.

Method Notes
Limit(count) / Offset(count) Paging.
From(timestamp) / To(timestamp) Time filters.
Get(ctx) (*DLQResponse, error) DLQResponse{Messages []Message, Total int}

Read-only. Admin.RetryMessage exists but returns an error without making a request, so replay from Go means reading the row and pushing the payload again, or calling POST /api/v1/messages/{partitionId}/{transactionId}/retry directly.

*TransactionBuilder

client.Transaction(). Pushes and acks in one PostgreSQL transaction, all-or-nothing.

Method Notes
Ack(messages interface{}, status string, opts AckOptions) *TransactionBuilder Same message forms as Queen.Ack. A LeaseID is collected into requiredLeases.
Queue(name string) *TransactionQueueBuilder Sub-builder with Partition(name) and Push(payload) *TransactionBuilder. Pass a PushItem to control that item’s TransactionID, TraceID and Partition; pass anything else and it is treated as a bare payload with a minted id.
Commit(ctx) (*TransactionResponse, error) POST /api/v1/transaction. Check resp.Success as well as err.
resp, err := client.Transaction().
    Ack(msg, queen.AckStatusCompleted, queen.AckOptions{ConsumerGroup: "workers"}).
    Queue("orders.enriched").Push(payload).
    Commit(ctx)

*Admin

client.Admin(). Every method takes a context and returns map[string]interface{}, the broker’s JSON body verbatim, except where noted.

Resources

Method Route
GetOverview(ctx) GET /api/v1/resources/overview
GetNamespaces(ctx) GET /api/v1/resources/namespaces
GetTasks(ctx) GET /api/v1/resources/tasks
ListQueues(ctx, params ListQueuesParams) GET /api/v1/resources/queues
GetQueue(ctx, name) GET /api/v1/resources/queues/:name
GetPartitions(ctx, queueName) GET /api/v1/status/queues/:queue?includePartitions=true

ListQueuesParams carries Namespace, Task, Limit, Offset. GetPartitions requires a queue name and returns the queue-detail payload: the broker bundles partition data there rather than exposing a partitions collection.

Messages, DLQ and traces

Method Route
ListMessages(ctx, params ListMessagesParams) GET /api/v1/messages
GetMessage(ctx, partitionID, transactionID) GET /api/v1/messages/:partitionId/:transactionId
DeleteMessage(ctx, partitionID, transactionID) DELETE /api/v1/messages/:partitionId/:transactionId
ListDLQ(ctx, params ListDLQParams) GET /api/v1/dlq
GetTraceNames(ctx, limit, offset) GET /api/v1/traces/names
GetTracesByName(ctx, traceName, limit, offset) GET /api/v1/traces/by-name/:traceName
GetTracesForMessage(ctx, partitionID, transactionID) GET /api/v1/traces/:partitionId/:transactionId

ListMessagesParams carries Queue, Partition, Status, ConsumerGroup, Limit, Offset. ListDLQParams carries Queue, ConsumerGroup, Partition, From, To, Limit, Offset.

Status and analytics

Method Route
GetStatus(ctx, params GetStatusParams) GET /api/v1/status
GetQueueStats(ctx, namespace, task) GET /api/v1/status/queues
GetQueueDetail(ctx, name, includePartitions bool) GET /api/v1/status/queues/:name
GetAnalytics(ctx, from, to) GET /api/v1/status/analytics
GetQueueLagAnalytics(ctx, from, to) GET /api/v1/analytics/queue-lag
GetQueueOpsAnalytics(ctx, from, to) GET /api/v1/analytics/queue-ops
GetQueueParkedReplicas(ctx, from, to) GET /api/v1/analytics/queue-parked-replicas
GetRetentionAnalytics(ctx, from, to) GET /api/v1/analytics/retention
GetSystemMetrics(ctx, from, to) GET /api/v1/analytics/system-metrics
GetWorkerMetrics(ctx, from, to) GET /api/v1/analytics/worker-metrics
GetPostgresStats(ctx) GET /api/v1/analytics/postgres-stats

GetStatusParams carries Queue, Namespace, Task.

Consumer groups

Method Route
ListConsumerGroups(ctx) GET /api/v1/consumer-groups
GetConsumerGroup(ctx, name) GET /api/v1/consumer-groups/:name
GetLaggingConsumers(ctx, minLagSeconds) GET /api/v1/consumer-groups/lagging?minLagSeconds=
DeleteConsumerGroupForQueue(ctx, group, queue, deleteMetadata) DELETE /api/v1/consumer-groups/:group/queues/:queue
SeekConsumerGroup(ctx, group, queue, opts SeekConsumerGroupOptions) POST /api/v1/consumer-groups/:group/queues/:queue/seek
SeekConsumerGroupPartition(ctx, group, queue, partition, opts) POST /api/v1/consumer-groups/:group/queues/:queue/partitions/:partition/seek
RefreshConsumerStats(ctx) POST /api/v1/stats/refresh

SeekConsumerGroupOptions carries Timestamp string (RFC3339) and ToEnd bool, which are mutually exclusive. The wire body is {"toEnd": true} or {"timestamp": "…"}.

System

Method Route Returns
Health(ctx) GET /health map
Metrics(ctx) GET /metrics (string, error): the raw body when it is not JSON, "" when it parsed as JSON
PrometheusMetrics(ctx) GET /metrics/prometheus (string, error)
GetMaintenanceMode(ctx) GET /api/v1/system/maintenance map
SetMaintenanceMode(ctx, enabled) POST /api/v1/system/maintenance map
GetPopMaintenanceMode(ctx) GET /api/v1/system/maintenance/pop map
SetPopMaintenanceMode(ctx, enabled) POST /api/v1/system/maintenance/pop map

Errors

Non-2xx responses come back as *HTTPError:

var he *queen.HTTPError
if errors.As(err, &he) && he.StatusCode == 429 { /* ... */ }
Field Notes
StatusCode int HTTP status.
Body string Raw body.
Code string The body’s code field. Proxy contract: rate_limited, quota_exceeded on 429; cluster_suspended, storage_quota_exceeded, feature_gated, forbidden on 403. Empty when the body has none.
RetryAfterSeconds *float64 Parsed from the Retry-After header on a 429; nil when absent or non-numeric.

IsClusterSuspended() reports the terminal 403 that no amount of retrying resolves. Error() renders as HTTP <status> [<code>]: <body>.

A 204 response yields a nil map, because the broker sends no body at all on 204.

*HttpClient

Reachable via GetHttpClient() when you need a route the SDK does not wrap.

Method Notes
Get(ctx, path string, timeoutMs int, affinityKey string, opts ...RequestOption) timeoutMs of 0 uses the client default.
Post(ctx, path, body interface{}, opts ...RequestOption)
PostWithAffinity(ctx, path, body, affinityKey, opts ...RequestOption)
Delete(ctx, path, opts ...RequestOption)
GetLoadBalancer()
Close()

All four return (map[string]interface{}, error); a top-level JSON array is wrapped as {"data": [...]}. WithLongPollRetry() is the only RequestOption: it marks the request as a long-poll pop so a 429 backs off indefinitely instead of using the bounded budget.

Logging

Set by the QUEEN_CLIENT_LOG environment variable, read once at package init:

Value Level
unset LogLevelNone
debug, true, 1 LogLevelDebug
info LogLevelInfo
warn, warning LogLevelWarn
error LogLevelError
anything else LogLevelNone

SetLogLevel(level LogLevel) and GetLogLevel() LogLevel change it at runtime.

Helpers

Function Notes
GenerateUUID() string The id used for transactionId.
GenerateUUIDv4() string, GenerateUUIDv7() (string, error) Explicit versions.
ParseUUID(s) (uuid.UUID, error), MustParseUUID(s) uuid.UUID
IsValidUUID(s) bool
ValidateURL(raw) (string, error), ValidateURLs(urls) ([]string, error)
ValidateQueueName, ValidatePartitionName, ValidateConsumerGroup (string, error): normalise and reject.
ValidateBatchSize, ValidateConcurrency, ValidateTimeout error.
ValidateMessage(*Message) error, ValidateMessages([]*Message) error Run by Ack before any request.

Message

Field JSON Notes
TransactionID transactionId
PartitionID partitionId Required for every ack.
LeaseID leaseId Empty when the pop used server-side auto-ack.
Queue, Partition queue, partition
Data data map[string]interface{}.
CreatedAt createdAt
ErrorMessage errorMessage
RetryCount retryCount
ProducerSub producerSub The authenticated producer identity, stamped by the broker from the JWT sub at push time. Present only when JWT auth is on. Clients cannot set it.

Streaming SDK

import "github.com/smartpricing/queen/clients/client-go/streams"

s := streams.From(queen.Queue("orders").AsStreamSource())

AsStreamSource() on a *QueueBuilder adapts it to runtime.Source. Every combinator returns a new *Stream.

Group Methods
Stateless Map(MapFn), Filter(FilterFn), FlatMap(FlatMapFn)
Keying KeyBy(KeyFn)
Windows WindowTumbling(seconds float64, opts ...WindowOption), WindowSliding(size, slide float64, opts ...), WindowSession(gap float64, opts ...), WindowCron(every string, opts ...)
Reduce Reduce(ReduceFn, initial interface{}), Aggregate(map[string]ExtractorFn, fieldOrder ...string)
Gate Gate(GateFn)
Sink To(Sink), Foreach(ForeachFn)
Terminal Run(ctx, opts RunOptions) (*Runner, error), Compile() (*runtime.CompiledStream, error)

Window options: WithGracePeriod, WithIdleFlushMs, WithEventTime, WithAllowedLateness, WithOnLate. streams/helpers provides TokenBucketGate(TokenBucketGateOptions) and SlidingWindowGate(SlidingWindowGateOptions), both returning a GateFn.

RunOptions

Field Default
QueryID required, the durable identity
URL required, the broker base URL for the /streams/v1/* calls
BearerToken ""
BatchSize 200
MaxPartitions 4
MaxWaitMillis 1000
SubscriptionMode ""
SubscriptionFrom ""
Reset false (wipe state on a config-hash mismatch)
ConsumerGroup streams.<QueryID>
Logger nil

*Runner exposes Metrics() Metrics with the counters CyclesTotal, FlushCyclesTotal, MessagesTotal, PushItemsTotal, StateOpsTotal, LateEventsTotal, ErrorsTotal, GateAllowsTotal, GateDenialsTotal and LastError.

The chain shape is fingerprinted into a config hash; re-deploying a different chain under the same QueryID is rejected at registration unless Reset is set. Only operator kinds and their structural config are hashed, not the bodies of your functions.

Same surface, other languages

Navigation

Type to search…

↑↓ navigate↵ selectEsc close