Search by

crazy-goat / rabbit-stream

crazy-goat

Package info

github.com/crazy-goat/rabbit-stream

pkg:composer/crazy-goat/rabbit-stream

Statistics

Installs: 4

Dependents: 0

Suggesters: 0

Stars: 0

Open Issues: 70

v1.4.0 2026-09-21 21:47 UTC

This package is auto-updated.

Last update: 2026-10-04 20:09:01 UTC


README

A PHP library implementing the RabbitMQ Streams Protocol client.

It provides low-level TCP communication with a RabbitMQ broker over the native Stream protocol (port 5552), including binary frame serialization/deserialization.

Requirements

  • PHP 8.1+, 64-bit build (stream offsets are uint64; see Requirements)
  • RabbitMQ with the rabbitmq_stream plugin enabled

Installation

composer require crazy-goat/rabbit-stream

Quick Start

Publishing

use CrazyGoat\RabbitStream\Client\Connection;

$connection = Connection::create(host: 'localhost', port: 5552);

$producer = $connection->createProducer('my-stream', name: 'my-producer');
$producer->send('hello world');
$producer->waitForConfirms(timeout: 5);
$producer->close();

$connection->close();

Message bodies are plain strings — Producer::send() and sendBatch() automatically wrap them in an AMQP 1.0 Data section on the wire, and the consumer returns them unwrapped (see Publishing).

TLS transport (encrypted connections)

By default connections are plaintext tcp:// (port 5552). Pass a CrazyGoat\RabbitStream\VO\TlsConfig to use the encrypted ssl:// transport (RabbitMQ stream listener on port 5551):

use CrazyGoat\RabbitStream\Client\Connection;
use CrazyGoat\RabbitStream\VO\TlsConfig;

$connection = Connection::create(
    host: 'localhost',
    port: 5551, // TLS stream port
    tls: new TlsConfig(
        cafile: '/etc/ssl/ca.pem', // optional: custom CA bundle
        // localCert: '/etc/ssl/client.crt',   // optional: client certificate
        // localPk: '/etc/ssl/client.key',     // (e.g. for EXTERNAL SASL)
    ),
);

Peer certificate and hostname verification are on by default (verify_peer/verify_peer_name). For a self-signed development broker you can disable them explicitly — do not do this in production:

$connection = Connection::create(
    host: 'localhost',
    port: 5551,
    tls: new TlsConfig(verifyPeer: false, verifyPeerName: false),
);

Consuming

use CrazyGoat\RabbitStream\Client\Connection;
use CrazyGoat\RabbitStream\VO\OffsetSpec;

$connection = Connection::create(host: 'localhost', port: 5552);

$consumer = $connection->createConsumer('my-stream', offset: OffsetSpec::first());
// read() returns [] when nothing arrives within the timeout, so this loop
// stops at the first quiet 5 seconds. That suits draining a backlog.
while ($messages = $consumer->read(timeout: 5)) {
    foreach ($messages as $msg) {
        echo $msg->getBody() . "\n";
    }
}
$consumer->close();

$connection->close();

To keep consuming a live stream that can have gaps, do not stop on an empty batch. Loop on a flag and treat [] as "nothing yet" (see examples/consumer.php for a full version with signal handling):

while ($running) {
    foreach ($consumer->read(timeout: 5) as $msg) {
        echo $msg->getBody() . "\n";
    }
}

Documentation

Full documentation lives in docs/: getting started, guides, API reference, protocol notes and examples. Translations are listed in docs/LANGUAGES.md (currently English only).

Usage

High-level API (Recommended)

use CrazyGoat\RabbitStream\Client\Connection;
use CrazyGoat\RabbitStream\Client\ConfirmationStatus;

// Connect (handshake and authentication handled automatically)
$connection = Connection::create(
    host: '127.0.0.1',
    user: 'guest',
    password: 'guest'
);

// Create a producer for 'my-stream'
$producer = $connection->createProducer(
    stream: 'my-stream',
    onConfirm: function (ConfirmationStatus $status): void {
        if ($status->isConfirmed()) {
            echo "Message {$status->getPublishingId()} confirmed\n";
        }
    }
);

// Send a message
$producer->send("Hello, RabbitMQ Stream!");

// Drive the loop to receive confirmations (optional, blocking)
$connection->readLoop(maxFrames: 1);

// Close producer and connection
$producer->close();
$connection->close();

Consuming with Message Decoding

use CrazyGoat\RabbitStream\Client\AmqpMessageDecoder;
use CrazyGoat\RabbitStream\Client\OsirisChunkParser;

// ... subscribe to stream and receive Deliver response

$chunk = $deliverResponse->getChunkBytes();
$entries = OsirisChunkParser::parse($chunk);

// Decode AMQP 1.0 messages into Message objects
$messages = AmqpMessageDecoder::decodeAll($entries);

foreach ($messages as $message) {
    echo "Offset: {$message->getOffset()}\n";
    echo "Body: {$message->getBody()}\n";
    echo "Content-Type: {$message->getContentType()}\n";
    echo "Message-ID: {$message->getMessageId()}\n";
}

Consumer with Auto-Commit

use CrazyGoat\RabbitStream\Client\Connection;
use CrazyGoat\RabbitStream\VO\OffsetSpec;

$connection = Connection::create(host: 'localhost', port: 5552);

// Named consumer with auto-commit every 100 messages
// The name is used to persist the offset on the server
$consumer = $connection->createConsumer(
    stream: 'my-stream',
    offset: OffsetSpec::first(),
    name: 'my-consumer-group',
    autoCommit: 100,
);

// Stops at the first empty batch (nothing within 5 seconds); see "Consuming".
while ($messages = $consumer->read(timeout: 5)) {
    foreach ($messages as $msg) {
        echo $msg->getBody() . "\n";
    }
}
$consumer->close(); // stores final offset automatically

// On next startup, resume from the stored offset. queryOffset() returns null
// when nothing has been stored yet (a normal first run, not an error).
$storedOffset = $connection->queryOffset('my-consumer-group', 'my-stream');
$consumer = $connection->createConsumer(
    stream: 'my-stream',
    offset: $storedOffset === null ? OffsetSpec::first() : OffsetSpec::offset($storedOffset),
    name: 'my-consumer-group',
    autoCommit: 100,
);

$connection->close();

Note: autoCommit triggers storeOffset every N messages. The offset is also stored on close(). A named consumer is required for offset persistence — unnamed consumers cannot use storeOffset or queryOffset. queryOffset() returns null (rather than throwing) when no offset has been stored for the name/stream pair.

See examples/consumer_auto_commit.php for a full working example.

Low-level Connection API

StreamConnection exposes the raw protocol frames (requests and responses) for cases the high-level Connection does not cover. It is documented in the StreamConnection API reference; for most applications use the high-level API above.

Protocol Implementation Status

Protocol reference: https://github.com/rabbitmq/rabbitmq-server/blob/main/deps/rabbitmq_stream/docs/PROTOCOL.adoc

Connection & Authentication

Command Key Request Response
PeerProperties 0x0011 ✅ ✅
SaslHandshake 0x0012 ✅ ✅
SaslAuthenticate 0x0013 ✅ ✅
Tune 0x0014 ✅ ✅
Open 0x0015 ✅ ✅

Publishing

Command Key Request Response
DeclarePublisher 0x0001 ✅ ✅
Publish 0x0002 ✅ —
PublishConfirm 0x0003 — ✅
PublishError 0x0004 — ✅
QueryPublisherSequence 0x0005 ✅ ✅
DeletePublisher 0x0006 ✅ ✅

Consuming

Command Key Request Response
Subscribe 0x0007 ✅ ✅
Deliver 0x0008 — ✅
Credit 0x0009 ✅ ✅
StoreOffset 0x000a ✅ —
QueryOffset 0x000b ✅ ✅
Unsubscribe 0x000c ✅ ✅
ConsumerUpdate 0x001a ✅ ✅

Stream Management

Command Key Request Response
Create 0x000d ✅ ✅
Delete 0x000e ✅ ✅
Metadata 0x000f ✅ ✅
MetadataUpdate 0x0010 — ✅
CreateSuperStream 0x001d ✅ ✅
DeleteSuperStream 0x001e ✅ ✅
StreamStats 0x001c ✅ ✅

Routing (Super Streams)

Command Key Request Response
Route 0x0018 ✅ ✅
Partitions 0x0019 ✅ ✅

Connection Management

Command Key Request Response
Close 0x0016 ✅ ✅
Heartbeat 0x0017 ✅ —
ExchangeCommandVersions 0x001b ✅ ✅
ResolveOffsetSpec 0x001f ✅ ✅

Legend: ✅ implemented, ❌ not implemented, — not applicable (one-direction command)