Skip to content

Commit 53831d4

Browse files
committed
feat: unify blocking and streaming HTTP modes
1 parent 4fec952 commit 53831d4

4 files changed

Lines changed: 47 additions & 5 deletions

File tree

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
package io.github.easy4j.opencode;
2+
3+
/**
4+
* HTTP 对话响应模式。
5+
*/
6+
public enum HttpResponseMode {
7+
8+
/** 等待完整响应后一次性返回。 */
9+
BLOCKING,
10+
11+
/** 消费 Provider SSE 并逐段回调。 */
12+
STREAM,
13+
14+
/** 由调用方根据结构化输出等请求特征选择。 */
15+
AUTO
16+
}

src/main/java/io/github/easy4j/opencode/OpenCodeClient.java

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
import java.util.concurrent.ThreadPoolExecutor;
2222
import java.util.concurrent.TimeUnit;
2323
import java.util.concurrent.atomic.AtomicInteger;
24+
import java.util.function.Consumer;
2425

2526
/**
2627
* OpenCode 客户端门面:HTTP Server + SSE 事件流 + 本地 CLI。
@@ -226,6 +227,7 @@ public boolean isCliEnabled() {
226227
// ============================================================
227228

228229
private void copyHttpConfig(OpenCodeHttpClientConfig src) {
230+
this.config.getHttp().setMode(src.getMode());
229231
this.config.getHttp().setEnabled(src.isEnabled());
230232
this.config.getHttp().setStartupCheckEnabled(src.isStartupCheckEnabled());
231233
this.config.getHttp().setFailFastOnUnavailable(src.isFailFastOnUnavailable());
@@ -244,7 +246,7 @@ private void copyHttpConfig(OpenCodeHttpClientConfig src) {
244246
this.config.getHttp().setStreamMaxPoolSize(src.getStreamMaxPoolSize());
245247
this.config.getHttp().setStreamQueueCapacity(src.getStreamQueueCapacity());
246248
this.config.getHttp().setStreamKeepAliveMillis(src.getStreamKeepAliveMillis());
247-
this.config.getHttp().setSseEventQueueCapacity(src.getSseEventQueueCapacity());
249+
this.config.getHttp().setStreamEventQueueCapacity(src.getStreamEventQueueCapacity());
248250
this.config.getHttp().setRetryOnConnectionFailure(src.isRetryOnConnectionFailure());
249251
this.config.getHttp().setVerifySsl(src.isVerifySsl());
250252
this.config.getHttp().setDefaultModel(src.getDefaultModel());
@@ -360,10 +362,19 @@ public ChatStreamingResponse chatCompletionStream(ChatRequest request, String se
360362

361363
public ChatStreamingResponse chatCompletionStream(ChatRequest request, String sessionKey,
362364
OpenCodeRequestContext context) {
365+
return chatCompletionStream(request, sessionKey, context, null);
366+
}
367+
368+
/**
369+
* 流式对话,并在订阅启动前绑定增量回调,避免丢失首批分片。
370+
*/
371+
public ChatStreamingResponse chatCompletionStream(ChatRequest request, String sessionKey,
372+
OpenCodeRequestContext context,
373+
Consumer<String> deltaConsumer) {
363374
String sessionId = httpClient.ensureSession(sessionKey, context);
364375
PromptRequest promptRequest = ChatMessageMapper.toPromptRequest(request);
365376

366-
ChatStreamingResponse stream = new ChatStreamingResponse();
377+
ChatStreamingResponse stream = new ChatStreamingResponse().onDelta(deltaConsumer);
367378

368379
OpenCodeSseClient.QueueSubscription subscription =
369380
sseClient.subscribeQueueSubscription(context);

src/main/java/io/github/easy4j/opencode/OpenCodeHttpClientConfig.java

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,9 @@
1111
@Data
1212
public class OpenCodeHttpClientConfig {
1313

14+
/** 对话响应模式,默认保持兼容的完整响应模式。 */
15+
private HttpResponseMode mode = HttpResponseMode.BLOCKING;
16+
1417
/**
1518
* 是否启用 HTTP 子系统。
1619
* <p>为 false 时跳过 HTTP 客户端初始化和检查。</p>
@@ -97,14 +100,26 @@ public class OpenCodeHttpClientConfig {
97100
/** 流式事件处理线程空闲保活时间(毫秒)。 */
98101
private long streamKeepAliveMillis = 60_000L;
99102

100-
/** 单个 SSE 订阅的事件缓存上限。 */
101-
private int sseEventQueueCapacity = 1_024;
103+
/** 单个流式订阅的事件缓存上限。 */
104+
private int streamEventQueueCapacity = 1_024;
102105

103106
/**
104107
* 遇到失效连接等传输故障时是否允许 OkHttp 自动恢复。
105108
*/
106109
private boolean retryOnConnectionFailure = true;
107110

111+
/** @deprecated 使用 {@link #getStreamEventQueueCapacity()}。 */
112+
@Deprecated
113+
public int getSseEventQueueCapacity() {
114+
return streamEventQueueCapacity;
115+
}
116+
117+
/** @deprecated 使用 {@link #setStreamEventQueueCapacity(int)}。 */
118+
@Deprecated
119+
public void setSseEventQueueCapacity(int value) {
120+
this.streamEventQueueCapacity = value;
121+
}
122+
108123
/**
109124
* 是否校验 HTTPS 证书;为 false 时关闭校验(仅建议开发环境)。
110125
*/

src/main/java/io/github/easy4j/opencode/api/OpenCodeSseClient.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -118,7 +118,7 @@ public BlockingQueue<Event> subscribeQueue(OpenCodeRequestContext context) {
118118

119119
public QueueSubscription subscribeQueueSubscription(OpenCodeRequestContext context) {
120120
BlockingQueue<Event> queue = new ArrayBlockingQueue<>(
121-
Math.max(1, config.getSseEventQueueCapacity()));
121+
Math.max(1, config.getStreamEventQueueCapacity()));
122122
EventSource source = subscribe(event -> offerLatest(queue, event), context);
123123
return new QueueSubscription(queue, source);
124124
}

0 commit comments

Comments
 (0)