Skip to content

Commit 19cf880

Browse files
committed
Fix AbstractOplogTest
1 parent 09b11f7 commit 19cf880

8 files changed

Lines changed: 100 additions & 52 deletions

File tree

‎gradle.lockfile‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@ io.netty:netty-tcnative-classes:2.0.48.Final=default
2323
io.netty:netty-transport-native-unix-common:4.1.80.Final=allCodeCoverageReportClassDirectories,runtimeClasspath,testRuntimeClasspath
2424
io.netty:netty-transport:4.1.74.Final=default
2525
io.netty:netty-transport:4.1.80.Final=allCodeCoverageReportClassDirectories,runtimeClasspath,testRuntimeClasspath
26-
io.reactivex.rxjava2:rxjava:2.2.21=allCodeCoverageReportClassDirectories
2726
net.bytebuddy:byte-buddy-agent:1.12.13=allCodeCoverageReportClassDirectories
2827
net.bytebuddy:byte-buddy:1.12.13=allCodeCoverageReportClassDirectories
2928
org.assertj:assertj-core:3.23.1=allCodeCoverageReportClassDirectories

‎h2-backend/gradle.lockfile‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@ io.netty:netty-tcnative-classes:2.0.48.Final=default
1919
io.netty:netty-transport-native-unix-common:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
2020
io.netty:netty-transport:4.1.74.Final=default
2121
io.netty:netty-transport:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
22-
io.reactivex.rxjava2:rxjava:2.2.21=testCompileClasspath,testRuntimeClasspath
2322
net.bytebuddy:byte-buddy-agent:1.12.13=testCompileClasspath,testRuntimeClasspath
2423
net.bytebuddy:byte-buddy:1.12.13=testCompileClasspath,testRuntimeClasspath
2524
org.apiguardian:apiguardian-api:1.1.2=testCompileClasspath

‎memory-backend/gradle.lockfile‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@ io.netty:netty-tcnative-classes:2.0.48.Final=default
1717
io.netty:netty-transport-native-unix-common:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
1818
io.netty:netty-transport:4.1.74.Final=default
1919
io.netty:netty-transport:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
20-
io.reactivex.rxjava2:rxjava:2.2.21=testCompileClasspath,testRuntimeClasspath
2120
net.bytebuddy:byte-buddy-agent:1.12.13=testCompileClasspath,testRuntimeClasspath
2221
net.bytebuddy:byte-buddy:1.12.13=testCompileClasspath,testRuntimeClasspath
2322
org.apiguardian:apiguardian-api:1.1.2=testCompileClasspath

‎postgresql-backend/gradle.lockfile‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,6 @@ io.netty:netty-tcnative-classes:2.0.48.Final=default
3131
io.netty:netty-transport-native-unix-common:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
3232
io.netty:netty-transport:4.1.74.Final=default
3333
io.netty:netty-transport:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
34-
io.reactivex.rxjava2:rxjava:2.2.21=testCompileClasspath,testRuntimeClasspath
3534
junit:junit:4.13.2=testCompileClasspath,testRuntimeClasspath
3635
net.bytebuddy:byte-buddy-agent:1.12.13=testCompileClasspath,testRuntimeClasspath
3736
net.bytebuddy:byte-buddy:1.12.13=testCompileClasspath,testRuntimeClasspath

‎test-common/build.gradle‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,6 @@ dependencies {
1616
api "org.mockito:mockito-core:latest.release"
1717
api "org.mockito:mockito-junit-jupiter:latest.release"
1818

19-
api group: 'io.reactivex.rxjava2', name: 'rxjava', version: 'latest.release'
20-
2119
runtimeOnly group: 'ch.qos.logback', name: 'logback-classic', version: '[1.3.0, 1.4.0)'
2220

2321
testImplementation "org.testcontainers:testcontainers:latest.release"

‎test-common/gradle.lockfile‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ io.netty:netty-tcnative-classes:2.0.48.Final=default
2323
io.netty:netty-transport-native-unix-common:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
2424
io.netty:netty-transport:4.1.74.Final=default
2525
io.netty:netty-transport:4.1.80.Final=runtimeClasspath,testRuntimeClasspath
26-
io.reactivex.rxjava2:rxjava:2.2.21=compileClasspath,default,runtimeClasspath,testCompileClasspath,testRuntimeClasspath
26+
io.reactivex.rxjava2:rxjava:2.2.21=default
2727
junit:junit:4.13.2=testCompileClasspath,testRuntimeClasspath
2828
net.bytebuddy:byte-buddy-agent:1.12.13=compileClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath
2929
net.bytebuddy:byte-buddy-agent:1.12.7=default

‎test-common/src/main/java/de/bwaldvogel/mongo/backend/AbstractOplogTest.java‎

Lines changed: 43 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -17,13 +17,13 @@
1717
import java.util.List;
1818
import java.util.NoSuchElementException;
1919
import java.util.UUID;
20-
import java.util.concurrent.TimeUnit;
2120

2221
import org.bson.BsonDocument;
2322
import org.bson.BsonInt32;
2423
import org.bson.BsonTimestamp;
2524
import org.bson.Document;
2625
import org.bson.conversions.Bson;
26+
import org.junit.jupiter.api.AfterEach;
2727
import org.junit.jupiter.api.BeforeEach;
2828
import org.junit.jupiter.api.Disabled;
2929
import org.junit.jupiter.api.Test;
@@ -40,7 +40,6 @@
4040
import com.mongodb.reactivestreams.client.Success;
4141

4242
import de.bwaldvogel.mongo.oplog.OperationType;
43-
import io.reactivex.subscribers.TestSubscriber;
4443

4544
public abstract class AbstractOplogTest extends AbstractTest {
4645

@@ -52,6 +51,7 @@ public void beforeEach() {
5251
backend.enableOplog();
5352
}
5453

54+
@AfterEach
5555
@Override
5656
void assertNoOpenCursors() throws Exception {
5757
// workaround to close all change stream publishers
@@ -384,88 +384,86 @@ public void testChangeStreamStartAtOperationTime() {
384384
ChangeStreamDocument<Document> document = cursor.next();
385385
BsonTimestamp startAtOperationTime = document.getClusterTime();
386386

387-
MongoChangeStreamCursor<ChangeStreamDocument<Document>> cursor2 = collection.watch().startAtOperationTime(startAtOperationTime).cursor();
388-
ChangeStreamDocument<Document> document2 = cursor2.next();
389-
assertThat(document2.getFullDocument().get("a")).isEqualTo(2);
390-
document2 = cursor2.next();
391-
assertThat(document2.getFullDocument().get("a")).isEqualTo(3);
387+
try (MongoChangeStreamCursor<ChangeStreamDocument<Document>> cursor2 = collection.watch().startAtOperationTime(startAtOperationTime).cursor()) {
388+
ChangeStreamDocument<Document> document2 = cursor2.next();
389+
assertThat(document2.getFullDocument().get("a")).isEqualTo(2);
390+
document2 = cursor2.next();
391+
assertThat(document2.getFullDocument().get("a")).isEqualTo(3);
392+
}
392393
}
393394
}
394395

395396
@Test
396-
public void testChangeStreamAndReplaceOneWithUpsertTrue() throws Exception {
397+
void testChangeStreamAndReplaceOneWithUpsertTrue() throws Throwable {
397398
TestSubscriber<ChangeStreamDocument<Document>> streamSubscriber = new TestSubscriber<>();
398399
asyncCollection.watch().fullDocument(FullDocument.UPDATE_LOOKUP).subscribe(streamSubscriber);
399400

400401
TestSubscriber<UpdateResult> replaceOneSubscriber = new TestSubscriber<>();
401402
asyncCollection.replaceOne(json("a: 1"), json("a: 1"), new ReplaceOptions().upsert(true))
402403
.subscribe(replaceOneSubscriber);
403404

404-
replaceOneSubscriber.awaitTerminalEvent();
405-
replaceOneSubscriber.assertNoErrors();
405+
replaceOneSubscriber.awaitSingleValue();
406406

407407
TestSubscriber<Document> findSubscriber = new TestSubscriber<>();
408408
asyncCollection.find(json("a:1")).subscribe(findSubscriber);
409-
findSubscriber.awaitTerminalEvent();
410-
assertThat(getSingleValue(findSubscriber).get("a")).isEqualTo(1);
409+
assertThat(findSubscriber.awaitSingleValue().get("a")).isEqualTo(1);
411410

412-
streamSubscriber.awaitCount(1).assertValueCount(1).cancel();
413-
assertThat(getSingleValue(streamSubscriber).getOperationType().getValue()).isEqualTo("insert");
414-
assertThat(getSingleValue(streamSubscriber).getFullDocument()).isEqualTo(getSingleValue(findSubscriber));
411+
ChangeStreamDocument<Document> value = streamSubscriber.awaitSingleValue();
412+
assertThat(value.getOperationType().getValue()).isEqualTo("insert");
413+
assertThat(value.getFullDocument()).isEqualTo(findSubscriber.awaitSingleValue());
415414
}
416415

417416
@Test
418-
public void testSimpleChangeStreamWithFilter() throws Exception {
419-
TestSubscriber<Success> insertSubscriber1 = new TestSubscriber<>();
420-
TestSubscriber<Success> insertSubscriber2 = new TestSubscriber<>();
417+
void testSimpleChangeStreamWithFilter() throws Exception {
421418
TestSubscriber<ChangeStreamDocument<Document>> streamSubscriber = new TestSubscriber<>();
422419

423-
asyncCollection.insertOne(json("_id: 2")).subscribe(insertSubscriber1);
424-
insertSubscriber1.awaitTerminalEvent();
420+
insertOne(asyncCollection, json("_id: 1"));
425421

426422
Bson filter = match(Filters.eq("fullDocument.bu", "abc"));
427423
List<Bson> pipeline = singletonList(filter);
428424

429425
asyncCollection.watch(pipeline).subscribe(streamSubscriber);
430426

431-
insertSubscriber1 = new TestSubscriber<>();
432-
asyncCollection.insertOne(json("_id: 2, bu: 'abc'")).subscribe(insertSubscriber1);
433-
asyncCollection.insertOne(json("_id: 3, bu: 'xyz'")).subscribe(insertSubscriber2);
434-
insertSubscriber1.awaitTerminalEvent();
435-
insertSubscriber2.awaitTerminalEvent();
427+
insertOne(asyncCollection, json("_id: 2, bu: 'abc'"));
428+
insertOne(asyncCollection, json("_id: 3, bu: 'xyz'"));
436429

437-
streamSubscriber.awaitCount(1).assertValueCount(1).cancel();
438-
assertThat(getSingleValue(streamSubscriber).getFullDocument().get("bu")).isEqualTo("abc");
430+
ChangeStreamDocument<Document> changeStreamDocument = streamSubscriber.awaitSingleValue();
431+
assertThat(changeStreamDocument.getFullDocument().get("bu")).isEqualTo("abc");
439432
}
440433

441434
@Test
442-
public void testOplogShouldFilterNamespaceOnChangeStreams() throws Exception {
443-
TestSubscriber<Success> insertSubscriber = new TestSubscriber<>();
435+
void testOplogSubscription() throws Exception {
436+
TestSubscriber<ChangeStreamDocument<Document>> streamSubscriber = new TestSubscriber<>();
437+
asyncCollection.watch().subscribe(streamSubscriber);
438+
439+
insertOne(asyncCollection, json("_id: 1"));
440+
441+
ChangeStreamDocument<Document> changeStreamDocument = streamSubscriber.awaitSingleValue();
442+
assertThat(changeStreamDocument.getOperationType()).isEqualTo(com.mongodb.client.model.changestream.OperationType.INSERT);
443+
assertThat(changeStreamDocument.getFullDocument()).isEqualTo(json("_id: 1"));
444+
}
445+
446+
@Test
447+
void testOplogShouldFilterNamespaceOnChangeStreams() throws Exception {
444448
com.mongodb.reactivestreams.client.MongoCollection<Document> asyncCollection1 =
445449
asyncDb.getCollection(asyncCollection.getNamespace().getCollectionName() + "1");
446450

447-
asyncCollection.insertOne(json("_id: 1")).subscribe(insertSubscriber);
448-
asyncCollection1.insertOne(json("_id: 1")).subscribe(insertSubscriber);
449-
450-
insertSubscriber.awaitTerminalEvent(1, TimeUnit.SECONDS);
451+
insertOne(asyncCollection, json("_id: 1"));
452+
insertOne(asyncCollection1, json("_id: 1"));
451453

452454
TestSubscriber<ChangeStreamDocument<Document>> streamSubscriber = new TestSubscriber<>();
453-
Bson filter = match(Filters.eq("fullDocument.a", 1));
454-
asyncCollection.watch(singletonList(filter)).subscribe(streamSubscriber);
455-
456-
// Necessary to give time for the change stream to start before the below insert operation is executed.
457-
Thread.sleep(50);
455+
asyncCollection.watch().subscribe(streamSubscriber);
458456

459-
insertSubscriber = new TestSubscriber<>();
460-
asyncCollection1.insertOne(json("_id: 2, a: 1")).subscribe(insertSubscriber);
461-
insertSubscriber.awaitTerminalEvent(1, TimeUnit.SECONDS);
457+
insertOne(asyncCollection1, json("_id: 2"));
458+
insertOne(asyncCollection, json("_id: 2"));
462459

463-
streamSubscriber.awaitTerminalEvent(1, TimeUnit.SECONDS);
464-
assertThat(streamSubscriber.values()).isEmpty();
460+
streamSubscriber.awaitSingleValue();
465461
}
466462

467-
private static <T> T getSingleValue(TestSubscriber<T> subscriber) {
468-
return subscriber.values().get(0);
463+
private static void insertOne(com.mongodb.reactivestreams.client.MongoCollection<Document> collection, Document document) throws Exception {
464+
TestSubscriber<Success> insertSubscriber = new TestSubscriber<>();
465+
collection.insertOne(document).subscribe(insertSubscriber);
466+
insertSubscriber.awaitSingleValue();
469467
}
470468

471469
}
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
package de.bwaldvogel.mongo.backend;
2+
3+
import java.util.concurrent.CountDownLatch;
4+
import java.util.concurrent.TimeUnit;
5+
6+
import org.reactivestreams.Subscriber;
7+
import org.reactivestreams.Subscription;
8+
import org.slf4j.Logger;
9+
import org.slf4j.LoggerFactory;
10+
11+
class TestSubscriber<T> implements Subscriber<T> {
12+
13+
private static final Logger log = LoggerFactory.getLogger(TestSubscriber.class);
14+
15+
private final CountDownLatch countDownLatch = new CountDownLatch(1);
16+
private T value;
17+
private Throwable throwable;
18+
private Subscription subscription;
19+
20+
@Override
21+
public void onSubscribe(Subscription subscription) {
22+
this.subscription = subscription;
23+
subscription.request(1);
24+
}
25+
26+
@Override
27+
public void onNext(T value) {
28+
log.debug("onNext: {}", value);
29+
Assert.isNull(this.value, () -> "Got a second value: " + this.value + " and " + value);
30+
this.value = value;
31+
subscription.cancel();
32+
countDownLatch.countDown();
33+
}
34+
35+
@Override
36+
public void onError(Throwable throwable) {
37+
log.error("onError", throwable);
38+
this.throwable = throwable;
39+
}
40+
41+
@Override
42+
public void onComplete() {
43+
log.info("onComplete", throwable);
44+
}
45+
46+
T awaitSingleValue() throws Exception {
47+
boolean success = countDownLatch.await(30, TimeUnit.SECONDS);
48+
Assert.isTrue(success, () -> "Failed waiting countdown latch");
49+
if (throwable != null) {
50+
throw new RuntimeException(throwable);
51+
}
52+
Assert.notNull(value, () -> "Got no value yet");
53+
return value;
54+
}
55+
56+
}

0 commit comments

Comments
 (0)