Search by

roadrunner / centrifugo

roxblnfk

Centrifugo bridge for RoadRunner: handle Centrifugo proxy events in PHP workers and call the Centrifugo server API over RPC

2.5.0 2026-10-10 20:27 UTC

README

RoadRunner

Centrifugo bridge for RoadRunner: proxy events and server API

Documentation Sponsor

Psalm Level Type Coverage Mutation testing badge


PHP bridge for the RoadRunner centrifuge plugin: handle Centrifugo proxy events (connect, subscribe, publish, RPC, …) in PHP workers and call the Centrifugo server API over RoadRunner RPC.

Get Started

Installation

composer require roadrunner/centrifugo

PHP Latest Version on Packagist License Total Downloads

You can use the convenient installer to download the latest available compatible version of RoadRunner assembly:

composer require roadrunner/cli --dev
vendor/bin/rr get

Configuration

Add the centrifuge section to your RoadRunner configuration (.rr.yaml):

rpc:
  listen: tcp://127.0.0.1:6001

server:
  command: "php app.php"
  relay: pipes

centrifuge:
  # RoadRunner listens here for proxy requests from Centrifugo
  proxy_address: "tcp://0.0.0.0:10001"
  # Centrifugo gRPC API address (used by the server API client)
  grpc_api_address: "tcp://127.0.0.1:10000"

and point the Centrifugo proxy endpoints to it:

{
  "admin": true,
  "api_key": "secret",
  "admin_password": "password",
  "admin_secret": "admin_secret",
  "allowed_origins": [
    "*"
  ],
  "token_hmac_secret_key": "test",
  "publish": true,
  "proxy_publish": true,
  "proxy_subscribe": true,
  "proxy_connect": true,
  "allow_subscribe_for_client": true,
  "proxy_connect_endpoint": "grpc://127.0.0.1:10001",
  "proxy_connect_timeout": "10s",
  "proxy_publish_endpoint": "grpc://127.0.0.1:10001",
  "proxy_publish_timeout": "10s",
  "proxy_subscribe_endpoint": "grpc://127.0.0.1:10001",
  "proxy_subscribe_timeout": "10s",
  "proxy_refresh_endpoint": "grpc://127.0.0.1:10001",
  "proxy_refresh_timeout": "10s",
  "proxy_sub_refresh_endpoint": "grpc://127.0.0.1:10001",
  "proxy_sub_refresh_timeout": "1s",
  "proxy_rpc_endpoint": "grpc://127.0.0.1:10001",
  "proxy_rpc_timeout": "10s"
}

Note proxy_connect_endpoint, proxy_publish_endpoint, proxy_subscribe_endpoint, proxy_refresh_endpoint, proxy_sub_refresh_endpoint, proxy_rpc_endpoint - endpoint address of roadrunner server with activated centrifuge plugin.

Handling proxy events

Create a worker (app.php) that waits for proxy requests and responds to them:

<?php

require __DIR__ . '/vendor/autoload.php';

use RoadRunner\Centrifugo\CentrifugoWorker;
use RoadRunner\Centrifugo\Payload;
use RoadRunner\Centrifugo\Request;
use RoadRunner\Centrifugo\Request\RequestFactory;
use Spiral\RoadRunner\Worker;

$worker = Worker::create();
$requestFactory = new RequestFactory($worker);

// Create a new Centrifugo Worker from global environment
$centrifugoWorker = new CentrifugoWorker($worker, $requestFactory);

while ($request = $centrifugoWorker->waitRequest()) {

    if ($request instanceof Request\Invalid) {
        $errorMessage = $request->getException()->getMessage();

        if ($request->getException() instanceof \RoadRunner\Centrifugo\Exception\InvalidRequestTypeException) {
            $payload = $request->getException()->payload;
        }

        // Handle invalid request
        // $logger->error($errorMessage, $payload ?? []);

        continue;
    }

    if ($request instanceof Request\Connect) {
        try {
            // Authenticate the connection, e.g. using $request->getData() or $request->headers
            $request->respond(new Payload\ConnectResponse(
                user: '1',
                channels: ['news'],
            ));
        } catch (\Throwable $e) {
            $request->error($e->getCode(), $e->getMessage());
        }

        continue;
    }

    if ($request instanceof Request\Refresh) {
        try {
            // Do something
            $request->respond(new Payload\RefreshResponse(
                // ...
            ));
        } catch (\Throwable $e) {
            $request->error($e->getCode(), $e->getMessage());
        }

        continue;
    }

    if ($request instanceof Request\Subscribe) {
        try {
            // Do something
            $request->respond(new Payload\SubscribeResponse(
                // ...
            ));

            // You can also disconnect connection
            $request->disconnect(4500, 'Connection is not allowed.');
        } catch (\Throwable $e) {
            $request->error($e->getCode(), $e->getMessage());
        }

        continue;
    }

    if ($request instanceof Request\Publish) {
        try {
            // Do something
            $request->respond(new Payload\PublishResponse(
                // ...
            ));

            // You can also disconnect connection
            $request->disconnect(4500, 'Connection is not allowed.');
        } catch (\Throwable $e) {
            $request->error($e->getCode(), $e->getMessage());
        }

        continue;
    }

    if ($request instanceof Request\RPC) {
        try {
            // Handle $request->method with $request->getData() as params
            $response = ['user' => ['id' => 1, 'username' => 'john_smith']];

            $request->respond(new Payload\RPCResponse(
                data: $response,
            ));
        } catch (\Throwable $e) {
            $request->error($e->getCode(), $e->getMessage());
        }

        continue;
    }
}

Proxy events

It's possible to proxy some client connection events from Centrifugo to the RoadRunner application server and react to them in a custom way. For example, it's possible to authenticate connection via request from Centrifugo to application backend, refresh client sessions and answer to RPC calls sent by a client over bidirectional connection.

The list of events that can be proxied:

  • connect – called when a client connects to Centrifugo, so it's possible to authenticate user, return custom data to a client, subscribe connection to several channels, attach meta information to the connection, and so on. Works for bidirectional and unidirectional transports.
  • refresh - called when a client session is going to expire, so it's possible to prolong it or just let it expire. Can also be used just as a periodical connection liveness callback from Centrifugo to app backend. Works for bidirectional and unidirectional transports.
  • sub_refresh - called when it's time to refresh the subscription. Centrifugo itself will ask your backend about subscription validity instead of subscription refresh workflow on the client-side.
  • subscribe - called when clients try to subscribe on a channel, so it's possible to check permissions and return custom initial subscription data. Works for bidirectional transports only.
  • publish - called when a client tries to publish into a channel, so it's possible to check permissions and optionally modify publication data. Works for bidirectional transports only.
  • rpc - called when a client sends RPC, you can do whatever logic you need based on a client-provided RPC method and params. Works for bidirectional transports only.

Note You can find additional information about proxy events here.

Centrifugo server API

RPCCentrifugoApi calls the Centrifugo server API through the RoadRunner centrifuge plugin, so grpc_api_address must point to the Centrifugo gRPC API (enabled with "grpc_api": true in the Centrifugo config).

use RoadRunner\Centrifugo\RPCCentrifugoApi;
use Spiral\Goridge\RPC\RPC;

$api = new RPCCentrifugoApi(RPC::create('tcp://127.0.0.1:6001'));

$api->publish(channel: 'news', message: \json_encode(['text' => 'Hello']));
$api->broadcast(channels: ['news', 'updates'], message: \json_encode(['text' => 'Hello']));
$api->disconnect(user: '1');

$clients = $api->presence(channel: 'news');
$channels = $api->channels();
try Spiral Framework