Skip to content

Commit f5fc18f

Browse files
committed
perf: avoid unused sse event buffering
1 parent 826d1d4 commit f5fc18f

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
@@ -445,8 +445,7 @@ public StreamingChatResponse chatCompletionStream(ChatRequest request,
445445
Map<String, String> headers,
446446
Consumer<String> deltaConsumer) {
447447
Objects.requireNonNull(request, "request");
448-
StreamingChatResponse stream = new StreamingChatResponse(
449-
config.getHttp().getStreamEventQueueCapacity());
448+
StreamingChatResponse stream = new StreamingChatResponse();
450449
stream.onDelta(deltaConsumer);
451450
SseSubscription subscription = sseClient.subscribeChat(
452451
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
@@ -22,29 +22,28 @@
2222
*/
2323
public class HermesChatClient extends HermesHttpClient {
2424

25-
private final HermesHttpClientConfig config;
2625
private final HermesSseClient sseClient;
2726
private final boolean ownsSseClient;
2827

2928
public HermesChatClient(HermesHttpClientConfig config) {
3029
super(config);
31-
this.config = Objects.requireNonNull(config, "config");
30+
Objects.requireNonNull(config, "config");
3231
this.sseClient = new HermesSseClient(config, getObjectMapper(), getOkHttpClient());
3332
this.ownsSseClient = true;
3433
}
3534

3635
public HermesChatClient(HermesHttpClientConfig config, ObjectMapper objectMapper,
3736
OkHttpClient httpClient) {
3837
super(config, objectMapper, httpClient);
39-
this.config = Objects.requireNonNull(config, "config");
38+
Objects.requireNonNull(config, "config");
4039
this.sseClient = new HermesSseClient(config, getObjectMapper(), getOkHttpClient());
4140
this.ownsSseClient = true;
4241
}
4342

4443
public HermesChatClient(HermesHttpClientConfig config, ObjectMapper objectMapper,
4544
OkHttpClient httpClient, HermesSseClient sseClient) {
4645
super(config, objectMapper, httpClient);
47-
this.config = Objects.requireNonNull(config, "config");
46+
Objects.requireNonNull(config, "config");
4847
this.sseClient = Objects.requireNonNull(sseClient, "sseClient");
4948
this.ownsSseClient = false;
5049
}
@@ -63,8 +62,7 @@ public StreamingChatResponse chatCompletionStream(ChatRequest request,
6362
Map<String, String> headers,
6463
Consumer<String> deltaConsumer) {
6564
Objects.requireNonNull(request, "request");
66-
StreamingChatResponse stream = new StreamingChatResponse(
67-
config.getStreamEventQueueCapacity()).onDelta(deltaConsumer);
65+
StreamingChatResponse stream = new StreamingChatResponse().onDelta(deltaConsumer);
6866
SseSubscription subscription = sseClient.subscribeChat(
6967
request.withStream(), headers, stream::accept, stream::finish, stream::fail);
7068
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
@@ -3,10 +3,6 @@
33
* @author <a href="https://github.com/loong10k">@Loong Wan</a>
44
*/
55

6-
import lombok.Getter;
7-
8-
import java.util.concurrent.ArrayBlockingQueue;
9-
import java.util.concurrent.BlockingQueue;
106
import java.util.concurrent.CompletableFuture;
117
import java.util.concurrent.atomic.AtomicReference;
128
import java.util.Objects;
@@ -20,15 +16,9 @@ public class StreamingChatResponse extends CompletableFuture<String> {
2016
private final StringBuilder content = new StringBuilder();
2117
private Consumer<String> deltaConsumer;
2218
private final AtomicReference<Runnable> cancellation = new AtomicReference<>();
23-
@Getter
24-
private final BlockingQueue<SseEvent> eventQueue;
2519

20+
/** 创建不缓存原始事件的流式响应。 */
2621
public StreamingChatResponse() {
27-
this(1_024);
28-
}
29-
30-
public StreamingChatResponse(int eventQueueCapacity) {
31-
this.eventQueue = new ArrayBlockingQueue<>(Math.max(1, eventQueueCapacity));
3222
}
3323

3424
public StreamingChatResponse onDelta(Consumer<String> consumer) {
@@ -37,10 +27,6 @@ public StreamingChatResponse onDelta(Consumer<String> consumer) {
3727
}
3828

3929
public void accept(SseEvent event) {
40-
if (!eventQueue.offer(event)) {
41-
eventQueue.poll();
42-
eventQueue.offer(event);
43-
}
4430
String delta = event.deltaText();
4531
if (delta != null) {
4632
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)