No description
Find a file
Repository files (latest commit first)
Filename Latest commit message Latest commit date
Ralph Schaer 7485995bc8 upgrade
2026-10-04 10:06:35 +02:00
.github/workflows upgrade 2026-08-30 20:53:37 +02:00
.mvn/wrapper upgrade 2026-05-31 15:54:45 +02:00
src Pin Centrifugo image for stable Swagger tests 2026-09-05 06:30:21 +02:00
.gitignore Work 2025-07-27 06:34:57 +02:00
config.json work 2025-07-28 10:58:01 +02:00
docker-compose.yml work 2025-08-01 14:31:19 +02:00
LICENSE Work 2025-07-27 06:34:57 +02:00
mvnw Work 2025-07-27 06:34:57 +02:00
mvnw.cmd Work 2025-07-27 06:34:57 +02:00
mySettings.xml Work 2025-07-27 06:34:57 +02:00
pom.xml upgrade 2026-10-04 10:06:35 +02:00
README.md Pin Centrifugo image for stable Swagger tests 2026-09-05 06:30:21 +02:00
swagger.json Add support for new Centrifugo 6.8.x functions 2026-06-27 17:26:48 +02:00

Centrifugo Java Client

An unofficial Java client for the Centrifugo HTTP Server API.

Features

  • Full coverage for the Centrifugo v6 HTTP Server API exposed by Swagger, including experimental and Centrifugo PRO contracts
  • Strongly typed request and response models
  • Fluent builders for request creation
  • OpenFeign-based HTTP client integration
  • Configurable timeouts, retrying, logging, and transport
  • Integration-tested against Centrifugo v6 with Testcontainers

Installation

Add the dependency to your pom.xml:

<dependency>
    <groupId>ch.rasc</groupId>
    <artifactId>jcentserver-client</artifactId>
    <version>2.0.3</version>
</dependency>

Quick Start

import java.util.Map;

import ch.rasc.jcentserverclient.CentrifugoServerApiClient;
import ch.rasc.jcentserverclient.models.PublishRequest;
import ch.rasc.jcentserverclient.models.PublishResponse;

CentrifugoServerApiClient client = CentrifugoServerApiClient.create(config -> config
    .apiKey("your-centrifugo-api-key")
    .baseUrl("http://localhost:8000/api"));

PublishRequest request = PublishRequest.builder()
    .channel("chat:general")
    .data(Map.of("message", "Hello, World!"))
    .build();

PublishResponse response = client.publication().publish(request);

if (!response.hasError()) {
    System.out.println("Published successfully");
}

API Categories

The client facade exposes the following API groups:

  • client.publication() for publish() and broadcast()
  • client.connection() for subscribe(), unsubscribe(), disconnect(), and refresh()
  • client.history() for history() and historyRemove()
  • client.presence() for presence() and presenceStats()
  • client.rpc() for rpc()
  • client.stats() for channels(), connections(), and info()
  • client.channels() as a convenience alias for channel lookups
  • client.userStatus() for updateUserStatus(), getUserStatus(), and deleteUserStatus()
  • client.userBlock() for blockUser() and unblockUser()
  • client.token() for revokeToken() and invalidateUserTokens()
  • client.device() for device and topic registration, updates, removal, and listing
  • client.push() for push notification sending, status updates, and cancellation
  • client.map() for map publishing/removal, state, stream, stats, clearing, and shared poll publishing
  • client.batch() for batched commands

Examples

Publishing

import java.util.Map;

import ch.rasc.jcentserverclient.models.PublishRequest;

PublishRequest request = PublishRequest.builder()
    .channel("notifications")
    .data(Map.of("text", "New notification", "urgency", "high"))
    .build();

client.publication().publish(request);

Publishing with Stream Position Metadata

PublishRequest request = PublishRequest.builder()
    .channel("docs:42")
    .data(Map.of("content", "next revision"))
    .version(8L)
    .versionEpoch("revision-series-1")
    .build();

client.publication().publish(request);

Broadcasting

import java.util.List;
import java.util.Map;

import ch.rasc.jcentserverclient.models.BroadcastRequest;

BroadcastRequest request = BroadcastRequest.builder()
    .channels(List.of("chat:room1", "chat:room2", "chat:room3"))
    .data(Map.of("announcement", "Server maintenance in 10 minutes"))
    .build();

client.publication().broadcast(request);

Map API

import java.util.Map;

import ch.rasc.jcentserverclient.models.MapPublishRequest;

client.map().mapPublish(MapPublishRequest.builder()
    .channel("documents")
    .key("doc-42")
    .data(Map.of("title", "API notes"))
    .score(System.currentTimeMillis())
    .build());

var state = client.map().mapReadState(builder -> builder
    .channel("documents")
    .limit(100)
    .asc(true));

System.out.println(state.result().entries().size());

Connection Management

import ch.rasc.jcentserverclient.models.DisconnectRequest;
import ch.rasc.jcentserverclient.models.SubscribeOptionOverride;
import ch.rasc.jcentserverclient.models.SubscribeRequest;

SubscribeRequest subscribeRequest = SubscribeRequest.builder()
    .user("user123")
    .channel("personal:user123")
    .override(SubscribeOptionOverride.builder()
        .presence(true)
        .joinLeave(true)
        .build())
    .build();

client.connection().subscribe(subscribeRequest);

DisconnectRequest disconnectRequest = DisconnectRequest.builder()
    .user("user123")
    .build();

client.connection().disconnect(disconnectRequest);

Label-based Targeting

import ch.rasc.jcentserverclient.models.FilterNode;
import ch.rasc.jcentserverclient.models.RefreshRequest;

RefreshRequest request = RefreshRequest.builder()
    .allUsers(true)
    .labelFilter(FilterNode.builder()
        .op("eq")
        .key("tenant")
        .val("acme")
        .build())
    .build();

client.connection().refresh(request);

Custom RPC

import java.util.Map;

import ch.rasc.jcentserverclient.models.RpcRequest;

RpcRequest request = RpcRequest.builder()
    .method("ping")
    .params(Map.of("value", "hello"))
    .build();

var response = client.rpc().rpc(request);

if (response.error() != null) {
    System.err.println("RPC failed: " + response.error().message());
}

User Blocking and Tokens

import ch.rasc.jcentserverclient.models.InvalidateUserTokensRequest;

client.userBlock().blockUser(builder -> builder
    .user("user123")
    .expireAt(System.currentTimeMillis() / 1000 + 3600));

client.token().revokeToken("token-uid");

client.token().invalidateUserTokens(InvalidateUserTokensRequest.builder()
    .user("user123")
    .channel("private:user123")
    .build());

Server Info

import ch.rasc.jcentserverclient.models.InfoRequest;

var response = client.stats().info(InfoRequest.builder().build());
System.out.println(response.result().nodes().size());

Push Notifications

import java.util.List;
import java.util.Map;

import ch.rasc.jcentserverclient.models.PushNotification;
import ch.rasc.jcentserverclient.models.PushRecipient;
import ch.rasc.jcentserverclient.models.SendPushNotificationRequest;
import ch.rasc.jcentserverclient.models.WebPushPushNotification;

SendPushNotificationRequest request = SendPushNotificationRequest.builder()
    .uid("push-1")
    .recipient(new PushRecipient(null, null, null, null, null, null, null, null,
        List.of("webpush-token")))
    .notification(new PushNotification(null, null, null,
        new WebPushPushNotification(Map.of("TTL", "60"), Map.of("title", "Hello")),
        System.currentTimeMillis() / 1000 + 3600))
    .build();

client.push().sendPushNotification(request);

Configuration

Basic Configuration

CentrifugoServerApiClient client = CentrifugoServerApiClient.create(config -> config
    .apiKey("your-api-key")
    .baseUrl("http://localhost:8000/api"));

Advanced Configuration

import ch.rasc.jcentserverclient.Configuration;

import feign.Logger;
import feign.http2client.Http2Client;

Configuration configuration = Configuration.builder()
    .apiKey("your-api-key")
    .baseUrl("https://your-centrifugo-server.com/api")
    .client(new Http2Client())
    .transportErrorMode(true)
    .logLevel(Logger.Level.BASIC)
    .build();

CentrifugoServerApiClient client = CentrifugoServerApiClient.create(configuration);

Error Handling

import ch.rasc.jcentserverclient.ApiException;

try {
    var response = client.publication().publish(request);

    if (response.hasError()) {
        System.err.println(
            "Error: " + response.error().message() + " (Code: " + response.error().code() + ")");
    }
}
catch (ApiException e) {
    System.err.println("API Error: " + e.getMessage());
    System.err.println("Status: " + e.status());
}

Requirements

  • Java 17 or higher
  • Centrifugo with the HTTP Server API enabled

Authentication

When apiKey is configured, the client authenticates with the X-API-Key header. Configure the same key on both sides.

Example Centrifugo config:

{
  "http_api": {
    "key": "your-secret-api-key"
  }
}

The API key is optional so the client can also connect to a server configured with http_api.insecure (development only), or use an additionalRequestInterceptor for another authentication mechanism.

Set transportErrorMode(true) to send X-Centrifugo-Error-Mode: transport on every request. In that mode, ordinary API errors are decoded into ApiException. Batch and broadcast calls can contain per-operation errors and must still be inspected individually.

Testing

The integration test suite starts a centrifugo/centrifugo:v6.9.3 container automatically via Testcontainers. The version is pinned so that the checked-in Swagger snapshot is compared with the same API contract on every platform, even when an older v6 image is cached locally. LiveSwaggerAlignmentIntegrationTest compares the checked-in swagger.json with the document served by that live container. SwaggerAlignmentTest additionally verifies every endpoint's request/response binding and every model's JSON property names and Java types.

Run the full suite:

./mvnw test

Run a focused test class:

./mvnw -Dtest=RequestSerializationTest test

Contributing

Contributions are welcome. Open an issue or submit a pull request.

License

This project is licensed under the Apache License 2.0. See LICENSE for details.