Skip to content

Commit b4f8ef6

Browse files
committed
feat(sse): preserve frames and classify consumer failures
1 parent 65c90b5 commit b4f8ef6

1 file changed

Lines changed: 19 additions & 4 deletions

File tree

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

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,10 @@
44
import io.github.easy4j.hermes.HermesHttpClientConfig;
55
import io.github.easy4j.hermes.HermesOkHttpClientFactory;
66
import io.github.easy4j.hermes.api.model.ChatRequest;
7+
import io.github.easy4j.hermes.api.sse.EndpointEventDecoder;
8+
import io.github.easy4j.hermes.api.sse.SseConsumerException;
79
import io.github.easy4j.hermes.api.sse.SseEvent;
10+
import io.github.easy4j.hermes.api.sse.SseFrame;
811
import io.github.easy4j.hermes.api.sse.SseQueueSubscription;
912
import io.github.easy4j.hermes.api.sse.SseSubscription;
1013
import io.github.easy4j.hermes.exception.HermesHttpException;
@@ -74,6 +77,8 @@ public class HermesSseClient implements AutoCloseable {
7477
* SSE 事件反序列化使用的 ObjectMapper。
7578
*/
7679
private final ObjectMapper mapper;
80+
/** 将保留的原始帧转换为当前兼容事件视图。 */
81+
private final EndpointEventDecoder<SseEvent> eventDecoder;
7782
/**
7883
* 执行 HTTP 请求的 OkHttpClient。
7984
*/
@@ -115,6 +120,10 @@ public HermesSseClient(HermesHttpClientConfig config, ObjectMapper objectMapper,
115120
this.config = Objects.requireNonNull(config, "config");
116121
this.mapper = Objects.isNull(objectMapper) ? new ObjectMapper()
117122
.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false) : objectMapper;
123+
this.eventDecoder = frame -> {
124+
this.mapper.readTree(frame.getData());
125+
return SseEvent.fromFrame(frame);
126+
};
118127
this.ownsHttpClient = Objects.isNull(httpClient);
119128
this.ownsDispatcher = Objects.nonNull(httpClient);
120129
OkHttpClient baseClient = this.ownsHttpClient ? HermesOkHttpClientFactory.create(config) : httpClient;
@@ -277,8 +286,10 @@ public void onEvent(EventSource eventSource, String id, String type, String data
277286
if (data == null || data.isEmpty()) {
278287
return;
279288
}
289+
SseFrame frame = new SseFrame(id, type, data, System.currentTimeMillis());
290+
SseEvent event;
280291
try {
281-
mapper.readTree(data);
292+
event = eventDecoder.decode(frame);
282293
} catch (Exception error) {
283294
if (config.getDebug().allows(HttpLogLevel.BODY)) {
284295
log.debug("Hermes SSE parse failed: label={}, data={}", label, truncate(data), error);
@@ -289,18 +300,22 @@ public void onEvent(EventSource eventSource, String id, String type, String data
289300
return;
290301
}
291302

292-
SseEvent event = new SseEvent();
293-
event.setEvent(type);
294-
event.setData(data);
295303
try {
296304
consumer.accept(event);
297305
} catch (Exception error) {
306+
SseConsumerException failure = new SseConsumerException(frame, error);
298307
if (config.getDebug().allows(HttpLogLevel.BODY)) {
299308
log.debug("Hermes SSE consumer failed: label={}, data={}", label, truncate(data), error);
300309
} else {
301310
debug(HttpLogLevel.BASIC, "Hermes SSE consumer failed: label={}, dataLength={}, error={}",
302311
label, data.length(), error.getMessage());
303312
}
313+
terminalSignal.set(true);
314+
try {
315+
onError.accept(failure);
316+
} finally {
317+
finish(subscription);
318+
}
304319
}
305320
}
306321

0 commit comments

Comments
 (0)