1717import java .util .List ;
1818import java .util .NoSuchElementException ;
1919import java .util .UUID ;
20- import java .util .concurrent .TimeUnit ;
2120
2221import org .bson .BsonDocument ;
2322import org .bson .BsonInt32 ;
2423import org .bson .BsonTimestamp ;
2524import org .bson .Document ;
2625import org .bson .conversions .Bson ;
26+ import org .junit .jupiter .api .AfterEach ;
2727import org .junit .jupiter .api .BeforeEach ;
2828import org .junit .jupiter .api .Disabled ;
2929import org .junit .jupiter .api .Test ;
4040import com .mongodb .reactivestreams .client .Success ;
4141
4242import de .bwaldvogel .mongo .oplog .OperationType ;
43- import io .reactivex .subscribers .TestSubscriber ;
4443
4544public 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}
0 commit comments