@@ -495,4 +495,41 @@ void chatTransportFailureDoesNotReplayWithoutVerifiedIdempotency() throws Except
495495 }
496496 }
497497
498+
499+ @ Test
500+ void saturatedQueueFailsExplicitlyWithoutDroppingOldestEvent () throws Exception {
501+ try (MockWebServer server = new MockWebServer ()) {
502+ server .enqueue (new MockResponse ()
503+ .setResponseCode (200 )
504+ .setHeader ("Content-Type" , "text/event-stream" )
505+ .setBody ("id: evt-1\n data: {\" delta\" :\" first\" }\n \n "
506+ + "id: evt-2\n data: {\" delta\" :\" second\" }\n \n "
507+ + "data: [DONE]\n \n " ));
508+ server .start ();
509+
510+ HermesHttpClientConfig config = new HermesHttpClientConfig ()
511+ .setEndpointPolicy (io .github .easy4j .hermes .security .EndpointPolicy
512+ .trustedLocal ("127.0.0.1" , server .getPort ()))
513+ .setBaseUrl ("http://127.0.0.1:" + server .getPort ());
514+ config .setStreamEventQueueCapacity (1 );
515+ config .setStreamReconnectMaxAttempts (0 );
516+
517+ try (HermesSseClient sse = new HermesSseClient (config , null , null );
518+ SseQueueSubscription queued = sse .subscribeRunEventsQueue ("run-overflow" )) {
519+ long deadline = System .nanoTime () + TimeUnit .SECONDS .toNanos (3 );
520+ while (queued .getSubscription ().isActive () && System .nanoTime () < deadline ) {
521+ Thread .yield ();
522+ }
523+
524+ assertFalse (queued .getSubscription ().isActive ());
525+ assertTrue (queued .getSubscription ().getTerminalError ()
526+ instanceof SseQueueOverflowException ,
527+ "queue saturation must be reported explicitly" );
528+ assertEquals (1 , queued .getQueue ().size ());
529+ assertEquals ("evt-1" , queued .getQueue ().peek ().getId (),
530+ "overflow must not silently discard the oldest undelivered event" );
531+ }
532+ }
533+ }
534+
498535}
0 commit comments