diff --git a/netty-socketio-core/pom.xml b/netty-socketio-core/pom.xml index 1e1da25f..471d2cc8 100644 --- a/netty-socketio-core/pom.xml +++ b/netty-socketio-core/pom.xml @@ -143,6 +143,11 @@ + + org.mongodb + mongodb-driver-reactivestreams + provided + @@ -285,6 +290,7 @@ **/store/RedissonStoreTest.java **/store/event/HazelcastRingBufferEventStoreTest.java **/store/event/RedisPubSubEventStoreTest.java + **/store/event/MongoPubSubEventStoreTest.java @@ -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 +