2121import okhttp3 .extension .logging .HttpLogLevel ;
2222
2323import static org .junit .jupiter .api .Assertions .assertEquals ;
24+ import static org .junit .jupiter .api .Assertions .assertFalse ;
2425import static org .junit .jupiter .api .Assertions .assertNotNull ;
2526import static org .junit .jupiter .api .Assertions .assertTrue ;
2627
@@ -29,7 +30,7 @@ class SseLifecycleContractTest {
2930 @ Test
3031 void standardChatChunkDoesNotRequireNestedDataField () throws Exception {
3132 OkHttpClient client = new OkHttpClient .Builder ().addInterceptor (chain -> {
32- String body = "data: {\" choices\" :[{\" delta\" :{\" content\" :\" hello \" }}]}\n \n "
33+ String body = "data: {\" choices\" :[{\" delta\" :{\" content\" :\" 你好 \" }}]}\n \n "
3334 + "data: [DONE]\n \n " ;
3435 return new Response .Builder ()
3536 .request (chain .request ())
@@ -53,7 +54,9 @@ void standardChatChunkDoesNotRequireNestedDataField() throws Exception {
5354
5455 assertTrue (complete .await (3 , TimeUnit .SECONDS ));
5556 assertNotNull (received .get ());
56- assertEquals ("hello" , received .get ().deltaText ());
57+ assertEquals ("你好 " , received .get ().deltaText ());
58+ assertEquals ("{\" choices\" :[{\" delta\" :{\" content\" :\" 你好 \" }}]}" ,
59+ received .get ().getData ());
5760 }
5861 }
5962
@@ -83,15 +86,97 @@ void malformedJsonAndConsumerFailureRemainIsolated() throws Exception {
8386 AtomicInteger consumerCalls = new AtomicInteger ();
8487
8588 try (HermesSseClient sse = new HermesSseClient (config , null , client )) {
86- CountDownLatch complete = new CountDownLatch (1 );
87- sse .subscribeChat (request , ignored -> {
89+ CountDownLatch failed = new CountDownLatch (1 );
90+ AtomicReference <Throwable > failure = new AtomicReference <>();
91+ SseSubscription subscription = sse .subscribeChat (request , ignored -> {
8892 consumerCalls .incrementAndGet ();
8993 throw new IllegalStateException ("consumer failure" );
90- }, complete ::countDown , ignored -> { });
94+ }, () -> { }, error -> {
95+ failure .set (error );
96+ failed .countDown ();
97+ });
9198
92- assertTrue (complete .await (3 , TimeUnit .SECONDS ));
99+ assertTrue (failed .await (3 , TimeUnit .SECONDS ));
93100 assertEquals (1 , consumerCalls .get (),
94101 "malformed JSON must not be delivered to the business consumer" );
102+ assertTrue (failure .get () instanceof SseConsumerException ,
103+ "consumer callback failures must not be reported as JSON parsing failures" );
104+ assertFalse (subscription .isActive ());
105+ assertEquals (0 , sse .activeSubscriptionCount ());
106+ }
107+ }
108+
109+ @ Test
110+ void multiLineCrLfCommentAndFragmentedUtf8PreserveRawFrame () throws Exception {
111+ try (MockWebServer server = new MockWebServer ()) {
112+ String body = ": keepalive\r \n "
113+ + "id: evt-1\r \n "
114+ + "event: assistant.delta\r \n "
115+ + "data: {\" choices\" :[{\" delta\" :\r \n "
116+ + "data: {\" content\" :\" 你 好 \" }}]}\r \n \r \n "
117+ + "data: [DONE]\r \n \r \n " ;
118+ server .enqueue (new MockResponse ()
119+ .setResponseCode (200 )
120+ .setHeader ("Content-Type" , "text/event-stream" )
121+ .setChunkedBody (body , 1 ));
122+ server .start ();
123+
124+ HermesHttpClientConfig config = new HermesHttpClientConfig ()
125+ .setEndpointPolicy (io .github .easy4j .hermes .security .EndpointPolicy
126+ .trustedLocal ("127.0.0.1" , server .getPort ()))
127+ .setBaseUrl ("http://127.0.0.1:" + server .getPort ());
128+ ChatRequest request = new ChatRequest ();
129+ request .setMessages (Collections .singletonList (
130+ new ChatRequest .Message ("user" , "hello" )));
131+
132+ try (HermesSseClient sse = new HermesSseClient (config , null , null )) {
133+ CountDownLatch complete = new CountDownLatch (1 );
134+ AtomicReference <SseEvent > received = new AtomicReference <>();
135+ sse .subscribeChat (request , received ::set , complete ::countDown , ignored -> { });
136+
137+ assertTrue (complete .await (3 , TimeUnit .SECONDS ));
138+ SseEvent event = received .get ();
139+ assertNotNull (event );
140+ assertEquals ("evt-1" , event .getId ());
141+ assertEquals ("assistant.delta" , event .getEvent ());
142+ assertEquals ("你 好 " , event .deltaText ());
143+ assertTrue (event .getData ().contains ("\n " ),
144+ "multiple data lines must remain visible in the raw SSE data" );
145+ }
146+ }
147+ }
148+
149+ @ Test
150+ void unknownEventPreservesServerIdentityAndRawData () throws Exception {
151+ OkHttpClient client = new OkHttpClient .Builder ().addInterceptor (chain -> {
152+ String body = "id: mystery-7\n "
153+ + "event: mystery.data\n "
154+ + "data: {\" foo\" :\" bar\" }\n \n "
155+ + "data: [DONE]\n \n " ;
156+ return new Response .Builder ()
157+ .request (chain .request ())
158+ .protocol (Protocol .HTTP_1_1 )
159+ .code (200 )
160+ .message ("OK" )
161+ .body (ResponseBody .create (body , MediaType .get ("text/event-stream" )))
162+ .build ();
163+ }).build ();
164+
165+ HermesHttpClientConfig config = new HermesHttpClientConfig ();
166+ config .markUnsafeBaseUrlOverriddenForTest (true );
167+ ChatRequest request = new ChatRequest ();
168+ request .setMessages (Collections .singletonList (new ChatRequest .Message ("user" , "hello" )));
169+
170+ try (HermesSseClient sse = new HermesSseClient (config , null , client )) {
171+ CountDownLatch complete = new CountDownLatch (1 );
172+ AtomicReference <SseEvent > received = new AtomicReference <>();
173+ sse .subscribeChat (request , received ::set , complete ::countDown , ignored -> { });
174+
175+ assertTrue (complete .await (3 , TimeUnit .SECONDS ));
176+ assertNotNull (received .get ());
177+ assertEquals ("mystery-7" , received .get ().getId ());
178+ assertEquals ("mystery.data" , received .get ().getEvent ());
179+ assertEquals ("{\" foo\" :\" bar\" }" , received .get ().getData ());
95180 }
96181 }
97182
0 commit comments