Skip to content

Commit 089a050

Browse files
committed
perf: avoid unused sse event buffering
1 parent 2013f22 commit 089a050

5 files changed

Lines changed: 22 additions & 30 deletions

File tree

‎src/main/java/io/github/easy4j/hermes/HermesClient.java‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -442,8 +442,7 @@ public StreamingChatResponse chatCompletionStream(ChatRequest request,
442442
Map<String, String> headers,
443443
Consumer<String> deltaConsumer) {
444444
Objects.requireNonNull(request, "request");
445-
StreamingChatResponse stream = new StreamingChatResponse(
446-
config.getHttp().getStreamEventQueueCapacity());
445+
StreamingChatResponse stream = new StreamingChatResponse();
447446
stream.onDelta(deltaConsumer);
448447
SseSubscription subscription = sseClient.subscribeChat(
449448
request.withStream(), headers, stream::accept, stream::finish, stream::fail);

‎src/main/java/io/github/easy4j/hermes/api/HermesChatClient.java‎

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -19,29 +19,28 @@
1919
*/
2020
public class HermesChatClient extends HermesHttpClient {
2121

22-
private final HermesHttpClientConfig config;
2322
private final HermesSseClient sseClient;
2423
private final boolean ownsSseClient;
2524

2625
public HermesChatClient(HermesHttpClientConfig config) {
2726
super(config);
28-
this.config = Objects.requireNonNull(config, "config");
27+
Objects.requireNonNull(config, "config");
2928
this.sseClient = new HermesSseClient(config, getObjectMapper(), getOkHttpClient());
3029
this.ownsSseClient = true;
3130
}
3231

3332
public HermesChatClient(HermesHttpClientConfig config, ObjectMapper objectMapper,
3433
OkHttpClient httpClient) {
3534
super(config, objectMapper, httpClient);
36-
this.config = Objects.requireNonNull(config, "config");
35+
Objects.requireNonNull(config, "config");
3736
this.sseClient = new HermesSseClient(config, getObjectMapper(), getOkHttpClient());
3837
this.ownsSseClient = true;
3938
}
4039

4140
public HermesChatClient(HermesHttpClientConfig config, ObjectMapper objectMapper,
4241
OkHttpClient httpClient, HermesSseClient sseClient) {
4342
super(config, objectMapper, httpClient);
44-
this.config = Objects.requireNonNull(config, "config");
43+
Objects.requireNonNull(config, "config");
4544
this.sseClient = Objects.requireNonNull(sseClient, "sseClient");
4645
this.ownsSseClient = false;
4746
}
@@ -60,8 +59,7 @@ public StreamingChatResponse chatCompletionStream(ChatRequest request,
6059
Map<String, String> headers,
6160
Consumer<String> deltaConsumer) {
6261
Objects.requireNonNull(request, "request");
63-
StreamingChatResponse stream = new StreamingChatResponse(
64-
config.getStreamEventQueueCapacity()).onDelta(deltaConsumer);
62+
StreamingChatResponse stream = new StreamingChatResponse().onDelta(deltaConsumer);
6563
SseSubscription subscription = sseClient.subscribeChat(
6664
request.withStream(), headers, stream::accept, stream::finish, stream::fail);
6765
stream.onCancel(subscription::close);

‎src/main/java/io/github/easy4j/hermes/api/sse/StreamingChatResponse.java‎

Lines changed: 1 addition & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,5 @@
11
package io.github.easy4j.hermes.api.sse;
22

3-
import lombok.Getter;
4-
5-
import java.util.concurrent.ArrayBlockingQueue;
6-
import java.util.concurrent.BlockingQueue;
73
import java.util.concurrent.CompletableFuture;
84
import java.util.concurrent.atomic.AtomicReference;
95
import java.util.Objects;
@@ -17,15 +13,9 @@ public class StreamingChatResponse extends CompletableFuture<String> {
1713
private final StringBuilder content = new StringBuilder();
1814
private Consumer<String> deltaConsumer;
1915
private final AtomicReference<Runnable> cancellation = new AtomicReference<>();
20-
@Getter
21-
private final BlockingQueue<SseEvent> eventQueue;
2216

17+
/** 创建不缓存原始事件的流式响应。 */
2318
public StreamingChatResponse() {
24-
this(1_024);
25-
}
26-
27-
public StreamingChatResponse(int eventQueueCapacity) {
28-
this.eventQueue = new ArrayBlockingQueue<>(Math.max(1, eventQueueCapacity));
2919
}
3020

3121
public StreamingChatResponse onDelta(Consumer<String> consumer) {
@@ -34,10 +24,6 @@ public StreamingChatResponse onDelta(Consumer<String> consumer) {
3424
}
3525

3626
public void accept(SseEvent event) {
37-
if (!eventQueue.offer(event)) {
38-
eventQueue.poll();
39-
eventQueue.offer(event);
40-
}
4127
String delta = event.deltaText();
4228
if (delta != null) {
4329
content.append(delta);

‎src/test/java/io/github/easy4j/hermes/HermesSseAndModelTest.java‎

Lines changed: 14 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,6 @@ void shouldParseSseEventsAndAccumulateStreamingContent() throws Exception {
6363
stream.accept(chat);
6464
stream.accept(direct);
6565
stream.accept(event("{}"));
66-
assertEquals(3, stream.getEventQueue().size());
6766
assertEquals("hello world", stream.getAccumulatedContent());
6867
stream.finish();
6968
assertEquals("hello world", stream.get(1, TimeUnit.SECONDS));
@@ -261,6 +260,20 @@ void shouldKeepLatestQueueEventAndExposeSubscriptionLifecycle() throws Exception
261260
assertEquals(0, sse.activeSubscriptionCount());
262261
}
263262

263+
server.enqueue(new MockResponse().setResponseCode(200)
264+
.setHeader("Content-Type", "text/event-stream")
265+
.setBody("data: [DONE]\n\n"));
266+
try (HermesSseClient sse = new HermesSseClient(config, null, null)) {
267+
SseSubscription completed = sse.subscribeSessionEvents(
268+
"session-complete", "hello", event -> { });
269+
assertNotNull(server.takeRequest(3, TimeUnit.SECONDS));
270+
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(3);
271+
while (completed.isActive() && System.nanoTime() < deadline) {
272+
Thread.yield();
273+
}
274+
assertFalse(completed.isActive());
275+
}
276+
264277
server.enqueue(new MockResponse().setResponseCode(200)
265278
.setHeader("Content-Type", "text/event-stream")
266279
.setBodyDelay(3, TimeUnit.SECONDS).setBody("data: [DONE]\n\n"));
@@ -294,12 +307,6 @@ void shouldCoverChatStreamingOverloadsAndCancellationPaths() throws Exception {
294307
HermesOkHttpClientFactory.shutdown(client);
295308
}
296309

297-
StreamingChatResponse bounded = new StreamingChatResponse(1);
298-
bounded.accept(event("{\"delta\":\"first\"}"));
299-
SseEvent latest = event("{\"delta\":\"latest\"}");
300-
bounded.accept(latest);
301-
assertEquals(latest, bounded.getEventQueue().poll());
302-
303310
AtomicInteger cancellations = new AtomicInteger();
304311
StreamingChatResponse cancelBeforeBind = new StreamingChatResponse();
305312
cancelBeforeBind.cancel(false);

‎src/test/java/io/github/easy4j/hermes/HermesSseApiShapeTest.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,8 @@ void shouldExposeOnlyTheUnifiedSseTopology() {
3333
assertFalse(hasMethod(HermesSseClient.class, "subscribeSessionStream"));
3434
assertFalse(hasMethod(HermesSseClient.class, "stop"));
3535
assertFalse(hasMethod(HermesChatClient.class, "events"));
36+
assertFalse(hasMethod(io.github.easy4j.hermes.api.sse.StreamingChatResponse.class,
37+
"getEventQueue"));
3638
for (Field field : HermesChatClient.class.getDeclaredFields()) {
3739
assertFalse("eventClient".equals(field.getName()));
3840
}

0 commit comments

Comments
 (0)