The Composer package is smartpricing/queen-mq, PSR-4 namespace Queen\. It requires PHP 8.3 or
newer plus Guzzle 7 and ramsey/uuid. The client is synchronous: there is no promise or
async variant. Laravel integration (service provider, facade, artisan command) ships in the same
package and is optional. Version 1.0.0, aligned with the broker.
use Queen\Queen;
$queen = new Queen('http://localhost:6632');Constructor
new Queen(string|array $config = [])A string is a single URL. An array with a numeric key 0 is a list of URLs. Anything else is a
config array merged over Defaults::CLIENT_DEFAULTS; without url or urls it throws
InvalidArgumentException('Must provide urls or url in configuration').
| Key | Default | Effect |
|---|---|---|
url |
none | Single broker URL. |
urls |
none | Multiple brokers. More than one entry activates the load balancer. |
timeoutMillis |
30000 |
Per-request deadline. |
retryAttempts |
3 |
Attempts for 5xx and network failures. |
retryDelayMillis |
1000 |
First retry delay; doubles per attempt. |
loadBalancingStrategy |
'affinity' |
'affinity', 'round-robin' or 'session'. |
affinityHashRing |
128 |
Virtual nodes per backend on the hash ring. |
enableFailover |
true |
On a 5xx or network error, try the next backend. |
healthRetryAfterMillis |
5000 |
How long a backend stays marked unhealthy. |
bearerToken |
null |
Sent as Authorization: Bearer <token>. |
headers |
[] |
Extra headers on every request. |
retry429 |
[] |
HTTP 429 backoff policy. |
The load-balancing keys are inert with a single URL: the LoadBalancer is only constructed when
urls holds more than one entry.
retry429
Separate from retryAttempts, which covers 5xx and network faults. Constants live on
Queen\Http\Retry429Policy.
| Key | Default | Effect |
|---|---|---|
maxAttempts |
10 (DEFAULT_MAX_ATTEMPTS) for ordinary requests; UNBOUNDED (0) for a long-poll pop |
Total attempts including the first. Setting it applies the same bound to both kinds. |
baseMs |
500 (DEFAULT_BASE_MILLIS) |
First delay absent a Retry-After header. Doubles per attempt. |
capMs |
30000 (DEFAULT_CAP_MILLIS) |
Ceiling for the exponential delay. |
A Retry-After header (seconds) wins over the computed delay. A 429 is retried against the same
backend and never triggers failover.
Queen
| Method | Signature |
|---|---|
queue |
queue(?string $name = null): QueueBuilder |
admin |
admin(): Admin. Lazily created, one per client |
transaction |
transaction(): TransactionBuilder |
ack |
ack(array|string $message, bool|string $status = true, array $context = []): array |
renew |
renew(string|array $messageOrLeaseId): array |
flushAllBuffers |
flushAllBuffers(): void |
getBufferStats |
getBufferStats(): array |
deleteConsumerGroup |
deleteConsumerGroup(string $consumerGroup, bool $deleteMetadata = true): mixed |
updateConsumerGroupTimestamp |
updateConsumerGroupTimestamp(string $consumerGroup, string $timestamp): mixed |
close |
close(): void. Flushes buffers (best effort) and cleans up |
ack
$message is a message array, a transaction-id string, or an array of messages. $status is
true → completed, false → failed, or a status string passed through verbatim. $context
accepts group and error.
partitionId is mandatory on every message array; without it the builder throws
InvalidArgumentException. A leaseId is forwarded when present. Per-message statuses inside a
batch work by tagging items with _status and _error.
One message posts to /api/v1/ack, more than one to /api/v1/ack/batch. Both respond with one
result per ack in request order.
renew
Accepts a lease-id string, a message array, or an array of either. Posts to
/api/v1/lease/{leaseId}/extend once per distinct lease id.
QueueBuilder
$queen->queue($name). Configuration methods return static; terminal methods are push(),
pop(), consume(), getConsumer(), create(), delete(), dlq() and flushBuffer().
Addressing
| Method | Default | Notes |
|---|---|---|
partition(string $name) |
'Default' |
The ordered lane. Anything else switches pop to the partition-scoped route. |
namespace(string $name) |
null |
Grouping label at create time; also a pop filter. |
task(string $name) |
null |
Second grouping label; same dual role. |
group(string $name) |
null |
Consumer group. Absent, the broker uses __QUEUE_MODE__. |
Queue lifecycle
| Method | Wire call | Returns |
|---|---|---|
config(array $options) |
none | static. Merged over QUEUE_DEFAULTS. |
create() |
POST /api/v1/configure |
OperationBuilder |
delete() |
DELETE /api/v1/resources/queues/:queue |
OperationBuilder; throws RuntimeException with no queue name |
OperationBuilder exposes onSuccess(Closure), onError(Closure) and execute(): mixed.
Defaults::QUEUE_DEFAULTS:
| Key | Default | Unit |
|---|---|---|
leaseTime |
300 |
seconds |
retryLimit |
3 |
attempts |
delayedProcessing |
0 |
seconds |
windowBuffer |
0 |
seconds |
retentionSeconds |
0 |
seconds |
completedRetentionSeconds |
0 |
seconds |
encryptionEnabled |
false |
none |
Extra keys you pass to config() are sent verbatim, so any option /configure accepts is
reachable. /configure is a full replace: keys absent from the body reset to the broker’s
defaults, so send the complete intended configuration every time. The option list is in
Reference.
Producing
| Method | Notes |
|---|---|
buffer(array $options) |
['messageCount' => 100, 'timeMillis' => 1000]. Turns push() into an enqueue. |
push(array $payload): PushBuilder |
One item array or a list of them. Throws RuntimeException with no queue name. |
Each item may carry data, payload, or be the payload itself. The client mints a UUIDv7
transactionId via Queen\Support\Uuid::v7() when the item has none.
PushBuilder exposes onSuccess(Closure), onError(Closure), onDuplicate(Closure) and
execute(): mixed. execute() returns the broker’s per-item array for a direct push; a buffered
push returns immediately without sending.
$queen->queue('orders')
->partition('customer-42')
->push([['data' => ['orderId' => 8891]]])
->execute();Consuming
| Method | pop() default |
consume() default |
Notes |
|---|---|---|---|
batch(int $size) |
1 |
1 |
With partitions($n), the total budget across all claimed partitions. |
partitions(int $n) |
1 |
1 |
Claim up to N partitions per call; all share one leaseId. Sent only when > 1. |
wait(bool $enabled) |
false |
true |
Long-poll. |
timeoutMillis(int $millis) |
30000 |
30000 |
Server-side deadline; the HTTP call adds 5 s of slack. |
autoAck(bool $enabled) |
false |
true |
The two meanings differ. See the warning. |
subscriptionMode(string $mode) |
null |
null |
all | new. |
subscriptionFrom(string $from) |
null |
null |
'now' or an ISO timestamp. |
concurrency(int $count) |
n/a | 1 |
consume() only. |
limit(int $count) |
n/a | null |
consume() only. |
idleMillis(int $millis) |
n/a | null |
consume() only. |
renewLease(bool $enabled, ?int $intervalMillis = null) |
n/a | false |
consume() only. |
each() |
n/a | batch mode | consume() only. One message per handler call. |
pop(): array
Returns a list of message arrays, [] when there is nothing to return. 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(Closure $handler): ConsumeBuilder
ConsumeBuilder exposes onSuccess(Closure), onError(Closure) and execute(): void.
getConsumer(): HighLevelConsumer
A pull-loop consumer modelled on php-rdkafka’s KafkaConsumer, for code that wants to own its
own loop instead of handing over a closure.
| Method | Notes |
|---|---|
subscribe(): void |
Resolves the pop path and query from the builder’s options. Must be called before consuming. |
consume(int $timeoutMs = 1000): ?array |
One message or null. Forces batch=1, wait=true and timeout=$timeoutMs. |
consumeBatch(int $timeoutMs = 1000, int $maxMessages = 10): array |
Up to $maxMessages; [] when none. |
ack(array $message, bool $success = true): array |
|
nack(array $message): array |
|
renewLease(array|string $messageOrLeaseId): array |
|
isClosed(): bool |
|
close(): void |
consume() and consumeBatch() call pcntl_signal_dispatch() when the extension is available,
so a SIGTERM handler you installed runs between polls. Long-poll timeouts and connection
refusals are absorbed into null / []; every other exception propagates.
$message['trace']
Messages returned by HighLevelConsumer carry a trace closure:
$message['trace'](['traceName' => 'tenant-acme', 'eventType' => 'info', 'data' => ['step' => 1]]);data is required; traceName accepts a string or an array of strings; eventType defaults to
'info'. It never throws: failures return ['success' => false, 'error' => ...].
Buffering
flushBuffer(): void flushes this builder’s queue/partition buffer; throws RuntimeException
with no queue name. Buffers are keyed on "<queue>/<partition>", so builders addressing the same
pair share one. A buffer flushes at messageCount messages or timeMillis after its first
message.
Dead-letter queue
dlq(?string $consumerGroup = null): DLQBuilder. Throws RuntimeException with no queue name.
| Method | Notes |
|---|---|
limit(int $count) / offset(int $count) |
Paging. |
from(string $timestamp) / to(string $timestamp) |
Time filters. |
get(): array |
['messages' => [...], 'total' => n] |
Read-only. Replay lives on the admin façade: retryMessage(string $partitionId, string $transactionId).
TransactionBuilder
$queen->transaction(). Pushes and acks in one PostgreSQL transaction, all-or-nothing.
| Method | Notes |
|---|---|
ack(array|object $messages, string $status = 'completed', array $context = []): static |
Requires transactionId and partitionId; a leaseId is collected into requiredLeases. |
queue(string $queueName): TransactionQueueBuilder |
Sub-builder with partition(string) and push(array): TransactionBuilder. |
addPushOperation(string $queueName, ?string $partition, array $items): void |
Lower-level entry used by the sub-builder. |
commit(): array |
POST /api/v1/transaction. |
$queen->transaction()
->ack($message)
->queue('orders.enriched')->push([['data' => ['id' => 1]]])
->commit();Admin
$queen->admin(). Every method returns the broker’s decoded body. Methods taking array $params
turn it into a query string.
Resources
| Method | Route |
|---|---|
getOverview() |
GET /api/v1/resources/overview |
getNamespaces() |
GET /api/v1/resources/namespaces |
getTasks() |
GET /api/v1/resources/tasks |
listQueues(array $params = []) |
GET /api/v1/resources/queues |
getQueue(string $name) |
GET /api/v1/resources/queues/:name |
Messages and traces
| Method | Route |
|---|---|
listMessages(array $params = []) |
GET /api/v1/messages |
getMessage(string $partitionId, string $transactionId) |
GET /api/v1/messages/:partitionId/:transactionId |
deleteMessage(string $partitionId, string $transactionId) |
DELETE /api/v1/messages/:partitionId/:transactionId |
retryMessage(string $partitionId, string $transactionId) |
POST /api/v1/messages/:partitionId/:transactionId/retry |
getTraceNames(array $params = []) |
GET /api/v1/traces/names |
getTracesByName(string $traceName, array $params = []) |
GET /api/v1/traces/by-name/:traceName |
getTracesForMessage(string $partitionId, string $transactionId) |
GET /api/v1/traces/:partitionId/:transactionId |
Status and analytics
| Method | Route |
|---|---|
getStatus(array $params = []) |
GET /api/v1/status |
getQueueStats(array $params = []) |
GET /api/v1/status/queues |
getQueueDetail(string $name, array $params = []) |
GET /api/v1/status/queues/:name |
getAnalytics(array $params = []) |
GET /api/v1/status/analytics |
getSystemMetrics(array $params = []) |
GET /api/v1/analytics/system-metrics |
getWorkerMetrics(array $params = []) |
GET /api/v1/analytics/worker-metrics |
getPostgresStats() |
GET /api/v1/analytics/postgres-stats |
Consumer groups
| Method | Route |
|---|---|
listConsumerGroups() |
GET /api/v1/consumer-groups |
getConsumerGroup(string $name) |
GET /api/v1/consumer-groups/:name |
getLaggingConsumers(int $minLagSeconds = 60) |
GET /api/v1/consumer-groups/lagging?minLagSeconds= |
deleteConsumerGroupForQueue(string $consumerGroup, string $queueName, bool $deleteMetadata = true) |
DELETE /api/v1/consumer-groups/:group/queues/:queue |
seekConsumerGroup(string $consumerGroup, string $queueName, array $options = []) |
POST /api/v1/consumer-groups/:group/queues/:queue/seek |
refreshConsumerStats() |
POST /api/v1/stats/refresh |
seekConsumerGroup posts $options verbatim: ['toEnd' => true] or ['timestamp' => '<RFC3339>'].
System
| Method | Route |
|---|---|
health() |
GET /health |
metrics() |
GET /metrics |
getMaintenanceMode() |
GET /api/v1/system/maintenance |
setMaintenanceMode(bool $enabled) |
POST /api/v1/system/maintenance |
getPopMaintenanceMode() |
GET /api/v1/system/maintenance/pop |
setPopMaintenanceMode(bool $enabled) |
POST /api/v1/system/maintenance/pop |
Errors
Non-2xx responses throw Queen\Exceptions\HttpException, with isClusterSuspended() and
isRateLimited() helpers. The machine-readable codes are constants on
Queen\Exceptions\ErrorCode:
| Constant | Value | Status |
|---|---|---|
RATE_LIMITED |
rate_limited |
429 |
QUOTA_EXCEEDED |
quota_exceeded |
429 |
CLUSTER_SUSPENDED |
cluster_suspended |
403 |
STORAGE_QUOTA_EXCEEDED |
storage_quota_exceeded |
403 |
FEATURE_GATED |
feature_gated |
403 |
FORBIDDEN |
forbidden |
403 |
A 204 response yields null, because the broker sends no body at all on 204.
Laravel integration
The package auto-registers Queen\Laravel\QueenServiceProvider and the Queen facade alias.
Publish the config with the queen-config tag; it lands at config/queen.php. The provider binds
Queen\Queen as a singleton and aliases it to queen.
Config keys and env vars
| Config key | Env var | Default |
|---|---|---|
url |
QUEEN_URL |
http://localhost:6632 |
urls |
QUEEN_URLS |
unset (comma-separated list, wins over url) |
bearer_token |
QUEEN_BEARER_TOKEN |
unset |
timeout |
QUEEN_TIMEOUT |
30000 |
retry_attempts |
QUEEN_RETRY_ATTEMPTS |
3 |
retry_delay |
QUEEN_RETRY_DELAY |
1000 |
load_balancing_strategy |
QUEEN_LB_STRATEGY |
affinity |
enable_failover |
QUEEN_ENABLE_FAILOVER |
true |
affinity_hash_ring |
QUEEN_AFFINITY_HASH_RING |
150 |
health_retry_after |
QUEEN_HEALTH_RETRY_AFTER |
30000 |
headers |
none | [] |
retry_429.maxAttempts |
QUEEN_RETRY_429_MAX_ATTEMPTS |
unset |
retry_429.baseMs |
QUEEN_RETRY_429_BASE_MS |
unset |
retry_429.capMs |
QUEEN_RETRY_429_CAP_MS |
unset |
Facade
use Queen\Laravel\QueenFacade as Queen;
Queen::queue('orders')->push([['data' => ['id' => 1]]])->execute();The facade proxies queue(), transaction(), admin(), ack(), renew(),
flushAllBuffers(), getBufferStats() and close().
Artisan command
Registered only when running in console.
php artisan queen:consume orders "App\\Handlers\\OrderHandler"| Argument / option | Notes |
|---|---|
queue |
Queue name to consume from. |
handler |
Fully qualified class name with a handle() method. |
--group= |
Consumer group. |
--batch=1 |
Messages per batch. |
--auto-ack |
Enable auto-acknowledgment. |
--subscription-mode= |
Subscription mode. |
--subscription-from= |
Subscription start point. |
--timeout=30000 |
Long-poll timeout in milliseconds. |
--idle-timeout= |
Stop after N milliseconds of inactivity. |
--limit= |
Stop after processing N messages. |
The command errors out when the handler class does not exist.
Not in this client
There is no streaming SDK for PHP. The Stream builder, windows, gates and the
/streams/v1/* runtime exist in the JavaScript, Python and Go clients only.