@@ -1105,17 +1105,22 @@ public void pubsub_sharded_max_size_message_callback(boolean standalone) {
11051105 @ SneakyThrows
11061106 @ ParameterizedTest (name = "standalone = {0}" )
11071107 @ ValueSource (booleans = {true , false })
1108- @ Disabled (
1109- "This test is for demonstrating behavior when there's an exception from a user callback and"
1110- + " will always fail." )
11111108 public void pubsub_test_callback_exception (boolean standalone ) {
11121109 final GlideString channel = gs (UUID .randomUUID ().toString ());
1113- final GlideString message = gs ("message" );
1110+ final GlideString message1 = gs ("message1" );
1111+ final GlideString message2 = gs ("message2" );
1112+ final GlideString errorMsg = gs ("errorMsg" );
1113+ final GlideString message3 = gs ("message3" );
11141114
11151115 ArrayList <PubSubMessage > callbackMessages = new ArrayList <>();
11161116 final MessageCallback callback =
11171117 (pubSubMessage , context ) -> {
1118- throw new RuntimeException ("Test exception." );
1118+ new Exception ().printStackTrace ();
1119+ if (pubSubMessage .getMessage ().equals (errorMsg )) {
1120+ throw new RuntimeException ("Test callback error message" );
1121+ }
1122+ ArrayList <PubSubMessage > receivedMessages = (ArrayList <PubSubMessage >) context ;
1123+ receivedMessages .add (pubSubMessage );
11191124 };
11201125
11211126 Map <? extends ChannelMode , Set <GlideString >> subscriptions =
@@ -1130,22 +1135,30 @@ public void pubsub_test_callback_exception(boolean standalone) {
11301135 Optional .ofNullable (callback ),
11311136 Optional .of (callbackMessages ));
11321137 var sender =
1133- createClientWithSubscriptions (
1134- standalone ,
1135- subscriptions ,
1136- Optional .ofNullable (callback ),
1137- Optional .of (callbackMessages ));
1138+ createClient (standalone );
11381139
11391140 try {
1140- assertEquals (OK , sender .publish (message , channel ).get ());
1141+ assertEquals (OK , sender .publish (message1 , channel ).get ());
1142+ assertEquals (OK , sender .publish (message2 , channel ).get ());
1143+ // assertEquals(OK, sender.publish(errorMsg, channel).get());
1144+ assertEquals (OK , sender .publish (message3 , channel ).get ());
11411145
11421146 // Allow the message to propagate.
11431147 Thread .sleep (MESSAGE_DELIVERY_DELAY );
11441148
1145- assertEquals (1 , callbackMessages .size ());
1146- assertEquals (message , callbackMessages .get (0 ).getMessage ());
1149+ assertEquals (3 , callbackMessages .size ());
1150+ assertEquals (message1 , callbackMessages .get (0 ).getMessage ());
11471151 assertEquals (channel , callbackMessages .get (0 ).getChannel ());
1148- assertNull (callbackMessages .get (0 ).getPattern ());
1152+ assertTrue (callbackMessages .get (0 ).getPattern ().isEmpty ());
1153+
1154+ assertEquals (message2 , callbackMessages .get (1 ).getMessage ());
1155+ assertEquals (channel , callbackMessages .get (1 ).getChannel ());
1156+ assertTrue (callbackMessages .get (1 ).getPattern ().isEmpty ());
1157+
1158+ // Ensure we can receive message 3 which is after the message that triggers a throw.
1159+ assertEquals (message3 , callbackMessages .get (2 ).getMessage ());
1160+ assertEquals (channel , callbackMessages .get (2 ).getChannel ());
1161+ assertTrue (callbackMessages .get (2 ).getPattern ().isEmpty ());
11491162 } finally {
11501163 if (!standalone ) {
11511164 // Since all tests run on the same cluster, when closing the client, garbage collector can
0 commit comments