Skip to content

Clients

The six SDKs, curl, queenctl and the broker embedded in a Rust program: how to install and connect each, and where they differ.

Updated View as Markdown

Queen speaks JSON over HTTP, on port 6632, so every client is a thin layer over the same routes, and curl is a complete client. What an SDK adds is the part that is tedious to write by hand: consume loops that ack and nack for you, lease renewal, push buffering, retries with failover across the nodes of a cluster, and builders for transactions, KV and timers. Pick the one for your language. The tutorials in examples/tutorials/ use these same packages, one program per language, and check themselves against a live broker.

Install and connect

npm install queen-mq
import { Queen } from 'queen-mq'

const queen = new Queen({ url: 'http://localhost:6632' })
// ...
await queen.close()

Node.js 24 or later. With urls: [...] instead of url, the client sends each consumer’s pops to the node that a consistent hash of its queue, partition and group picks (the default affinity strategy; round-robin and session are the others), and fails over to another node on a network error or a 5xx. README

pip install queen-mq
from queen import Queen

async def main():
    queen = Queen(url="http://localhost:6632")
    # ... every call is awaited
    await queen.close()

Python 3.9 or later. The client is asyncio from end to end, and url is a keyword argument. README

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

client, err := queen.New("http://localhost:6632")
if err != nil {
	return err
}
defer client.Close(context.Background())

Go 1.24 or later. queen.New also takes a slice of URLs or a ClientConfig. The repository lives at github.com/queen-mq/queen, but the module path still carries the organization it started in, so import it exactly as written above. README

cargo add queen-mq serde_json
cargo add tokio --features rt-multi-thread,macros
use queen_mq::{Config, Queen};

let queen = Queen::connect(Config::new("http://localhost:6632"))?;

Rust 1.86 or later, on tokio. connect checks the configuration and opens no socket; the first request does. README

The client is one header, clients/client-cpp/queen_client.hpp, used from a checkout of the repository: it includes the thread pool in clients/server/include/ by a relative path, and clients/server/vendor/ carries the nlohmann/json it needs. You install cpp-httplib and OpenSSL (brew install cpp-httplib openssl on macOS).

#include "queen_client.hpp"

queen::QueenClient client("http://localhost:6632");
// ...
client.close();
g++ -std=c++17 -O2 -pthread \
    -I clients/client-cpp -I clients/server/vendor \
    -I /opt/homebrew/include -I "$(brew --prefix openssl)/include" \
    app.cpp -o app \
    -L "$(brew --prefix openssl)/lib" -lssl -lcrypto -lpthread

C++17 or later. Link OpenSSL even over plain HTTP: the header turns on cpp-httplib’s TLS support itself. README

composer require queen-mq/php-client
use Queen\Queen;

$queen = new Queen('http://localhost:6632');

PHP 8.3 or later. The same package carries the Laravel integration: package discovery registers a queen queue connection, and your jobs stay ordinary Laravel jobs.

QUEUE_CONNECTION=queen
QUEEN_URL=http://127.0.0.1:6632

See Laravel. README

Every SDK is 2.0.0, like the broker, and all of them talk to a 2.0 broker. The Go client’s module path ends in /v2, as Go requires for a major version above 1.

curl is a client

Every SDK call is one of the broker’s routes, so anything that can send an HTTP request can push, pop, ack, commit a transaction, and read and write KV and timers, with nothing to install:

curl -s -X POST http://localhost:6632/api/v1/push -H 'content-type: application/json' \
  -d '{"items":[{"queue":"orders","partition":"customer-123","payload":{"orderId":8891}}]}'

The quickstart runs a whole step this way, with curl and jq, and the HTTP reference lists every route. queenctl, the operator CLI (status, tail, push, replay, lag, dead letters), ships inside the Docker image (docker exec queen queenctl status) and as release binaries for Linux, macOS and Windows. See queenctl.

What each SDK implements

The SDKs cover the same core (push, pop, consume, ack, transactions) and differ at the edges. This table lists every row where one of them stops short; where an SDK lacks a method, the route is still there over HTTP.

Parity here is a claim about the surface, not about shared code: each client is written natively in its own idiom. These sixteen rows are that surface, read out of clients/ on this tree. Six of them are full parity across all six clients; the other ten are where a client stops.

Capability JavaScript Python Go Rust PHP C++
Push, pop, ack, multi-partition claim yes yes yes yes yes yes
Client-side push buffering yes yes yes yes yes yes
Transactions through POST /api/v1/transaction yes yes yes yes yes yes
Dead-letter reader yes yes yes yes yes yes
HTTP 429 policy, separate from the retry counter yes yes yes yes yes yes
kv and timers riders on a transaction yes yes yes yes yes yes
Key/value: the seven operations yes yes yes yes yes five of seven
Timers: schedule, cancel, peek, list yes yes yes yes yes two of four
once, the idempotency marker in one call yes yes no no no no
Admin facade yes yes yes yes yes no
affinity load-balancing strategy yes yes yes yes yes no
Streams DSL over /streams/v1/* yes yes yes yes no no
Lease renewal inside the consume loop yes yes yes yes yes no
Per-message trace helper yes yes yes no yes no
Host override independent of the address dialled yes no no yes no no
pop reports a failure instead of an empty result no no yes yes yes no

The ten divergent rows, with their conditions:

  • Key/value. The seven operations are get, getMany, getPrefix, put, putIfAbsent, delete and incr. C++ wraps five of them: getMany and getPrefix have no method, and both remain reachable through POST /api/v1/kv on the shared HttpClient. getPrefix inside a transaction is refused by the broker in every client, which is its rule and not a client’s gap.
  • Timers. C++ has schedule and cancel only; peek and list are reads and stay on GET /api/v1/timers/{queue}[/{timerKey}]. There is no reschedule operation anywhere, because schedule is the upsert and status reports which happened; Python, Go, Rust and PHP still spell a reschedule alias over it, JavaScript and C++ do not.
  • once. JavaScript (kv.once, transaction().once) and Python (kv.once, transaction().once) fold putIfAbsent plus required into the question people actually ask, “did I win?”. The other four write the putIfAbsent themselves, which is the same wire and one line longer.
  • kv on a streams operator context. No client has it, so it is not a row. Inside a stream the state primitive is state_ops, which commits with the sink and the ack in the cycle’s own transaction; a key/value write from an operator would not, and that atomicity is the thing the stream already gives you for free.
  • Admin facade. C++ has no Admin class. Every management and observability route is still reachable, through the shared HttpClient returned by get_http_client().
  • affinity. The C++ LoadBalancer implements round-robin and "session" only, and load_balancing_strategy = "affinity" falls through to round-robin silently. The strategy matters because it keeps one consumer’s pops on one backend, which is what keeps two clients from contending on the same partition claim.
  • Streams. The Stream builder, the four window kinds, gates and the /streams/v1/* runtime are in the JavaScript, Python, Go and Rust clients. PHP and C++ have none of it.
  • Lease renewal. In C++, renew_lease() sets a flag the consumer loop does not act on: the worker carries an unimplemented placeholder where the timer would go. Call client.renew(message) yourself before the lease expires.
  • Per-message trace. Where it exists it is attached by the consume loop, so a message returned by a manual pop() never carries it. Rust records a trace through Admin::record_trace instead of hanging a method off the message. C++ cannot hang one off a JSON object at all, and its trace hook is an explicit no-op.
  • Host override. JavaScript (hostHeader) and Rust (host_header) advertise a request authority, and a TLS SNI, independent of the address the socket dials. That is what addresses a named tenant cluster behind a shared proxy endpoint. The other four clients send the dialled address.
  • pop failures. Go, Rust and PHP surface a 4xx, an exhausted 429 budget, a terminal 403 or a network fault to the caller, so an empty result means an empty queue. JavaScript, Python and C++ log the failure and return an empty result, which does not distinguish an empty queue from revoked credentials.

The broker inside a Rust program

A Rust program can run the broker in its own process, as the library queen::Broker. Every call goes through the same handler functions the HTTP router dispatches to, so behaviour and defaults are the broker’s, and the data lives in a directory on local disk that you name.

2.0 is not on crates.io yet (the queen-engine published there is 1.3, which ran on PostgreSQL), so depend on the repository, pinned to the 2.0.0 commit:

[dependencies]
queen-engine = { git = "https://github.com/queen-mq/queen", rev = "86562c25d", default-features = false }
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
serde_json = "1"

The crate imports as queen and needs Rust 1.88. default-features = false leaves out what only the standalone binary needs: the HTTP listener, the dashboard, the embedded proxy and the Kafka facade. Start it on a data directory:

server/tests/embedded_raft_smoke.rsrust
let broker = Broker::start(BrokerConfig::new().raft(&dir))
    .await
    .expect("broker start");

Push, with a transactionId on the first item:

server/tests/embedded_raft_smoke.rsrust
let first = broker
    .push(vec![
        qp::PushItem::new(q.clone(), serde_json::json!({"n": 1}))
            .transaction_id(txn_id.clone()),
        qp::PushItem::new(q.clone(), serde_json::json!({"n": 2})),
        qp::PushItem::new(q.clone(), serde_json::json!({"n": 3})),
    ])
    .await
    .expect("push");

Ack one message and push the next in one transaction:

server/tests/embedded_raft_smoke.rsrust
let txn = broker
    .transaction(
        &qp::TransactionRequest::new(vec![
            qp::TxnOperation::Ack(qp::TxnAckOperation {
                transaction_id: m0.transaction_id.clone(),
                partition_id: m0.partition_id.clone(),
                status: qp::AckStatus::Completed,
                consumer_group: Some(group.to_string()),
                lease_id: Some(m0.lease_id.clone()),
                error: None,
            }),
            qp::TxnOperation::Push {
                items: vec![qp::TxnPushItem::new(
                    q2.clone(),
                    serde_json::json!({"stage": 2}),
                )],
            },
        ])
        .with_required_leases([m0.lease_id.clone()]),
    )
    .await
    .expect("transaction");
assert!(txn.success, "transaction must commit: {txn:?}");

The embedded broker is one node and one tenant. It exposes push, pop, ack, lease renewal, transactions, queue configuration, the dead-letter queue and ephemeral queues. KV and timers are reachable only inside a transaction, so there is no standalone KV read or timer listing, and consumer-group administration, queue listings, traces and streams are not exposed. A panic in one of the broker’s core threads (planner, log writer, apply, checkpoint, raft) aborts the process, and the next start replays from the data directory; any other panic, your own included, unwinds as usual. Never open one data directory from two brokers at once.

Next

Navigation

Type to search…

↑↓ navigate↵ selectEsc close