- Java 100%
| Filename | Latest commit message | Latest commit date |
|---|---|---|
| .github/workflows | ||
| .mvn/wrapper | ||
| src | ||
| .gitignore | ||
| config.json | ||
| docker-compose.yml | ||
| LICENSE | ||
| mvnw | ||
| mvnw.cmd | ||
| mySettings.xml | ||
| pom.xml | ||
| README.md | ||
| swagger.json | ||
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()forpublish()andbroadcast()client.connection()forsubscribe(),unsubscribe(),disconnect(), andrefresh()client.history()forhistory()andhistoryRemove()client.presence()forpresence()andpresenceStats()client.rpc()forrpc()client.stats()forchannels(),connections(), andinfo()client.channels()as a convenience alias for channel lookupsclient.userStatus()forupdateUserStatus(),getUserStatus(), anddeleteUserStatus()client.userBlock()forblockUser()andunblockUser()client.token()forrevokeToken()andinvalidateUserTokens()client.device()for device and topic registration, updates, removal, and listingclient.push()for push notification sending, status updates, and cancellationclient.map()for map publishing/removal, state, stream, stats, clearing, and shared poll publishingclient.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.