Skip to content

Commit a8995e1

Browse files
committed
feat: 新增 ChatRequest/ChatResponse 对齐 OpenAI 标准 + chatCompletionStream 流式
- 新增 ChatRequest/ChatResponse/ChatMessage 模型类(OpenAI 标准) - 新增 ChatMessageMapper 转换器(ChatRequest↔PromptRequest, PromptResult→ChatResponse) - 新增 chatCompletion(ChatRequest)/chatCompletionWithSession(ChatRequest, sessionKey) - 新增 chatCompletionStream(ChatRequest, sessionKey) 基于 SSE 事件流
1 parent 5f2d3f6 commit a8995e1

10 files changed

Lines changed: 552 additions & 13 deletions

File tree

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

Lines changed: 145 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -150,8 +150,152 @@ public boolean chatCompletionWithSessionAsync(PromptRequest request, String sess
150150
return httpClient.chatCompletionWithSessionAsync(request, sessionKey);
151151
}
152152

153+
// ----------------------------------------------------------------
154+
// OpenAI 标准 ChatRequest/ChatResponse(对齐 OpenClaw/Hermes)
155+
// ----------------------------------------------------------------
156+
157+
/**
158+
* 按 sessionId 发送 OpenAI 标准请求并同步等待响应。
159+
* <p>内部自动将 {@link ChatRequest} 转换为 {@link PromptRequest},
160+
* 将 {@link PromptResult} 转换为 {@link ChatResponse}。</p>
161+
*
162+
* @param sessionId 会话 ID
163+
* @param request OpenAI 标准请求
164+
* @return OpenAI 标准响应
165+
*/
166+
public ChatResponse chatCompletion(String sessionId, ChatRequest request) {
167+
PromptRequest promptRequest = io.github.hiwepy.opencode.mapper.ChatMessageMapper.toPromptRequest(request);
168+
PromptResult result = httpClient.prompt(sessionId, promptRequest);
169+
return io.github.hiwepy.opencode.mapper.ChatMessageMapper.toChatResponse(result);
170+
}
171+
172+
/**
173+
* 按 sessionKey 发送 OpenAI 标准请求并同步等待响应。
174+
* <p>与 OpenClaw/Hermes 的 {@code chatCompletionWithSession(request, sessionKey)} 完全对称。</p>
175+
*
176+
* @param request OpenAI 标准请求
177+
* @param sessionKey 会话复用 key
178+
* @return OpenAI 标准响应
179+
*/
180+
public ChatResponse chatCompletionWithSession(ChatRequest request, String sessionKey) {
181+
PromptRequest promptRequest = io.github.hiwepy.opencode.mapper.ChatMessageMapper.toPromptRequest(request);
182+
PromptResult result = httpClient.chatCompletionWithSession(promptRequest, sessionKey);
183+
return io.github.hiwepy.opencode.mapper.ChatMessageMapper.toChatResponse(result);
184+
}
185+
186+
/**
187+
* 按 sessionKey 流式发送消息,返回 {@link ChatStreamingResponse}。
188+
* <p>
189+
* 内部流程:ensureSession → promptAsync(不阻塞)→ 订阅全局 SSE 事件流 →
190+
* 按 sessionId 过滤事件,累积 text delta → session idle 时完成 future。
191+
* </p>
192+
* <p>与 Hermes 的 {@code chatCompletionStream} 对称,调用方可通过
193+
* {@link ChatStreamingResponse#onDelta(Consumer)} 注册增量回调,
194+
* 或通过 {@link ChatStreamingResponse#get()} 阻塞等待完整文本。</p>
195+
*
196+
* @param request OpenAI 标准请求
197+
* @param sessionKey 会话复用 key
198+
* @return 流式响应(CompletableFuture,完成时携带完整文本)
199+
*/
200+
public ChatStreamingResponse chatCompletionStream(ChatRequest request, String sessionKey) {
201+
String sessionId = httpClient.ensureSession(sessionKey);
202+
PromptRequest promptRequest = io.github.hiwepy.opencode.mapper.ChatMessageMapper.toPromptRequest(request);
203+
204+
ChatStreamingResponse stream = new ChatStreamingResponse();
205+
206+
// 订阅全局 SSE,按 sessionId 过滤事件
207+
java.util.concurrent.BlockingQueue<io.github.hiwepy.opencode.model.Event> queue = sseClient.subscribeQueue();
208+
209+
// 异步消费事件
210+
java.util.concurrent.CompletableFuture.runAsync(() -> {
211+
try {
212+
long deadline = System.currentTimeMillis() + (config.getLocalTimeoutSeconds() * 1000L);
213+
while (!stream.isDone() && System.currentTimeMillis() < deadline) {
214+
io.github.hiwepy.opencode.model.Event event = queue.poll(3, java.util.concurrent.TimeUnit.SECONDS);
215+
if (event == null) {
216+
continue;
217+
}
218+
219+
// 按 sessionId 过滤
220+
String eventSessionId = event.getProperties() != null
221+
? Objects.toString(event.getProperties().get("sessionID"), null) : null;
222+
if (eventSessionId == null || !eventSessionId.equals(sessionId)) {
223+
continue;
224+
}
225+
226+
String type = event.getType();
227+
if (type == null) {
228+
continue;
229+
}
230+
231+
// text delta 事件
232+
if (type.contains("text.delta") || type.contains("message.part.updated")) {
233+
String delta = extractDeltaText(event);
234+
if (delta != null && !delta.isEmpty()) {
235+
stream.acceptDelta(delta);
236+
}
237+
}
238+
239+
// session idle = 完成
240+
if (type.contains("session.status") || type.contains("session.idle")) {
241+
String status = event.getProperties() != null
242+
? Objects.toString(event.getProperties().get("status"), null) : null;
243+
if ("idle".equals(status) || type.contains("idle")) {
244+
stream.finish();
245+
return;
246+
}
247+
}
248+
249+
// session error
250+
if (type.contains("session.error")) {
251+
String error = event.getProperties() != null
252+
? Objects.toString(event.getProperties().get("error"), "unknown error") : "unknown error";
253+
stream.fail(new RuntimeException(error));
254+
return;
255+
}
256+
}
257+
if (!stream.isDone()) {
258+
stream.fail(new RuntimeException("Stream timed out for session: " + sessionId));
259+
}
260+
} catch (InterruptedException e) {
261+
Thread.currentThread().interrupt();
262+
stream.fail(e);
263+
} catch (Exception e) {
264+
stream.fail(e);
265+
}
266+
});
267+
268+
// 触发异步 prompt(不阻塞,SSE 事件驱动结果)
269+
httpClient.promptAsync(sessionId, promptRequest);
270+
271+
return stream;
272+
}
273+
274+
/**
275+
* 从事件属性中提取增量文本。
276+
*/
277+
private static String extractDeltaText(io.github.hiwepy.opencode.model.Event event) {
278+
if (event.getProperties() == null) {
279+
return null;
280+
}
281+
// 尝试 part.text
282+
Object part = event.getProperties().get("part");
283+
if (part instanceof Map) {
284+
Object text = ((Map<?, ?>) part).get("text");
285+
if (text != null) {
286+
return text.toString();
287+
}
288+
}
289+
// 尝试 delta
290+
Object delta = event.getProperties().get("delta");
291+
if (delta != null) {
292+
return delta.toString();
293+
}
294+
return null;
295+
}
296+
153297
/**
154-
* 确保指定 sessionKey 对应的 session 存在,返回其 sessionId。
298+
* 确保 sessionKey 对应的 session 存在,返回其 sessionId。
155299
*/
156300
public String ensureSession(String sessionKey) {
157301
return httpClient.ensureSession(sessionKey);
Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,107 @@
1+
package io.github.hiwepy.opencode.mapper;
2+
3+
import io.github.hiwepy.opencode.model.ChatMessage;
4+
import io.github.hiwepy.opencode.model.ChatRequest;
5+
import io.github.hiwepy.opencode.model.ChatResponse;
6+
import io.github.hiwepy.opencode.model.Part;
7+
import io.github.hiwepy.opencode.model.PromptRequest;
8+
import io.github.hiwepy.opencode.model.PromptResult;
9+
10+
import java.util.Collections;
11+
import java.util.List;
12+
import java.util.UUID;
13+
14+
/**
15+
* {@link ChatRequest}/{@link ChatResponse}(OpenAI 标准)与 {@link PromptRequest}/{@link PromptResult}(OpenCode 会话模型)互转。
16+
*
17+
* @author wandl
18+
* @since 2.7.x
19+
*/
20+
public final class ChatMessageMapper {
21+
22+
private ChatMessageMapper() {
23+
}
24+
25+
/**
26+
* ChatRequest → PromptRequest。
27+
* <p>
28+
* 取最后一条 user 消息的 content 作为 text part;model "provider/model" 拆分为 ModelRef;
29+
* agent / system 直传。
30+
* </p>
31+
*/
32+
public static PromptRequest toPromptRequest(ChatRequest chatRequest) {
33+
if (chatRequest == null) {
34+
return null;
35+
}
36+
37+
// 取最后一条 user 消息内容
38+
String text = extractUserContent(chatRequest.getMessages());
39+
40+
PromptRequest promptRequest = PromptRequest.ofText(text);
41+
42+
// model: "provider/model" → ModelRef
43+
if (chatRequest.getModel() != null && chatRequest.getModel().contains("/")) {
44+
String[] parts = chatRequest.getModel().split("/", 2);
45+
promptRequest.setModel(new PromptRequest.ModelRef(parts[0], parts[1]));
46+
}
47+
48+
// agent
49+
if (chatRequest.getAgent() != null) {
50+
promptRequest.setAgent(chatRequest.getAgent());
51+
}
52+
53+
// system prompt
54+
if (chatRequest.getSystem() != null) {
55+
promptRequest.setSystem(chatRequest.getSystem());
56+
}
57+
58+
return promptRequest;
59+
}
60+
61+
/**
62+
* PromptResult → ChatResponse。
63+
* <p>getTextContent() → choices[0].message.content。</p>
64+
*/
65+
public static ChatResponse toChatResponse(PromptResult result) {
66+
if (result == null) {
67+
return null;
68+
}
69+
70+
ChatResponse response = new ChatResponse();
71+
response.setId(UUID.randomUUID().toString());
72+
response.setObject("chat.completion");
73+
response.setCreated(System.currentTimeMillis() / 1000);
74+
75+
String content = result.getTextContent();
76+
77+
ChatMessage message = new ChatMessage("assistant", content);
78+
79+
ChatResponse.Choice choice = new ChatResponse.Choice();
80+
choice.setIndex(0);
81+
choice.setMessage(message);
82+
choice.setFinishReason("stop");
83+
84+
response.setChoices(Collections.singletonList(choice));
85+
86+
return response;
87+
}
88+
89+
/**
90+
* 从消息列表中提取最后一条 user 消息的 content。
91+
*/
92+
private static String extractUserContent(List<ChatMessage> messages) {
93+
if (messages == null || messages.isEmpty()) {
94+
return "";
95+
}
96+
// 从后往前找最后一条 user 消息
97+
for (int i = messages.size() - 1; i >= 0; i--) {
98+
ChatMessage msg = messages.get(i);
99+
if ("user".equals(msg.getRole()) && msg.getContent() != null) {
100+
return msg.getContent();
101+
}
102+
}
103+
// 兜底:取最后一条消息
104+
ChatMessage last = messages.get(messages.size() - 1);
105+
return last.getContent() != null ? last.getContent() : "";
106+
}
107+
}
Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
package io.github.hiwepy.opencode.model;
2+
3+
import com.fasterxml.jackson.annotation.JsonInclude;
4+
import lombok.Getter;
5+
import lombok.NoArgsConstructor;
6+
import lombok.Setter;
7+
8+
/**
9+
* OpenAI Chat Completions API 消息对象(对齐 OpenClaw/Hermes)。
10+
*
11+
* @author wandl
12+
* @since 2.7.x
13+
*/
14+
@Getter
15+
@Setter
16+
@NoArgsConstructor
17+
@JsonInclude(JsonInclude.Include.NON_NULL)
18+
public class ChatMessage {
19+
20+
/** 消息角色:system / user / assistant / tool。 */
21+
private String role;
22+
23+
/** 消息文本内容。 */
24+
private String content;
25+
26+
public ChatMessage(String role, String content) {
27+
this.role = role;
28+
this.content = content;
29+
}
30+
31+
/** 快捷构造 user 消息。 */
32+
public static ChatMessage user(String content) {
33+
return new ChatMessage("user", content);
34+
}
35+
36+
/** 快捷构造 system 消息。 */
37+
public static ChatMessage system(String content) {
38+
return new ChatMessage("system", content);
39+
}
40+
41+
/** 快捷构造 assistant 消息。 */
42+
public static ChatMessage assistant(String content) {
43+
return new ChatMessage("assistant", content);
44+
}
45+
}
Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
1+
package io.github.hiwepy.opencode.model;
2+
3+
import com.fasterxml.jackson.annotation.JsonInclude;
4+
import lombok.Getter;
5+
import lombok.NoArgsConstructor;
6+
import lombok.Setter;
7+
8+
import java.util.List;
9+
import java.util.Map;
10+
11+
/**
12+
* OpenAI Chat Completions API 请求体(对齐 OpenClaw/Hermes)。
13+
* <p>
14+
* OpenCode 底层使用会话模型的 {@code PromptRequest},SDK 内部自动转换:
15+
* {@code messages} 最后一条 user 消息 → {@code parts} text;
16+
* {@code model "provider/model"} → {@code ModelRef}。
17+
* </p>
18+
*
19+
* @author wandl
20+
* @since 2.7.x
21+
*/
22+
@Getter
23+
@Setter
24+
@NoArgsConstructor
25+
@JsonInclude(JsonInclude.Include.NON_NULL)
26+
public class ChatRequest {
27+
28+
/** 模型标识,格式 {@code provider/model}(如 {@code anthropic/claude-sonnet-4-5})。 */
29+
private String model;
30+
31+
/** 消息数组(OpenAI 标准格式)。 */
32+
private List<ChatMessage> messages;
33+
34+
/** 是否启用 SSE 流式响应。 */
35+
private Boolean stream;
36+
37+
/** 流式选项。 */
38+
private Map<String, Object> streamOptions;
39+
40+
/** 指定 opencode agent 名称。 */
41+
private String agent;
42+
43+
/** 系统 prompt 附加内容。 */
44+
private String system;
45+
46+
/** 最大 token 数。 */
47+
private Integer maxTokens;
48+
49+
/** 采样温度(0-2)。 */
50+
private Double temperature;
51+
52+
/** nucleus 采样参数(0-1)。 */
53+
private Double topP;
54+
55+
/** 用户标识。 */
56+
private String user;
57+
58+
/** 快捷构造:单条 user 消息。 */
59+
public static ChatRequest ofUser(String content) {
60+
ChatRequest req = new ChatRequest();
61+
req.setMessages(List.of(ChatMessage.user(content)));
62+
return req;
63+
}
64+
65+
/** 快捷构造:单条 user 消息 + 模型。 */
66+
public static ChatRequest ofUser(String content, String model) {
67+
ChatRequest req = ofUser(content);
68+
req.setModel(model);
69+
return req;
70+
}
71+
}

0 commit comments

Comments
 (0)