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-mqimport { 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-mqfrom 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/v2import 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,macrosuse 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 -lpthreadC++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-clientuse 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:6632Every 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,deleteandincr. C++ wraps five of them:getManyandgetPrefixhave no method, and both remain reachable throughPOST /api/v1/kvon the sharedHttpClient.getPrefixinside a transaction is refused by the broker in every client, which is its rule and not a client’s gap. - Timers. C++ has
scheduleandcancelonly;peekandlistare reads and stay onGET /api/v1/timers/{queue}[/{timerKey}]. There is no reschedule operation anywhere, becausescheduleis the upsert andstatusreports which happened; Python, Go, Rust and PHP still spell areschedulealias over it, JavaScript and C++ do not. once. JavaScript (kv.once,transaction().once) and Python (kv.once,transaction().once) foldputIfAbsentplusrequiredinto the question people actually ask, “did I win?”. The other four write theputIfAbsentthemselves, which is the same wire and one line longer.kvon a streams operator context. No client has it, so it is not a row. Inside a stream the state primitive isstate_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.Adminfacade. C++ has noAdminclass. Every management and observability route is still reachable, through the sharedHttpClientreturned byget_http_client().affinity. The C++LoadBalancerimplements round-robin and"session"only, andload_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
Streambuilder, 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. Callclient.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 manualpop()never carries it. Rust records a trace throughAdmin::record_traceinstead 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. Hostoverride. 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.popfailures. 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:
let broker = Broker::start(BrokerConfig::new().raft(&dir))
.await
.expect("broker start");Push, with a transactionId on the first item:
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:
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.