@@ -311,6 +317,7 @@
**/store/RedissonStoreTest.java
**/store/event/HazelcastRingBufferEventStoreTest.java
**/store/event/RedisPubSubEventStoreTest.java
+ **/store/event/MongoPubSubEventStoreTest.java
diff --git a/netty-socketio-core/src/main/java/com/socketio4j/socketio/store/mongo/MongoEventStore.java b/netty-socketio-core/src/main/java/com/socketio4j/socketio/store/mongo/MongoEventStore.java
new file mode 100644
index 00000000..a410b8be
--- /dev/null
+++ b/netty-socketio-core/src/main/java/com/socketio4j/socketio/store/mongo/MongoEventStore.java
@@ -0,0 +1,1063 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.store.mongo;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Date;
+import java.util.List;
+import java.util.Objects;
+import java.util.Queue;
+import java.util.Set;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.Executors;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.bson.BsonDocument;
+import org.bson.BsonTimestamp;
+import org.bson.Document;
+import org.bson.conversions.Bson;
+import org.jetbrains.annotations.NotNull;
+import org.jetbrains.annotations.Nullable;
+import org.reactivestreams.Subscriber;
+import org.reactivestreams.Subscription;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.mongodb.MongoCommandException;
+import com.mongodb.ReadPreference;
+import com.mongodb.WriteConcern;
+import com.mongodb.client.model.Aggregates;
+import com.mongodb.client.model.Filters;
+import com.mongodb.client.model.IndexOptions;
+import com.mongodb.client.model.Indexes;
+import com.mongodb.client.model.changestream.ChangeStreamDocument;
+import com.mongodb.client.result.InsertOneResult;
+import com.mongodb.reactivestreams.client.ChangeStreamPublisher;
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.mongodb.reactivestreams.client.MongoCollection;
+import com.mongodb.reactivestreams.client.MongoDatabase;
+import com.socketio4j.socketio.store.event.EventListener;
+import com.socketio4j.socketio.store.event.EventMessage;
+import com.socketio4j.socketio.store.event.EventMessageJsonSupport;
+import com.socketio4j.socketio.store.event.EventStore;
+import com.socketio4j.socketio.store.event.EventStoreMode;
+import com.socketio4j.socketio.store.event.EventStoreType;
+import com.socketio4j.socketio.store.event.EventType;
+import com.socketio4j.socketio.store.event.PublishMode;
+
+/**
+ * MongoDB Change Streams based EventStore.
+ *
+ * Uses MongoDB Change Streams to watch for inserts on a collection and deliver
+ * events to subscribers. Each event type maps to its own collection
+ * (MULTI_CHANNEL)
+ * or all events go into one collection (SINGLE_CHANNEL).
+ *
+ * Built on the reactive streams driver: fully non-blocking and asynchronous.
+ * Both {@code publish0} and {@code subscribe0} (and {@code subscribeAsync})
+ * are entirely non-blocking. TTL index creation, reconciliation, and change stream
+ * initialization run asynchronously via Reactive Streams subscribers without blocking
+ * the calling threads.
+ *
+ * A TTL index is created on each collection to automatically expire documents
+ * after a configurable retention period (default 60 seconds), preventing
+ * unbounded data growth.
+ *
+ * Requires a MongoDB replica set (standalone does not support change streams).
+ */
+public class MongoEventStore implements EventStore {
+
+ private static final Logger log = LoggerFactory.getLogger(MongoEventStore.class);
+
+ private static final String DEFAULT_COLLECTION_PREFIX = "socketio_events_";
+ private static final long DEFAULT_TTL_SECONDS = 60;
+
+ /** Initial delay before reopening a change stream that ended or failed. */
+ private static final long INITIAL_REOPEN_DELAY_MILLIS = 500;
+
+ /** Maximum backoff delay for reopening a change stream. */
+ private static final long MAX_REOPEN_DELAY_MILLIS = 8000;
+
+ /**
+ * How long {@link #shutdown0()} waits for cancelled change streams to actually
+ * end.
+ * Cancelling is asynchronous, and the driver may still have a {@code getMore}
+ * in flight;
+ * bounded by {@link #CANCEL_GRACE_MILLIS} and returns early once cancelled.
+ */
+ private static final long CANCEL_GRACE_MILLIS = 1000;
+
+ /**
+ * MongoDB error code raised when an index exists with the same key but
+ * different options.
+ */
+ private static final int INDEX_OPTIONS_CONFLICT = 85;
+
+ /**
+ * Shared mapper: keeps byte[] payloads lossless, as the Kafka and NATS stores
+ * do.
+ */
+ private static final ObjectMapper MAPPER = EventMessageJsonSupport.createObjectMapper();
+
+ private final MongoDatabase database;
+ private final Long nodeId;
+ private final EventStoreMode eventStoreMode;
+ private final String collectionPrefix;
+ private final long ttlSeconds;
+ private final WriteConcern writeConcern;
+ private final ReadPreference readPreference;
+
+ private final AtomicBoolean running = new AtomicBoolean(true);
+
+ private final ConcurrentMap> collectionCache = new ConcurrentHashMap<>();
+
+ private final Set indexedCollections = ConcurrentHashMap.newKeySet();
+
+ private final ConcurrentMap> watchers = new ConcurrentHashMap<>();
+
+ private final ScheduledExecutorService watcherExecutor = Executors.newSingleThreadScheduledExecutor(r -> {
+ Thread t = new Thread(r);
+ t.setName("socketio-mongo-watcher-" + t.getId());
+ t.setDaemon(true);
+ return t;
+ });
+
+ /**
+ * Creates a new MongoEventStore.
+ *
+ * @param mongoClient shared MongoDB client
+ * @param databaseName database to use for event collections
+ * @param eventStoreMode SINGLE_CHANNEL or MULTI_CHANNEL (defaults to
+ * MULTI_CHANNEL)
+ * @param nodeId node identifier used to ignore self-published events
+ * @param collectionPrefix prefix for collection names (defaults to
+ * "socketio_events_")
+ * @param ttlSeconds TTL in seconds for automatic document expiry
+ * (defaults to 60)
+ */
+ public MongoEventStore(@NotNull MongoClient mongoClient,
+ @NotNull String databaseName,
+ @Nullable EventStoreMode eventStoreMode,
+ @Nullable Long nodeId,
+ @Nullable String collectionPrefix,
+ long ttlSeconds) {
+ this(mongoClient, databaseName, eventStoreMode, nodeId, collectionPrefix, ttlSeconds, null, null);
+ }
+
+ /**
+ * Creates a new MongoEventStore with custom write concern and read preference.
+ *
+ * @param mongoClient shared MongoDB client
+ * @param databaseName database to use for event collections
+ * @param eventStoreMode SINGLE_CHANNEL or MULTI_CHANNEL (defaults to
+ * MULTI_CHANNEL)
+ * @param nodeId node identifier used to ignore self-published events
+ * @param collectionPrefix prefix for collection names (defaults to
+ * "socketio_events_")
+ * @param ttlSeconds TTL in seconds for automatic document expiry
+ * (defaults to 60)
+ * @param writeConcern optional write concern for published events
+ * @param readPreference optional read preference for change stream operations
+ */
+ public MongoEventStore(@NotNull MongoClient mongoClient,
+ @NotNull String databaseName,
+ @Nullable EventStoreMode eventStoreMode,
+ @Nullable Long nodeId,
+ @Nullable String collectionPrefix,
+ long ttlSeconds,
+ @Nullable WriteConcern writeConcern,
+ @Nullable ReadPreference readPreference) {
+ Objects.requireNonNull(mongoClient, "mongoClient");
+ this.database = mongoClient.getDatabase(
+ Objects.requireNonNull(databaseName, "databaseName"));
+
+ if (nodeId != null) {
+ this.nodeId = nodeId;
+ } else {
+ this.nodeId = getNodeId();
+ }
+ if (eventStoreMode != null) {
+ this.eventStoreMode = eventStoreMode;
+ } else {
+ this.eventStoreMode = EventStoreMode.MULTI_CHANNEL;
+ }
+ if (collectionPrefix != null) {
+ this.collectionPrefix = collectionPrefix;
+ } else {
+ this.collectionPrefix = DEFAULT_COLLECTION_PREFIX;
+ }
+ if (ttlSeconds > 0) {
+ this.ttlSeconds = ttlSeconds;
+ } else {
+ this.ttlSeconds = DEFAULT_TTL_SECONDS;
+ }
+ this.writeConcern = writeConcern;
+ this.readPreference = readPreference;
+ }
+
+ @Override
+ public EventStoreMode getEventStoreMode() {
+ return eventStoreMode;
+ }
+
+ @Override
+ public EventStoreType getEventStoreType() {
+ return EventStoreType.PUBSUB;
+ }
+
+ @Override
+ public PublishMode getPublishMode() {
+ return PublishMode.UNRELIABLE;
+ }
+
+ public String getCollectionPrefix() {
+ return collectionPrefix;
+ }
+
+ public long getTtlSeconds() {
+ return ttlSeconds;
+ }
+
+ public WriteConcern getWriteConcern() {
+ return writeConcern;
+ }
+
+ public ReadPreference getReadPreference() {
+ return readPreference;
+ }
+
+ @Override
+ public void publish0(EventType type, EventMessage msg) {
+ if (!running.get()) {
+ log.warn("MongoEventStore is shut down; ignoring publish of {}", type);
+ return;
+ }
+
+ msg.setNodeId(nodeId);
+
+ String collectionName = getCollectionName(type);
+ String payload;
+ try {
+ payload = MAPPER.writeValueAsString(msg);
+ } catch (Exception e) {
+ throw new IllegalStateException("Failed to serialize EventMessage", e);
+ }
+ Document doc = new Document()
+ .append("nodeId", nodeId)
+ .append("eventType", type.name())
+ .append("createdAt", new Date())
+ .append("payload", payload);
+
+ // Fire and forget: this runs on the Netty event loop, so the insert is handed
+ // to
+ // the driver and its outcome only logged, the model KafkaEventStore.publish0
+ // uses.
+ getCollection(collectionName)
+ .insertOne(doc)
+ .subscribe(new PublishSubscriber(type, collectionName));
+ }
+
+ @Override
+ public void subscribe0(
+ EventType type,
+ final EventListener listener,
+ Class clazz) {
+ if (!running.get()) {
+ throw new IllegalStateException("MongoEventStore has been shut down");
+ }
+ subscribeAsync(type, listener, clazz);
+ }
+
+ /**
+ * Asynchronously subscribes to the specified event type without blocking.
+ *
+ * Initiates asynchronous TTL index creation, captures the cluster operation time
+ * asynchronously, and opens the change stream cursor.
+ *
+ * @param type event type to subscribe to
+ * @param listener listener for received messages
+ * @param clazz event message class
+ * @param message type
+ * @return a {@link CompletableFuture} that completes when the change stream subscription is established
+ */
+ public CompletableFuture subscribeAsync(
+ EventType type,
+ final EventListener listener,
+ Class clazz) {
+
+ if (!running.get()) {
+ throw new IllegalStateException("MongoEventStore has been shut down");
+ }
+
+ Objects.requireNonNull(listener, "listener cannot be null");
+ Objects.requireNonNull(clazz, "clazz cannot be null");
+
+ validateSubscribe(type);
+
+ String collectionName = getCollectionName(type);
+ MongoCollection collection = getCollection(collectionName);
+
+ if (indexedCollections.add(collectionName)) {
+ ensureTtlIndexAsync(collection, collectionName);
+ }
+
+ WatcherHandle handle = new WatcherHandle(null);
+
+ // Register before opening the stream in a single atomic map operation
+ watchers.compute(type, (k, queue) -> {
+ Queue q = queue;
+ if (q == null) {
+ q = new ConcurrentLinkedQueue<>();
+ }
+ q.add(handle);
+ return q;
+ });
+
+ CompletableFuture subscriptionFuture = new CompletableFuture<>();
+ startWatcherAsync(collection, type, handle, listener, clazz, subscriptionFuture);
+ return subscriptionFuture;
+ }
+
+ /**
+ * Builds the server-side Change Stream aggregation pipeline.
+ *
+ * Filters for inserts only (ignoring TTL deletes) and drops self-published
+ * events on the MongoDB server before network transmission.
+ */
+ private List createPipeline() {
+ return Collections.singletonList(
+ Aggregates.match(Filters.and(
+ Filters.eq("operationType", "insert"),
+ Filters.ne("fullDocument.nodeId", nodeId))));
+ }
+
+ /**
+ * Starts the watcher asynchronously by retrieving the server's operation time via ping
+ * and opening the change stream at that point without blocking the calling thread.
+ */
+ private void startWatcherAsync(
+ MongoCollection collection,
+ EventType type,
+ WatcherHandle handle,
+ EventListener listener,
+ Class clazz,
+ CompletableFuture subscriptionFuture) {
+
+ if (handle.stopped.get()) {
+ subscriptionFuture.complete(null);
+ return;
+ }
+
+ database.runCommand(new Document("ping", 1)).subscribe(new Subscriber() {
+ private BsonTimestamp opTime;
+
+ @Override
+ public void onSubscribe(Subscription s) {
+ s.request(1);
+ }
+
+ @Override
+ public void onNext(Document doc) {
+ if (doc != null) {
+ Object operationTime = doc.get("operationTime");
+ if (operationTime instanceof BsonTimestamp) {
+ opTime = (BsonTimestamp) operationTime;
+ }
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {
+ log.debug("Ping command failed to retrieve operationTime on {}, starting watch with server default: {}",
+ collection.getNamespace(), t.getMessage());
+ if (!handle.stopped.get()) {
+ watch(collection, type, handle, listener, clazz, subscriptionFuture);
+ } else {
+ subscriptionFuture.complete(null);
+ }
+ }
+
+ @Override
+ public void onComplete() {
+ if (handle.stopped.get()) {
+ subscriptionFuture.complete(null);
+ return;
+ }
+ if (opTime != null) {
+ handle.setStartAt(opTime);
+ }
+ watch(collection, type, handle, listener, clazz, subscriptionFuture);
+ }
+ });
+ }
+
+ /**
+ * Opens the change stream, resuming after the last delivered event when this is a
+ * reopen, and otherwise starting at the operation time captured while subscribing.
+ */
+ private void watch(MongoCollection collection,
+ EventType type,
+ WatcherHandle handle,
+ EventListener listener,
+ Class clazz,
+ @Nullable CompletableFuture subscriptionFuture) {
+ if (handle.stopped.get()) {
+ if (subscriptionFuture != null && !subscriptionFuture.isDone()) {
+ subscriptionFuture.complete(null);
+ }
+ return;
+ }
+
+ ChangeStreamPublisher stream = collection.watch(createPipeline());
+ BsonDocument resumeToken = handle.resumeToken();
+ if (resumeToken != null) {
+ stream = stream.resumeAfter(resumeToken);
+ } else if (handle.startAt() != null) {
+ stream = stream.startAtOperationTime(handle.startAt());
+ }
+
+ stream.subscribe(new ChangeSubscriber(collection, type, handle, listener, clazz, subscriptionFuture));
+ }
+
+ @Override
+ public void unsubscribe0(EventType type) {
+ Queue handles = watchers.remove(type);
+ if (handles == null || handles.isEmpty()) {
+ return;
+ }
+ for (WatcherHandle handle : handles) {
+ handle.stop();
+ }
+ }
+
+ @Override
+ public void shutdown0() {
+ shutdownAsync().join();
+ }
+
+ /**
+ * Asynchronously shuts down the Mongo event store.
+ *
+ * @return a {@link CompletableFuture} completing when all watchers are cancelled and executor is stopped
+ */
+ public CompletableFuture shutdownAsync() {
+ if (!running.compareAndSet(true, false)) {
+ return CompletableFuture.completedFuture(null);
+ }
+
+ List cancelled = new ArrayList();
+ for (Queue handles : watchers.values()) {
+ cancelled.addAll(handles);
+ }
+
+ Arrays.stream(EventType.values()).forEach(this::unsubscribe);
+ watchers.clear();
+
+ return CompletableFuture.runAsync(() -> {
+ awaitCancellation(cancelled);
+
+ watcherExecutor.shutdown();
+ try {
+ if (!watcherExecutor.awaitTermination(5, TimeUnit.SECONDS)) {
+ watcherExecutor.shutdownNow();
+ }
+ } catch (InterruptedException e) {
+ watcherExecutor.shutdownNow();
+ Thread.currentThread().interrupt();
+ }
+ });
+ }
+
+ /**
+ * Waits for the change streams cancelled by {@link #shutdown0()} to end, so that a
+ * caller closing its {@code MongoClient} right afterwards does not interrupt them
+ * mid-flight.
+ */
+ private void awaitCancellation(List handles) {
+ if (handles.isEmpty()) {
+ return;
+ }
+
+ long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(CANCEL_GRACE_MILLIS);
+ try {
+ for (WatcherHandle handle : handles) {
+ long remaining = deadline - System.nanoTime();
+ if (remaining <= 0) {
+ return;
+ }
+ handle.awaitTermination(remaining);
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ /**
+ * Creates or reconciles the TTL index used to expire published events asynchronously.
+ *
+ * {@code createIndex} creates the collection when it does not exist yet and is a
+ * no-op for an identical index. If an index already exists with a different retention period,
+ * it catches {@code INDEX_OPTIONS_CONFLICT} (error code 85) and asynchronously reconciles
+ * the TTL using {@code collMod}.
+ */
+ private void ensureTtlIndexAsync(MongoCollection collection, String collectionName) {
+ collection.createIndex(
+ Indexes.ascending("createdAt"),
+ new IndexOptions().expireAfter(ttlSeconds, TimeUnit.SECONDS))
+ .subscribe(new Subscriber() {
+ @Override
+ public void onSubscribe(Subscription s) {
+ s.request(1);
+ }
+
+ @Override
+ public void onNext(String indexName) {
+ log.debug("Ensured TTL index {} on {}", indexName, collection.getNamespace());
+ }
+
+ @Override
+ public void onError(Throwable t) {
+ if (t instanceof MongoCommandException
+ && ((MongoCommandException) t).getErrorCode() == INDEX_OPTIONS_CONFLICT) {
+ reconcileTtlIndexAsync(collection, collectionName);
+ } else {
+ indexedCollections.remove(collectionName);
+ log.error("Failed to apply TTL index of {}s on {}", ttlSeconds, collection.getNamespace(), t);
+ }
+ }
+
+ @Override
+ public void onComplete() {
+ }
+ });
+ }
+
+ private void reconcileTtlIndexAsync(MongoCollection collection, String collectionName) {
+ Document collModCmd = new Document("collMod", collection.getNamespace().getCollectionName())
+ .append("index", new Document("keyPattern", new Document("createdAt", 1))
+ .append("expireAfterSeconds", ttlSeconds));
+
+ database.runCommand(collModCmd).subscribe(new Subscriber() {
+ @Override
+ public void onSubscribe(Subscription s) {
+ s.request(1);
+ }
+
+ @Override
+ public void onNext(Document doc) {
+ log.info("Updated TTL index on {} to {} seconds via collMod", collection.getNamespace(), ttlSeconds);
+ }
+
+ @Override
+ public void onError(Throwable t) {
+ indexedCollections.remove(collectionName);
+ log.error("Failed to reconcile TTL index on {}", collection.getNamespace(), t);
+ }
+
+ @Override
+ public void onComplete() {
+ }
+ });
+ }
+
+ /**
+ * Rejects a subscription whose {@link EventType} does not match the configured
+ * mode.
+ *
+ * In SINGLE_CHANNEL mode every event type shares one collection, so only the
+ * {@code ALL_SINGLE_CHANNEL} subscription — the one {@code BaseStoreFactory}
+ * opens —
+ * is meaningful; a per-type subscription would silently watch the shared
+ * collection
+ * and receive unrelated event types. Mirrors {@code RedisStreamEventStore}.
+ */
+ private void validateSubscribe(EventType type) {
+ if (EventStoreMode.SINGLE_CHANNEL.equals(eventStoreMode)
+ && type != EventType.ALL_SINGLE_CHANNEL) {
+ throw new UnsupportedOperationException(
+ "Only ALL_SINGLE_CHANNEL allowed in SINGLE_CHANNEL mode");
+ }
+ if (EventStoreMode.MULTI_CHANNEL.equals(eventStoreMode)
+ && type == EventType.ALL_SINGLE_CHANNEL) {
+ throw new UnsupportedOperationException(
+ "ALL_SINGLE_CHANNEL not allowed in MULTI_CHANNEL mode");
+ }
+ }
+
+ private String getCollectionName(EventType type) {
+ if (EventStoreMode.SINGLE_CHANNEL.equals(eventStoreMode)) {
+ return collectionPrefix + EventType.ALL_SINGLE_CHANNEL.name();
+ }
+ return collectionPrefix + type.name();
+ }
+
+ private MongoCollection getCollection(String collectionName) {
+ return collectionCache.computeIfAbsent(collectionName, name -> {
+ MongoCollection col = database.getCollection(name);
+ if (writeConcern != null) {
+ col = col.withWriteConcern(writeConcern);
+ }
+ if (readPreference != null) {
+ col = col.withReadPreference(readPreference);
+ }
+ return col;
+ });
+ }
+
+ /**
+ * Logs a failed insert; a published event is never retried, as in the Kafka
+ * store.
+ */
+ private static final class PublishSubscriber implements Subscriber {
+
+ private final EventType type;
+ private final String collectionName;
+
+ PublishSubscriber(EventType type, String collectionName) {
+ this.type = type;
+ this.collectionName = collectionName;
+ }
+
+ @Override
+ public void onSubscribe(Subscription subscription) {
+ subscription.request(1);
+ }
+
+ @Override
+ public void onNext(InsertOneResult result) {
+ // the insert result carries nothing this store needs
+ }
+
+ @Override
+ public void onError(Throwable error) {
+ log.warn("Failed to publish {} to {}", type, collectionName, error);
+ }
+
+ @Override
+ public void onComplete() {
+ // nothing to do
+ }
+ }
+
+ /**
+ * Delivers change events to one listener and reopens the stream when it ends,
+ * which
+ * replaces the reconnect loop a blocking cursor needed.
+ */
+ final class ChangeSubscriber
+ implements Subscriber> {
+
+ private static final long DEMAND_BATCH_SIZE = 128;
+
+ private final MongoCollection collection;
+ private final EventType type;
+ private final WatcherHandle handle;
+ private final EventListener listener;
+ private final Class clazz;
+
+ private final CompletableFuture initialFuture;
+
+ ChangeSubscriber(MongoCollection collection,
+ EventType type,
+ WatcherHandle handle,
+ EventListener listener,
+ Class clazz) {
+ this(collection, type, handle, listener, clazz, null);
+ }
+
+ ChangeSubscriber(MongoCollection collection,
+ EventType type,
+ WatcherHandle handle,
+ EventListener listener,
+ Class clazz,
+ @Nullable CompletableFuture initialFuture) {
+ this.collection = collection;
+ this.type = type;
+ this.handle = handle;
+ this.listener = listener;
+ this.clazz = clazz;
+ this.initialFuture = initialFuture;
+ }
+
+ @Override
+ public void onSubscribe(Subscription subscription) {
+ handle.setSubscription(subscription);
+ subscription.request(DEMAND_BATCH_SIZE);
+ if (initialFuture != null && !initialFuture.isDone()) {
+ initialFuture.complete(null);
+ }
+ }
+
+ @Override
+ public void onNext(ChangeStreamDocument change) {
+ if (handle.stopped.get()) {
+ return;
+ }
+
+ handle.resetReconnectAttempts();
+ handle.setResumeToken(change.getResumeToken());
+
+ // Replenish demand to maintain backpressure flow control
+ Subscription sub = handle.getSubscription();
+ if (sub != null && !handle.stopped.get()) {
+ sub.request(1);
+ }
+
+ Document doc = change.getFullDocument();
+ if (doc == null) {
+ return;
+ }
+
+ // A node sees its own inserts too, so drop them on the document's own nodeId
+ // rather than paying for the deserialization first.
+ Object rawNodeId = doc.get("nodeId");
+ if (rawNodeId instanceof Number && nodeId.longValue() == ((Number) rawNodeId).longValue()) {
+ return;
+ }
+
+ if (EventStoreMode.MULTI_CHANNEL.equals(eventStoreMode)) {
+ String eventTypeName = doc.getString("eventType");
+ if (eventTypeName != null && !type.name().equals(eventTypeName)) {
+ return;
+ }
+ }
+
+ String payload = doc.getString("payload");
+ if (payload == null) {
+ return;
+ }
+
+ try {
+ T event = MAPPER.readValue(payload, clazz);
+ if (handle.stopped.get()) {
+ return;
+ }
+ // Belt and braces: a document written without the nodeId field still gets
+ // filtered on the value carried by the event itself.
+ if (event.getNodeId() == null || nodeId.longValue() != event.getNodeId().longValue()) {
+ listener.onMessage(event);
+ }
+ } catch (Exception e) {
+ log.warn("Failed to process change event on {}", collection.getNamespace(), e);
+ }
+ }
+
+ @Override
+ public void onError(Throwable error) {
+ if (handle.stopped.get()) {
+ handle.markTerminated();
+ if (initialFuture != null && !initialFuture.isDone()) {
+ initialFuture.complete(null);
+ }
+ return;
+ }
+
+ if (initialFuture != null && !initialFuture.isDone()) {
+ initialFuture.completeExceptionally(error);
+ }
+
+ if (isChangeStreamHistoryLost(error)) {
+ log.error("Change stream history lost on {}. Clearing resume point and restarting from current time.",
+ collection.getNamespace(), error);
+ handle.clearResumePoint();
+ } else if (isStandaloneError(error)) {
+ log.error("MongoDB change streams require a replica set or sharded cluster. "
+ + "Standalone deployments are not supported on {}: {}",
+ collection.getNamespace(), error.getMessage());
+ } else {
+ log.warn("Change stream on {} failed, reopening...", collection.getNamespace(), error);
+ }
+ reopen();
+ }
+
+ @Override
+ public void onComplete() {
+ if (handle.stopped.get()) {
+ handle.markTerminated();
+ return;
+ }
+ // The server ended the stream, for instance because the collection was dropped.
+ reopen();
+ }
+
+ private void reopen() {
+ int attempts = handle.incrementReconnectAttempts();
+ long delay = Math.min(
+ INITIAL_REOPEN_DELAY_MILLIS * (1L << Math.min(attempts - 1, 4)),
+ MAX_REOPEN_DELAY_MILLIS);
+ try {
+ watcherExecutor.schedule(
+ () -> watch(collection, type, handle, listener, clazz, null),
+ delay, TimeUnit.MILLISECONDS);
+ } catch (RejectedExecutionException e) {
+ log.debug("Not reopening the change stream on {}, the store is shutting down",
+ collection.getNamespace());
+ }
+ }
+
+ private boolean isChangeStreamHistoryLost(Throwable error) {
+ Throwable curr = error;
+ while (curr != null) {
+ if (curr instanceof MongoCommandException) {
+ int code = ((MongoCommandException) curr).getErrorCode();
+ // 286 = ChangeStreamHistoryLost, 40585 = ChangeStreamFatalError, 40579 =
+ // ChangeStreamTimedOut
+ if (code == 286 || code == 40585 || code == 40579) {
+ return true;
+ }
+ }
+ String msg = curr.getMessage();
+ if (msg != null && (msg.contains("history lost") || msg.contains("ChangeStreamHistoryLost")
+ || msg.contains("resume point may no longer be in the oplog"))) {
+ return true;
+ }
+ curr = curr.getCause();
+ }
+ return false;
+ }
+
+ private boolean isStandaloneError(Throwable error) {
+ Throwable curr = error;
+ while (curr != null) {
+ if (curr instanceof MongoCommandException) {
+ int code = ((MongoCommandException) curr).getErrorCode();
+ // 40573: The $changeStream stage is only supported on replica sets
+ if (code == 40573) {
+ return true;
+ }
+ }
+ String msg = curr.getMessage();
+ if (msg != null && (msg.contains("only supported on replica sets")
+ || msg.contains("Change streams are only supported on replica sets")
+ || msg.contains("The $changeStream stage is only supported on replica sets"))) {
+ return true;
+ }
+ curr = curr.getCause();
+ }
+ return false;
+ }
+ }
+
+ static final class WatcherHandle {
+
+ private static final Subscription CANCELLED = new Subscription() {
+ @Override
+ public void request(long n) {
+ }
+
+ @Override
+ public void cancel() {
+ }
+ };
+
+ final AtomicBoolean stopped = new AtomicBoolean(false);
+
+ private final CountDownLatch terminated = new CountDownLatch(1);
+ private final AtomicReference subRef = new AtomicReference<>();
+ private final AtomicInteger reconnectAttempts = new AtomicInteger(0);
+
+ private volatile BsonTimestamp startAt;
+ private volatile BsonDocument resumeToken;
+
+ WatcherHandle(BsonTimestamp startAt) {
+ this.startAt = startAt;
+ }
+
+ BsonTimestamp startAt() {
+ return startAt;
+ }
+
+ void setStartAt(BsonTimestamp startAt) {
+ this.startAt = startAt;
+ }
+
+ void clearResumePoint() {
+ this.resumeToken = null;
+ this.startAt = null;
+ }
+
+ int incrementReconnectAttempts() {
+ return reconnectAttempts.incrementAndGet();
+ }
+
+ void resetReconnectAttempts() {
+ reconnectAttempts.set(0);
+ }
+
+ /**
+ * Atomically publishes the subscription so {@link #stop()} can cancel it.
+ * If stop already happened, cancels the subscription immediately.
+ */
+ void setSubscription(Subscription subscription) {
+ while (!stopped.get()) {
+ Subscription current = subRef.get();
+ if (current == CANCELLED) {
+ subscription.cancel();
+ return;
+ }
+ if (subRef.compareAndSet(current, subscription)) {
+ if (stopped.get()) {
+ if (subRef.compareAndSet(subscription, CANCELLED)) {
+ subscription.cancel();
+ }
+ }
+ return;
+ }
+ }
+ subscription.cancel();
+ }
+
+ Subscription getSubscription() {
+ Subscription s = subRef.get();
+ if (s == CANCELLED) {
+ return null;
+ }
+ return s;
+ }
+
+ void setResumeToken(BsonDocument resumeToken) {
+ this.resumeToken = resumeToken;
+ }
+
+ BsonDocument resumeToken() {
+ return resumeToken;
+ }
+
+ void stop() {
+ stopped.set(true);
+ Subscription s = subRef.getAndSet(CANCELLED);
+ if (s != null && s != CANCELLED) {
+ try {
+ s.cancel();
+ } catch (Exception e) {
+ log.debug("Error cancelling change stream subscription", e);
+ }
+ }
+ terminated.countDown();
+ }
+
+ void markTerminated() {
+ terminated.countDown();
+ }
+
+ void awaitTermination(long nanos) throws InterruptedException {
+ terminated.await(nanos, TimeUnit.NANOSECONDS);
+ }
+ }
+
+ /**
+ * Builder for {@link MongoEventStore}.
+ */
+ public static final class Builder {
+
+ private final MongoClient mongoClient;
+ private final String databaseName;
+
+ private Long nodeId;
+ private EventStoreMode eventStoreMode = EventStoreMode.MULTI_CHANNEL;
+ private String collectionPrefix;
+ private long ttlSeconds = DEFAULT_TTL_SECONDS;
+ private WriteConcern writeConcern;
+ private ReadPreference readPreference;
+
+ /**
+ * Creates a new builder.
+ *
+ * @param mongoClient shared MongoDB client (must connect to a replica set)
+ * @param databaseName database to use for event collections
+ */
+ public Builder(@NotNull MongoClient mongoClient, @NotNull String databaseName) {
+ this.mongoClient = Objects.requireNonNull(mongoClient, "mongoClient");
+ this.databaseName = Objects.requireNonNull(databaseName, "databaseName");
+ }
+
+ public Builder nodeId(long nodeId) {
+ this.nodeId = nodeId;
+ return this;
+ }
+
+ public Builder eventStoreMode(@NotNull EventStoreMode mode) {
+ this.eventStoreMode = Objects.requireNonNull(mode, "eventStoreMode");
+ return this;
+ }
+
+ public Builder collectionPrefix(@NotNull String prefix) {
+ this.collectionPrefix = Objects.requireNonNull(prefix, "collectionPrefix");
+ return this;
+ }
+
+ /**
+ * Sets the TTL in seconds for automatic document expiry.
+ * Documents older than this are automatically removed by MongoDB.
+ * Default is 60 seconds.
+ *
+ * @param ttlSeconds retention period in seconds
+ * @return this builder
+ */
+ public Builder ttlSeconds(long ttlSeconds) {
+ this.ttlSeconds = ttlSeconds;
+ return this;
+ }
+
+ /**
+ * Sets the write concern to use for publishing events.
+ *
+ * @param writeConcern write concern (e.g. WriteConcern.W1 or UNACKNOWLEDGED)
+ * @return this builder
+ */
+ public Builder writeConcern(@NotNull WriteConcern writeConcern) {
+ this.writeConcern = Objects.requireNonNull(writeConcern, "writeConcern");
+ return this;
+ }
+
+ /**
+ * Sets the read preference for change stream cursors.
+ *
+ * @param readPreference read preference (e.g.
+ * ReadPreference.secondaryPreferred())
+ * @return this builder
+ */
+ public Builder readPreference(@NotNull ReadPreference readPreference) {
+ this.readPreference = Objects.requireNonNull(readPreference, "readPreference");
+ return this;
+ }
+
+ public MongoEventStore build() {
+ return new MongoEventStore(
+ mongoClient,
+ databaseName,
+ eventStoreMode,
+ nodeId,
+ collectionPrefix,
+ ttlSeconds,
+ writeConcern,
+ readPreference);
+ }
+ }
+}
diff --git a/netty-socketio-core/src/main/java/com/socketio4j/socketio/store/mongo/MongoStoreFactory.java b/netty-socketio-core/src/main/java/com/socketio4j/socketio/store/mongo/MongoStoreFactory.java
new file mode 100644
index 00000000..e46b1b11
--- /dev/null
+++ b/netty-socketio-core/src/main/java/com/socketio4j/socketio/store/mongo/MongoStoreFactory.java
@@ -0,0 +1,124 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.store.mongo;
+
+import java.util.Map;
+import java.util.Objects;
+import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+
+import org.jetbrains.annotations.NotNull;
+import org.jetbrains.annotations.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.socketio4j.socketio.store.Store;
+import com.socketio4j.socketio.store.event.BaseStoreFactory;
+import com.socketio4j.socketio.store.event.EventStore;
+import com.socketio4j.socketio.store.event.EventStoreMode;
+import com.socketio4j.socketio.store.memory.MemoryStore;
+
+/**
+ * A {@code StoreFactory} implementation that provides session-scoped storage
+ * and MongoDB Change Streams based event distribution.
+ *
+ * Session data is stored in memory via {@link MemoryStore} for low-latency,
+ * non-blocking access, while event propagation across cluster nodes is handled
+ * asynchronously by {@link MongoEventStore} using MongoDB Change Streams.
+ */
+public class MongoStoreFactory extends BaseStoreFactory {
+
+ private static final Logger log = LoggerFactory.getLogger(MongoStoreFactory.class);
+
+ private final EventStore eventStore;
+
+ /**
+ * Creates a {@code MongoStoreFactory} using default {@link MongoEventStore}
+ * with {@link EventStoreMode#MULTI_CHANNEL} mode.
+ *
+ * @param mongoClient shared MongoDB client (must connect to a replica set)
+ * @param databaseName database to use for event collections
+ */
+ public MongoStoreFactory(@NotNull MongoClient mongoClient, @NotNull String databaseName) {
+ this(mongoClient, databaseName, EventStoreMode.MULTI_CHANNEL);
+ }
+
+ /**
+ * Creates a {@code MongoStoreFactory} using default {@link MongoEventStore}
+ * with the specified {@link EventStoreMode}.
+ *
+ * @param mongoClient shared MongoDB client (must connect to a replica set)
+ * @param databaseName database to use for event collections
+ * @param eventStoreMode SINGLE_CHANNEL or MULTI_CHANNEL mode
+ */
+ public MongoStoreFactory(@NotNull MongoClient mongoClient,
+ @NotNull String databaseName,
+ @Nullable EventStoreMode eventStoreMode) {
+ this(createDefaultEventStore(mongoClient, databaseName, eventStoreMode));
+ }
+
+ private static EventStore createDefaultEventStore(MongoClient mongoClient,
+ String databaseName,
+ EventStoreMode mode) {
+ EventStoreMode targetMode = EventStoreMode.MULTI_CHANNEL;
+ if (mode != null) {
+ targetMode = mode;
+ }
+ return new MongoEventStore.Builder(mongoClient, databaseName)
+ .eventStoreMode(targetMode)
+ .build();
+ }
+
+ /**
+ * Creates a {@code MongoStoreFactory} using the provided {@link EventStore}.
+ *
+ * @param eventStore non-null event store implementation
+ */
+ public MongoStoreFactory(@NotNull EventStore eventStore) {
+ this.eventStore = Objects.requireNonNull(eventStore, "eventStore cannot be null");
+ }
+
+ @Override
+ public Store createStore(UUID sessionId) {
+ return new MemoryStore();
+ }
+
+ @Override
+ public EventStore eventStore() {
+ return eventStore;
+ }
+
+ @Override
+ public Map createMap(String name) {
+ return new ConcurrentHashMap<>();
+ }
+
+ @Override
+ public void shutdown() {
+ try {
+ eventStore.shutdown();
+ } catch (Exception e) {
+ log.error("Failed to shut down Mongo event store", e);
+ }
+ }
+
+ @Override
+ public String toString() {
+ return getClass().getSimpleName() + " (memory session store, MongoDB change streams publish/subscribe)";
+ }
+}
diff --git a/netty-socketio-core/src/main/java11/module-info.java b/netty-socketio-core/src/main/java11/module-info.java
index 2f541538..0983f6d7 100644
--- a/netty-socketio-core/src/main/java11/module-info.java
+++ b/netty-socketio-core/src/main/java11/module-info.java
@@ -51,6 +51,8 @@
exports com.socketio4j.socketio.store.redis_reliable;
exports com.socketio4j.socketio.store.redis_stream;
exports com.socketio4j.socketio.store.kafka;
+ exports com.socketio4j.socketio.store.nats_pubsub;
+ exports com.socketio4j.socketio.store.mongo;
// ============================================================
// Reflective-only packages (not exported)
@@ -76,6 +78,10 @@
requires static redisson;
requires static io.nats.jnats;
requires static kafka.clients;
+ requires static org.mongodb.bson;
+ requires static org.mongodb.driver.core;
+ requires static org.mongodb.driver.reactivestreams;
+ requires static org.reactivestreams;
// ============================================================
// Optional Netty native transports — only if available
diff --git a/netty-socketio-core/src/test/java/com/socketio4j/socketio/integration/cluster/DistributedMongoClusterTest.java b/netty-socketio-core/src/test/java/com/socketio4j/socketio/integration/cluster/DistributedMongoClusterTest.java
new file mode 100644
index 00000000..93d64bea
--- /dev/null
+++ b/netty-socketio-core/src/test/java/com/socketio4j/socketio/integration/cluster/DistributedMongoClusterTest.java
@@ -0,0 +1,171 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.integration.cluster;
+
+import com.socketio4j.socketio.TestResourceCleanup;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.parallel.ResourceLock;
+
+import com.mongodb.reactivestreams.client.MongoClient;
+
+import com.socketio4j.socketio.Configuration;
+import com.socketio4j.socketio.SocketIOServer;
+import com.socketio4j.socketio.store.container.CustomizedMongoContainer;
+import com.socketio4j.socketio.store.event.EventStoreMode;
+import com.socketio4j.socketio.store.mongo.MongoEventStore;
+import com.socketio4j.socketio.store.mongo.MongoStoreFactory;
+
+/**
+ * Runs {@link DistributedCommonTest} against all MongoDB-backed cluster variants while sharing
+ * one MongoDB Testcontainer for maximum execution speed and zero container setup overhead.
+ */
+@ResourceLock("EMBEDDED_MONGO")
+public class DistributedMongoClusterTest {
+
+ @SuppressWarnings("resource")
+ static final CustomizedMongoContainer MONGO_CONTAINER = new CustomizedMongoContainer();
+
+ private static final String DB_NAME = "socketio_test";
+
+ @BeforeAll
+ static void startMongo() {
+ if (!MONGO_CONTAINER.isRunning()) {
+ for (int attempt = 1; attempt <= 3; attempt++) {
+ try {
+ MONGO_CONTAINER.start();
+ break;
+ } catch (Exception e) {
+ if (attempt == 3) throw new RuntimeException("Failed to start MongoDB container", e);
+ try {
+ Thread.sleep(500);
+ } catch (InterruptedException error) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Interrupted while starting MongoDB test container", error);
+ }
+ }
+ }
+ }
+ }
+
+ @AfterAll
+ static void stopMongo() {
+ TestResourceCleanup.runAll("MongoDB test container cleanup",
+ () -> { if (MONGO_CONTAINER != null && MONGO_CONTAINER.isRunning()) MONGO_CONTAINER.stop(); });
+ }
+
+ /**
+ * Starts one node backed by its own store, so the two nodes talk to each other only
+ * through change streams — never through a shared in-process client.
+ */
+ private static SocketIOServer startNode(MongoEventStore store, Configuration cfg) {
+ DistributedClusterIntegrationSupport.applyReuseListenAddress(cfg);
+ cfg.setHostname("127.0.0.1");
+ cfg.setPort(0);
+ cfg.setStoreFactory(new MongoStoreFactory(store));
+
+ SocketIOServer node = new SocketIOServer(cfg);
+ DistributedClusterIntegrationSupport.attachDefaultRoomListeners(node);
+ node.start();
+ return node;
+ }
+
+ @Nested
+ @TestInstance(TestInstance.Lifecycle.PER_CLASS)
+ class SingleChannelMemoryTest extends DistributedCommonTest {
+ private MongoClient mc1;
+ private MongoClient mc2;
+ private MongoEventStore store1;
+ private MongoEventStore store2;
+
+ @BeforeAll
+ void setupNodes() {
+ Configuration cfg1 = new Configuration();
+ mc1 = MONGO_CONTAINER.createClient();
+ store1 = new MongoEventStore.Builder(mc1, DB_NAME)
+ .eventStoreMode(EventStoreMode.SINGLE_CHANNEL)
+ .build();
+ node1 = startNode(store1, cfg1);
+ port1 = cfg1.getPort();
+
+ Configuration cfg2 = new Configuration();
+ mc2 = MONGO_CONTAINER.createClient();
+ store2 = new MongoEventStore.Builder(mc2, DB_NAME)
+ .eventStoreMode(EventStoreMode.SINGLE_CHANNEL)
+ .build();
+ node2 = startNode(store2, cfg2);
+ port2 = cfg2.getPort();
+ }
+
+ @AfterAll
+ void tearDownNodes() {
+ TestResourceCleanup.runAll("MongoDB cluster node cleanup",
+ () -> { if (node1 != null) node1.stop(); },
+ () -> { if (node2 != null) node2.stop(); },
+ // Stopping a node leaves its event store running, so close the change
+ // streams before the clients they are reading from.
+ () -> { if (store1 != null) store1.shutdown(); },
+ () -> { if (store2 != null) store2.shutdown(); },
+ () -> { if (mc1 != null) mc1.close(); },
+ () -> { if (mc2 != null) mc2.close(); });
+ }
+ }
+
+ @Nested
+ @TestInstance(TestInstance.Lifecycle.PER_CLASS)
+ class MultiChannelMemoryTest extends DistributedCommonTest {
+ private MongoClient mc1;
+ private MongoClient mc2;
+ private MongoEventStore store1;
+ private MongoEventStore store2;
+
+ @BeforeAll
+ void setupNodes() {
+ Configuration cfg1 = new Configuration();
+ mc1 = MONGO_CONTAINER.createClient();
+ store1 = new MongoEventStore.Builder(mc1, DB_NAME)
+ .eventStoreMode(EventStoreMode.MULTI_CHANNEL)
+ .build();
+ node1 = startNode(store1, cfg1);
+ port1 = cfg1.getPort();
+
+ Configuration cfg2 = new Configuration();
+ mc2 = MONGO_CONTAINER.createClient();
+ store2 = new MongoEventStore.Builder(mc2, DB_NAME)
+ .eventStoreMode(EventStoreMode.MULTI_CHANNEL)
+ .build();
+ node2 = startNode(store2, cfg2);
+ port2 = cfg2.getPort();
+ }
+
+ @AfterAll
+ void tearDownNodes() {
+ TestResourceCleanup.runAll("MongoDB cluster node cleanup",
+ () -> { if (node1 != null) node1.stop(); },
+ () -> { if (node2 != null) node2.stop(); },
+ // Stopping a node leaves its event store running, so close the change
+ // streams before the clients they are reading from.
+ () -> { if (store1 != null) store1.shutdown(); },
+ () -> { if (store2 != null) store2.shutdown(); },
+ () -> { if (mc1 != null) mc1.close(); },
+ () -> { if (mc2 != null) mc2.close(); });
+ }
+ }
+}
diff --git a/netty-socketio-core/src/test/java/com/socketio4j/socketio/integration/interop/DistributedMongoJsClientInteropTest.java b/netty-socketio-core/src/test/java/com/socketio4j/socketio/integration/interop/DistributedMongoJsClientInteropTest.java
new file mode 100644
index 00000000..39361f9c
--- /dev/null
+++ b/netty-socketio-core/src/test/java/com/socketio4j/socketio/integration/interop/DistributedMongoJsClientInteropTest.java
@@ -0,0 +1,122 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.integration.interop;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.parallel.ResourceLock;
+
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.socketio4j.socketio.Configuration;
+import com.socketio4j.socketio.SocketIOServer;
+import com.socketio4j.socketio.TestResourceCleanup;
+import com.socketio4j.socketio.integration.cluster.DistributedClusterIntegrationSupport;
+import com.socketio4j.socketio.store.container.CustomizedMongoContainer;
+import com.socketio4j.socketio.store.event.EventStoreMode;
+import com.socketio4j.socketio.store.memory.MemoryStoreFactory;
+import com.socketio4j.socketio.store.mongo.MongoEventStore;
+
+/**
+ * Multi-Node JS Client Interoperability Test Suite backed by MongoDB Change Streams.
+ */
+@ResourceLock("EMBEDDED_MONGO")
+@DisplayName("Multi-Node Official JS Client Interoperability Suite (MongoDB Change Streams)")
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+public class DistributedMongoJsClientInteropTest extends AbstractDistributedJsClientInteropTest {
+
+ @SuppressWarnings("resource")
+ private static final CustomizedMongoContainer MONGO_CONTAINER = new CustomizedMongoContainer();
+
+ private static final String DB_NAME = "socketio_js_interop";
+
+ private MongoClient mc1;
+ private MongoClient mc2;
+ private MongoEventStore store1;
+ private MongoEventStore store2;
+
+ @BeforeAll
+ @Override
+ public void setupCluster() throws Exception {
+ if (!MONGO_CONTAINER.isRunning()) {
+ for (int attempt = 1; attempt <= 3; attempt++) {
+ try {
+ MONGO_CONTAINER.start();
+ break;
+ } catch (Exception e) {
+ if (attempt == 3) {
+ throw new RuntimeException("Failed to start MongoDB container", e);
+ }
+ try {
+ Thread.sleep(500);
+ } catch (InterruptedException error) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Interrupted while starting MongoDB test container", error);
+ }
+ }
+ }
+ }
+
+ // Server 1
+ mc1 = MONGO_CONTAINER.createClient();
+ store1 = new MongoEventStore.Builder(mc1, DB_NAME)
+ .eventStoreMode(EventStoreMode.MULTI_CHANNEL)
+ .collectionPrefix("interop_events_")
+ .build();
+ Configuration cfg1 = new Configuration();
+ DistributedClusterIntegrationSupport.applyReuseListenAddress(cfg1);
+ cfg1.setHostname("127.0.0.1");
+ cfg1.setPort(DistributedClusterIntegrationSupport.findAvailablePort());
+ cfg1.setStoreFactory(new MemoryStoreFactory(store1));
+ node1 = new SocketIOServer(cfg1);
+ attachDefaultRoomListeners(node1);
+ node1.start();
+ port1 = cfg1.getPort();
+
+ // Server 2
+ mc2 = MONGO_CONTAINER.createClient();
+ store2 = new MongoEventStore.Builder(mc2, DB_NAME)
+ .eventStoreMode(EventStoreMode.MULTI_CHANNEL)
+ .collectionPrefix("interop_events_")
+ .build();
+ Configuration cfg2 = new Configuration();
+ DistributedClusterIntegrationSupport.applyReuseListenAddress(cfg2);
+ cfg2.setHostname("127.0.0.1");
+ cfg2.setPort(DistributedClusterIntegrationSupport.findAvailablePort());
+ cfg2.setStoreFactory(new MemoryStoreFactory(store2));
+ node2 = new SocketIOServer(cfg2);
+ attachDefaultRoomListeners(node2);
+ node2.start();
+ port2 = cfg2.getPort();
+
+ initJsScript();
+ }
+
+ @AfterAll
+ @Override
+ public void teardownCluster() {
+ TestResourceCleanup.runAll("MongoDB distributed interop cleanup",
+ () -> { if (node1 != null) node1.stop(); },
+ () -> { if (node2 != null) node2.stop(); },
+ () -> { if (store1 != null) store1.shutdown(); },
+ () -> { if (store2 != null) store2.shutdown(); },
+ () -> { if (mc1 != null) mc1.close(); },
+ () -> { if (mc2 != null) mc2.close(); },
+ () -> { if (MONGO_CONTAINER != null && MONGO_CONTAINER.isRunning()) MONGO_CONTAINER.stop(); });
+ }
+}
diff --git a/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/container/CustomizedMongoContainer.java b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/container/CustomizedMongoContainer.java
new file mode 100644
index 00000000..42b7063f
--- /dev/null
+++ b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/container/CustomizedMongoContainer.java
@@ -0,0 +1,100 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.store.container;
+
+import java.time.Duration;
+import java.util.concurrent.TimeUnit;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.utility.DockerImageName;
+
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.mongodb.reactivestreams.client.MongoClients;
+
+/**
+ * CI-safe customized MongoDB container for socketio4j tests.
+ *
+ * Uses a single-node replica set so that change streams are available.
+ */
+public class CustomizedMongoContainer extends GenericContainer {
+
+ private static final Logger log =
+ LoggerFactory.getLogger(CustomizedMongoContainer.class);
+
+ private static final int MONGO_PORT = 27017;
+
+ public CustomizedMongoContainer() {
+ super(DockerImageName.parse("mongo:7.0.40"));
+
+ withExposedPorts(MONGO_PORT);
+ withCommand("--replSet", "rs0");
+ withReuse(false);
+ withStartupAttempts(3);
+ withStartupTimeout(Duration.ofMinutes(2));
+ }
+
+ @Override
+ public void start() {
+ super.start();
+ initReplicaSet();
+ log.info("MongoDB replica set ready at {}", getConnectionString());
+ }
+
+ private void initReplicaSet() {
+ try {
+ ExecResult result = execInContainer(
+ "mongosh", "--eval",
+ "rs.initiate({_id: 'rs0', members: [{_id: 0, host: 'localhost:" + MONGO_PORT + "'}]})"
+ );
+ log.debug("rs.initiate output: {}", result.getStdout());
+
+ long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(30);
+ while (System.currentTimeMillis() < deadline) {
+ // rs.initiate() only starts an election, so rs.status().ok == 1 does not
+ // yet mean this member can accept writes — wait for it to be PRIMARY.
+ // --quiet keeps banners out of stdout, hello() throws until the set is
+ // initiated (swallowed below), and the result is compared exactly: a
+ // substring match would accept any output that merely contains the value.
+ ExecResult status = execInContainer(
+ "mongosh", "--quiet", "--eval",
+ "try { db.hello().isWritablePrimary } catch (e) { false }"
+ );
+ String out = status.getStdout().trim();
+ if ("true".equals(out)) {
+ return;
+ }
+ Thread.sleep(500);
+ }
+ throw new RuntimeException("MongoDB replica set not ready after 30s");
+ } catch (RuntimeException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new RuntimeException("Failed to initialize MongoDB replica set", e);
+ }
+ }
+
+ public String getConnectionString() {
+ return "mongodb://" + getHost() + ":" + getMappedPort(MONGO_PORT)
+ + "/?replicaSet=rs0&directConnection=true";
+ }
+
+ public MongoClient createClient() {
+ return MongoClients.create(getConnectionString());
+ }
+}
diff --git a/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/event/MongoPubSubEventStoreTest.java b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/event/MongoPubSubEventStoreTest.java
new file mode 100644
index 00000000..c6e89297
--- /dev/null
+++ b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/event/MongoPubSubEventStoreTest.java
@@ -0,0 +1,164 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.store.event;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.parallel.ResourceLock;
+import org.testcontainers.containers.GenericContainer;
+
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.socketio4j.socketio.protocol.Packet;
+import com.socketio4j.socketio.protocol.PacketType;
+import com.socketio4j.socketio.store.container.CustomizedMongoContainer;
+import com.socketio4j.socketio.store.mongo.MongoEventStore;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Test class for MongoEventStore using testcontainers.
+ */
+@ResourceLock("EMBEDDED_MONGO")
+public class MongoPubSubEventStoreTest extends AbstractEventStoreTest {
+
+ private static final String DB_NAME = "socketio_event_store_test";
+
+ private final AtomicInteger testCounter = new AtomicInteger();
+ private MongoClient sharedClient;
+ private String currentPrefix;
+
+ @BeforeEach
+ @Override
+ public void setUp() throws Exception {
+ currentPrefix = "test_" + testCounter.incrementAndGet() + "_";
+ super.setUp();
+ }
+
+ @Override
+ protected GenericContainer> createContainer() {
+ return new CustomizedMongoContainer().withReuse(false);
+ }
+
+ private synchronized MongoClient getSharedClient() {
+ if (sharedClient == null) {
+ CustomizedMongoContainer mongoContainer = (CustomizedMongoContainer) container;
+ sharedClient = mongoContainer.createClient();
+ }
+ return sharedClient;
+ }
+
+ @Override
+ protected EventStore createEventStore(Long nodeId) throws Exception {
+ return new MongoEventStore.Builder(getSharedClient(), DB_NAME)
+ .nodeId(nodeId)
+ .collectionPrefix(currentPrefix)
+ .eventStoreMode(EventStoreMode.MULTI_CHANNEL)
+ .build();
+ }
+
+ @Override
+ protected void closeClients() {
+ // Shared client remains open across tests to prevent reconnection overhead and cursor race
+ }
+
+ @AfterAll
+ @Override
+ public void stopContainer() {
+ if (sharedClient != null) {
+ try {
+ sharedClient.close();
+ } catch (Exception ignored) {
+ }
+ sharedClient = null;
+ }
+ super.stopContainer();
+ }
+
+ @Test
+ public void testSingleChannelPubSub() throws Exception {
+ String singlePrefix = "single_" + testCounter.incrementAndGet() + "_";
+
+ MongoEventStore singlePubStore = new MongoEventStore.Builder(getSharedClient(), DB_NAME)
+ .nodeId(300L)
+ .eventStoreMode(EventStoreMode.SINGLE_CHANNEL)
+ .collectionPrefix(singlePrefix)
+ .build();
+
+ MongoEventStore singleSubStore = new MongoEventStore.Builder(getSharedClient(), DB_NAME)
+ .nodeId(301L)
+ .eventStoreMode(EventStoreMode.SINGLE_CHANNEL)
+ .collectionPrefix(singlePrefix)
+ .build();
+
+ try {
+ CountDownLatch latch = new CountDownLatch(1);
+ AtomicReference receivedRef = new AtomicReference<>();
+
+ singleSubStore.subscribe(
+ EventType.ALL_SINGLE_CHANNEL,
+ message -> {
+ if (message instanceof DispatchMessage) {
+ receivedRef.set((DispatchMessage) message);
+ latch.countDown();
+ }
+ },
+ DispatchMessage.class
+ );
+
+ Packet packet = new Packet(PacketType.MESSAGE);
+ packet.setSubType(PacketType.EVENT);
+ packet.setName("single-channel-event");
+ packet.setNsp("/");
+ packet.setData("hello-single-channel");
+
+ DispatchMessage outgoing = new DispatchMessage("roomA", packet, "/");
+ outgoing.setNodeId(300L);
+
+ singlePubStore.publish(EventType.DISPATCH, outgoing);
+
+ assertTrue(latch.await(5, TimeUnit.SECONDS), "Message should be received in SINGLE_CHANNEL mode");
+ DispatchMessage received = receivedRef.get();
+ assertNotNull(received);
+ assertEquals("roomA", received.getRoom());
+ assertEquals(300L, received.getNodeId());
+ assertEquals("single-channel-event", received.getPacket().getName());
+ assertEquals("hello-single-channel", received.getPacket().getData());
+
+ // Unsubscribe test
+ singleSubStore.unsubscribe(EventType.ALL_SINGLE_CHANNEL);
+ CountDownLatch unsubLatch = new CountDownLatch(1);
+ AtomicReference unsubReceived = new AtomicReference<>();
+
+ singlePubStore.publish(EventType.DISPATCH, outgoing);
+ assertFalse(unsubLatch.await(2, TimeUnit.SECONDS), "No message after unsubscribe");
+ assertNull(unsubReceived.get());
+ } finally {
+ singlePubStore.shutdown();
+ singleSubStore.shutdown();
+ }
+ }
+}
diff --git a/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/mongo/MongoEventStoreTest.java b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/mongo/MongoEventStoreTest.java
new file mode 100644
index 00000000..05ab3544
--- /dev/null
+++ b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/mongo/MongoEventStoreTest.java
@@ -0,0 +1,489 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.store.mongo;
+
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.bson.BsonDocument;
+import org.bson.BsonTimestamp;
+import org.bson.Document;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.reactivestreams.Publisher;
+import org.reactivestreams.Subscription;
+
+import com.mongodb.MongoCommandException;
+import com.mongodb.MongoNamespace;
+import com.mongodb.ReadPreference;
+import com.mongodb.WriteConcern;
+import com.mongodb.client.model.changestream.ChangeStreamDocument;
+import com.mongodb.client.result.InsertOneResult;
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.mongodb.reactivestreams.client.MongoCollection;
+import com.mongodb.reactivestreams.client.MongoDatabase;
+import com.socketio4j.socketio.protocol.Packet;
+import com.socketio4j.socketio.protocol.PacketType;
+import com.socketio4j.socketio.store.event.DispatchMessage;
+import com.socketio4j.socketio.store.event.EventListener;
+import com.socketio4j.socketio.store.event.EventMessage;
+import com.socketio4j.socketio.store.event.EventMessageJsonSupport;
+import com.socketio4j.socketio.store.event.EventStoreMode;
+import com.socketio4j.socketio.store.event.EventStoreType;
+import com.socketio4j.socketio.store.event.EventType;
+import com.socketio4j.socketio.store.event.PublishMode;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class MongoEventStoreTest {
+
+ private MongoClient mongoClient;
+ private MongoDatabase mongoDatabase;
+ private MongoCollection mongoCollection;
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ void setUp() {
+ mongoClient = mock(MongoClient.class);
+ mongoDatabase = mock(MongoDatabase.class);
+ mongoCollection = mock(MongoCollection.class);
+ when(mongoClient.getDatabase(anyString())).thenReturn(mongoDatabase);
+ }
+
+ @Test
+ void testBuilderDefaultsAndCustomValues() {
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(42L)
+ .collectionPrefix("custom_events_")
+ .ttlSeconds(120)
+ .eventStoreMode(EventStoreMode.SINGLE_CHANNEL)
+ .writeConcern(WriteConcern.W1)
+ .readPreference(ReadPreference.secondaryPreferred())
+ .build();
+
+ assertEquals(EventStoreMode.SINGLE_CHANNEL, store.getEventStoreMode());
+ assertEquals(EventStoreType.PUBSUB, store.getEventStoreType());
+ assertEquals(PublishMode.UNRELIABLE, store.getPublishMode());
+ assertEquals("custom_events_", store.getCollectionPrefix());
+ assertEquals(120, store.getTtlSeconds());
+ assertEquals(WriteConcern.W1, store.getWriteConcern());
+ assertEquals(ReadPreference.secondaryPreferred(), store.getReadPreference());
+ }
+
+ @Test
+ void testValidateSubscribeModes() {
+ MongoEventStore singleChannelStore = new MongoEventStore.Builder(mongoClient, "testdb")
+ .eventStoreMode(EventStoreMode.SINGLE_CHANNEL)
+ .build();
+
+ assertThrows(UnsupportedOperationException.class,
+ () -> singleChannelStore.subscribe0(EventType.DISPATCH, msg -> {
+ }, EventMessage.class));
+
+ MongoEventStore multiChannelStore = new MongoEventStore.Builder(mongoClient, "testdb")
+ .eventStoreMode(EventStoreMode.MULTI_CHANNEL)
+ .build();
+
+ assertThrows(UnsupportedOperationException.class,
+ () -> multiChannelStore.subscribe0(EventType.ALL_SINGLE_CHANNEL, msg -> {
+ }, EventMessage.class));
+ }
+
+ @Test
+ void testWatcherHandleLifecycleAndCancel() {
+ BsonTimestamp startAt = new BsonTimestamp(100, 1);
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(startAt);
+
+ assertEquals(startAt, handle.startAt());
+ assertNull(handle.resumeToken());
+
+ BsonDocument token = new BsonDocument();
+ handle.setResumeToken(token);
+ assertEquals(token, handle.resumeToken());
+
+ handle.clearResumePoint();
+ assertNull(handle.resumeToken());
+ assertNull(handle.startAt());
+
+ AtomicBoolean cancelled = new AtomicBoolean(false);
+ Subscription subscription = new Subscription() {
+ @Override
+ public void request(long n) {
+ }
+
+ @Override
+ public void cancel() {
+ cancelled.set(true);
+ }
+ };
+
+ handle.setSubscription(subscription);
+ assertEquals(subscription, handle.getSubscription());
+
+ handle.stop();
+ assertTrue(cancelled.get());
+ assertTrue(handle.stopped.get());
+ assertNull(handle.getSubscription());
+ }
+
+ @Test
+ void testFastShutdownWithoutArtificialDelay() {
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .build();
+
+ long start = System.currentTimeMillis();
+ store.shutdown0();
+ long elapsed = System.currentTimeMillis() - start;
+
+ // Shutdown should complete promptly, well below the previous 1000ms artificial
+ // delay
+ assertTrue(elapsed < 800, "Shutdown took " + elapsed + "ms, expected under 800ms");
+ }
+
+ @Test
+ void testPublishValidation() {
+ when(mongoDatabase.getCollection(anyString())).thenReturn(mongoCollection);
+ Publisher dummyPublisher = subscriber -> subscriber
+ .onSubscribe(mock(Subscription.class));
+ when(mongoCollection.insertOne(any(Document.class))).thenReturn(dummyPublisher);
+
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(10L)
+ .build();
+
+ EventMessage msg = new EventMessage() {
+ @Override
+ public String getType() {
+ return "DISPATCH";
+ }
+ };
+
+ store.publish0(EventType.DISPATCH, msg);
+ assertEquals(10L, msg.getNodeId());
+ }
+
+ @Test
+ void testOplogHistoryLostClearsResumePoint() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(1L)
+ .build();
+
+ BsonTimestamp startAt = new BsonTimestamp(100, 1);
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(startAt);
+ handle.setResumeToken(new BsonDocument());
+
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, msg -> {
+ }, DispatchMessage.class);
+
+ // Simulate MongoCommandException with error code 286 (ChangeStreamHistoryLost)
+ MongoCommandException historyLostException = mock(MongoCommandException.class);
+ when(historyLostException.getErrorCode()).thenReturn(286);
+ when(historyLostException.getMessage())
+ .thenReturn("ChangeStreamHistoryLost: resume point no longer in oplog");
+
+ subscriber.onError(historyLostException);
+
+ // Resume point must be cleared so the next watch() starts at current time
+ assertNull(handle.resumeToken());
+ assertNull(handle.startAt());
+
+ store.shutdown0();
+ }
+
+ @Test
+ void testTransientErrorPreservesResumePoint() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(1L)
+ .build();
+
+ BsonTimestamp startAt = new BsonTimestamp(100, 1);
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(startAt);
+ BsonDocument token = new BsonDocument();
+ handle.setResumeToken(token);
+
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, msg -> {
+ }, DispatchMessage.class);
+
+ // Simulate transient network exception
+ subscriber.onError(new RuntimeException("Transient connection reset"));
+
+ // Resume token must be preserved to resume after the last delivered event
+ assertEquals(token, handle.resumeToken());
+
+ store.shutdown0();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void testChangeSubscriberDropsSelfPublishedEvents() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(50L)
+ .build();
+
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(new BsonTimestamp(1, 1));
+ EventListener listener = mock(EventListener.class);
+
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, listener, DispatchMessage.class);
+
+ ChangeStreamDocument change = mock(ChangeStreamDocument.class);
+ Document doc = new Document()
+ .append("nodeId", 50L) // same node id
+ .append("eventType", "DISPATCH")
+ .append("payload",
+ "{\"room\":\"r\",\"namespace\":\"/\",\"packet\":{\"type\":2,\"data\":\"hi\"}}");
+
+ when(change.getFullDocument()).thenReturn(doc);
+ when(change.getResumeToken()).thenReturn(new BsonDocument());
+
+ subscriber.onNext(change);
+
+ // Must drop without calling listener
+ verify(listener, never()).onMessage(any());
+
+ store.shutdown0();
+ }
+
+ @Test
+ void testChangeSubscriberDeliversAndPreservesBinaryPayload() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(50L)
+ .build();
+
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(new BsonTimestamp(1, 1));
+ AtomicReference receivedRef = new AtomicReference<>();
+
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, receivedRef::set, DispatchMessage.class);
+
+ byte[] rawBytes = new byte[] { 0x10, 0x20, 0x30, 0x40, 0x50 };
+ Packet packet = new Packet(PacketType.BINARY_EVENT);
+ packet.setName("bin-event");
+ packet.setNsp("/binary");
+ packet.setData(rawBytes);
+
+ DispatchMessage msg = new DispatchMessage("binRoom", packet, "/binary");
+ msg.setNodeId(99L); // remote node
+
+ // Use store's serializer logic
+ String json;
+ try {
+ json = EventMessageJsonSupport.createObjectMapper().writeValueAsString(msg);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+
+ Document doc = new Document()
+ .append("nodeId", 99L)
+ .append("eventType", "DISPATCH")
+ .append("payload", json);
+
+ @SuppressWarnings("unchecked")
+ ChangeStreamDocument change = mock(ChangeStreamDocument.class);
+ when(change.getFullDocument()).thenReturn(doc);
+ when(change.getResumeToken()).thenReturn(new BsonDocument());
+
+ subscriber.onNext(change);
+
+ DispatchMessage received = receivedRef.get();
+ assertNotNull(received);
+ assertEquals(99L, received.getNodeId());
+ assertEquals("binRoom", received.getRoom());
+ assertNotNull(received.getPacket());
+ assertArrayEquals(rawBytes, (byte[]) received.getPacket().getData());
+
+ store.shutdown0();
+ }
+
+ @Test
+ void testPublishAndSubscribeAfterShutdown() {
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(1L)
+ .build();
+
+ store.shutdown0();
+
+ // Second shutdown must be safe and idempotent
+ assertDoesNotThrow(store::shutdown0);
+
+ // Publish after shutdown should safely no-op without exception
+ assertDoesNotThrow(() -> store.publish0(EventType.DISPATCH, new DispatchMessage()));
+
+ // Subscribe after shutdown must be rejected
+ assertThrows(IllegalStateException.class, () -> store.subscribe0(EventType.DISPATCH, msg -> {
+ }, DispatchMessage.class));
+ }
+
+ @Test
+ void testReconnectAttemptsPersistAcrossReopens() {
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(new BsonTimestamp(1, 1));
+
+ assertEquals(1, handle.incrementReconnectAttempts());
+ assertEquals(2, handle.incrementReconnectAttempts());
+ assertEquals(3, handle.incrementReconnectAttempts());
+
+ handle.resetReconnectAttempts();
+ assertEquals(1, handle.incrementReconnectAttempts());
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void testStoppedWatcherHandleDropsInFlightOnNext() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(50L)
+ .build();
+
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(new BsonTimestamp(1, 1));
+ EventListener listener = mock(EventListener.class);
+
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, listener, DispatchMessage.class);
+
+ // Stop handle before event arrives
+ handle.stop();
+
+ ChangeStreamDocument change = mock(ChangeStreamDocument.class);
+ Document doc = new Document()
+ .append("nodeId", 99L)
+ .append("eventType", "DISPATCH")
+ .append("payload",
+ "{\"room\":\"r\",\"namespace\":\"/\",\"packet\":{\"type\":2,\"data\":\"hi\"}}");
+
+ when(change.getFullDocument()).thenReturn(doc);
+ when(change.getResumeToken()).thenReturn(new BsonDocument());
+
+ subscriber.onNext(change);
+
+ // Must drop without calling listener because handle is stopped
+ verify(listener, never()).onMessage(any());
+
+ store.shutdown0();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void testChangeSubscriberCatchesExceptionFromListener() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(50L)
+ .build();
+
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(new BsonTimestamp(1, 1));
+ EventListener listener = mock(EventListener.class);
+ org.mockito.Mockito.doThrow(new RuntimeException("Error in listener")).when(listener).onMessage(any());
+
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, listener, DispatchMessage.class);
+
+ ChangeStreamDocument change = mock(ChangeStreamDocument.class);
+ Document doc = new Document()
+ .append("nodeId", 99L)
+ .append("eventType", "DISPATCH")
+ .append("payload",
+ "{\"room\":\"r\",\"namespace\":\"/\",\"packet\":{\"type\":2,\"data\":\"hi\"}}");
+
+ when(change.getFullDocument()).thenReturn(doc);
+ when(change.getResumeToken()).thenReturn(new BsonDocument());
+
+ // Should safely catch Throwable without propagating or failing
+ assertDoesNotThrow(() -> subscriber.onNext(change));
+
+ store.shutdown0();
+ }
+
+ @Test
+ void testStandaloneErrorDetectionInChangeSubscriber() {
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(1L)
+ .build();
+
+ MongoEventStore.WatcherHandle handle = new MongoEventStore.WatcherHandle(new BsonTimestamp(1, 1));
+ MongoEventStore.ChangeSubscriber subscriber = store.new ChangeSubscriber<>(
+ mongoCollection, EventType.DISPATCH, handle, msg -> {
+ }, DispatchMessage.class);
+
+ MongoCommandException standaloneEx = mock(MongoCommandException.class);
+ when(standaloneEx.getErrorCode()).thenReturn(40573);
+ when(standaloneEx.getMessage())
+ .thenReturn("The $changeStream stage is only supported on replica sets");
+
+ assertDoesNotThrow(() -> subscriber.onError(standaloneEx));
+
+ store.shutdown0();
+ }
+
+ @Test
+ void testShutdownAsyncCompletesPromptly() {
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .build();
+
+ java.util.concurrent.CompletableFuture future = store.shutdownAsync();
+ assertNotNull(future);
+ assertDoesNotThrow(() -> future.get(2, java.util.concurrent.TimeUnit.SECONDS));
+ assertTrue(future.isDone());
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void testSubscribe0IsNonBlocking() {
+ when(mongoDatabase.getCollection(anyString())).thenReturn(mongoCollection);
+ when(mongoCollection.getNamespace()).thenReturn(new MongoNamespace("testdb.events"));
+
+ // Mock ping publisher that never completes (simulating a slow network)
+ Publisher slowPingPublisher = subscriber -> {
+ // intentionally do not invoke onNext or onComplete
+ };
+ when(mongoDatabase.runCommand(any(Document.class))).thenReturn(slowPingPublisher);
+
+ Publisher indexPublisher = subscriber -> subscriber.onSubscribe(mock(Subscription.class));
+ when(mongoCollection.createIndex(any(org.bson.conversions.Bson.class), any(com.mongodb.client.model.IndexOptions.class)))
+ .thenReturn(indexPublisher);
+
+ MongoEventStore store = new MongoEventStore.Builder(mongoClient, "testdb")
+ .nodeId(1L)
+ .build();
+
+ long start = System.currentTimeMillis();
+ // subscribe0 must return immediately without waiting for ping or index
+ store.subscribe0(EventType.DISPATCH, msg -> {
+ }, DispatchMessage.class);
+ long elapsed = System.currentTimeMillis() - start;
+
+ assertTrue(elapsed < 200, "subscribe0 blocked for " + elapsed + "ms, expected non-blocking (<200ms)");
+ store.shutdown0();
+ }
+}
diff --git a/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/mongo/MongoStoreFactoryTest.java b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/mongo/MongoStoreFactoryTest.java
new file mode 100644
index 00000000..f98135e8
--- /dev/null
+++ b/netty-socketio-core/src/test/java/com/socketio4j/socketio/store/mongo/MongoStoreFactoryTest.java
@@ -0,0 +1,105 @@
+/**
+ * Copyright (c) 2025 The Socketio4j Project
+ * Parent project : Copyright (c) 2012-2025 Nikita Koksharov
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.socketio4j.socketio.store.mongo;
+
+import java.util.Map;
+import java.util.UUID;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import com.mongodb.reactivestreams.client.MongoClient;
+import com.mongodb.reactivestreams.client.MongoDatabase;
+import com.socketio4j.socketio.store.Store;
+import com.socketio4j.socketio.store.event.EventStore;
+import com.socketio4j.socketio.store.event.EventStoreMode;
+import com.socketio4j.socketio.store.memory.MemoryStore;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class MongoStoreFactoryTest {
+
+ private MongoClient mongoClient;
+ private MongoDatabase mongoDatabase;
+
+ @BeforeEach
+ void setUp() {
+ mongoClient = mock(MongoClient.class);
+ mongoDatabase = mock(MongoDatabase.class);
+ when(mongoClient.getDatabase(anyString())).thenReturn(mongoDatabase);
+ }
+
+ @Test
+ void testConstructorWithDatabaseName() {
+ MongoStoreFactory factory = new MongoStoreFactory(mongoClient, "testdb");
+ assertNotNull(factory.eventStore());
+ assertEquals(EventStoreMode.MULTI_CHANNEL, factory.eventStore().getEventStoreMode());
+ }
+
+ @Test
+ void testConstructorWithMode() {
+ MongoStoreFactory factory = new MongoStoreFactory(mongoClient, "testdb", EventStoreMode.SINGLE_CHANNEL);
+ assertNotNull(factory.eventStore());
+ assertEquals(EventStoreMode.SINGLE_CHANNEL, factory.eventStore().getEventStoreMode());
+ }
+
+ @Test
+ void testCreateStoreReturnsMemoryStore() {
+ MongoStoreFactory factory = new MongoStoreFactory(mongoClient, "testdb");
+ UUID sessionId = UUID.randomUUID();
+ Store store = factory.createStore(sessionId);
+ assertNotNull(store);
+ assertInstanceOf(MemoryStore.class, store);
+
+ store.set("key", "value");
+ assertEquals("value", store.get("key"));
+ assertTrue(store.has("key"));
+ store.del("key");
+ store.destroy();
+ }
+
+ @Test
+ void testCreateMap() {
+ MongoStoreFactory factory = new MongoStoreFactory(mongoClient, "testdb");
+ Map map = factory.createMap("myMap");
+ assertNotNull(map);
+ map.put("k", "v");
+ assertEquals("v", map.get("k"));
+ }
+
+ @Test
+ void testShutdownDelegatesToEventStore() {
+ EventStore mockEventStore = mock(EventStore.class);
+ MongoStoreFactory factory = new MongoStoreFactory(mockEventStore);
+
+ factory.shutdown();
+ verify(mockEventStore).shutdown();
+ }
+
+ @Test
+ void testToString() {
+ MongoStoreFactory factory = new MongoStoreFactory(mongoClient, "testdb");
+ assertTrue(factory.toString().contains("MongoStoreFactory"));
+ }
+}
diff --git a/pom.xml b/pom.xml
index 1209a3ca..2ef73630 100644
--- a/pom.xml
+++ b/pom.xml
@@ -86,6 +86,7 @@
4.3.0
3.20.0
2.25.3
+ 5.5.0
1.17.0
1.6.0
1.10.3
@@ -352,6 +353,12 @@
${kafka.version}
provided
+
+ org.mongodb
+ mongodb-driver-reactivestreams
+ ${mongodb.version}
+ provided
+