实现Debug EML实时Trace调试链路
This commit is contained in:
1 parent
eee37315a0
commit
8d8fae670a
20 files changed
+1571
-76
No files matched your search
+41
-1
@@ -1,5 +1,6 @@
|
||||
package cn.nianxx.thhotel.integrations.ai.superagent.common.result;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
@@ -15,6 +16,7 @@ import java.util.List;
|
||||
* @param outputTokens 输出 token 数
|
||||
* @param totalTokens 总 token 数
|
||||
* @param eventTypes SSE 事件类型列表
|
||||
* @param traceEvents SuperAgent 公开 Trace 事件列表,只包含可展示摘要字段
|
||||
*/
|
||||
public record SuperAgentOpenApiResult(
|
||||
String sessionId,
|
||||
@@ -26,6 +28,44 @@ public record SuperAgentOpenApiResult(
|
||||
Integer inputTokens,
|
||||
Integer outputTokens,
|
||||
Integer totalTokens,
|
||||
List<String> eventTypes
|
||||
List<String> eventTypes,
|
||||
@JsonProperty("trace_events")
|
||||
List<SuperAgentOpenApiTraceEvent> traceEvents
|
||||
) {
|
||||
|
||||
/**
|
||||
* 兼容未启用 Trace 的历史调用方。
|
||||
*/
|
||||
public SuperAgentOpenApiResult(
|
||||
String sessionId,
|
||||
String runId,
|
||||
String profileId,
|
||||
String profileVersionId,
|
||||
String modelName,
|
||||
String rawAnswer,
|
||||
Integer inputTokens,
|
||||
Integer outputTokens,
|
||||
Integer totalTokens,
|
||||
List<String> eventTypes) {
|
||||
this(
|
||||
sessionId,
|
||||
runId,
|
||||
profileId,
|
||||
profileVersionId,
|
||||
modelName,
|
||||
rawAnswer,
|
||||
inputTokens,
|
||||
outputTokens,
|
||||
totalTokens,
|
||||
eventTypes,
|
||||
List.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* 复制集合,避免结果对象被调用方后续修改。
|
||||
*/
|
||||
public SuperAgentOpenApiResult {
|
||||
eventTypes = eventTypes == null ? List.of() : List.copyOf(eventTypes);
|
||||
traceEvents = traceEvents == null ? List.of() : List.copyOf(traceEvents);
|
||||
}
|
||||
}
|
||||
+37
@@ -0,0 +1,37 @@
|
||||
package cn.nianxx.thhotel.integrations.ai.superagent.common.result;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
|
||||
/**
|
||||
* SuperAgent Open API 公开 Trace 事件。只保留对外策略允许的摘要字段,不保存原始工具入参、Secret 或 Cookie。
|
||||
*
|
||||
* @param event SuperAgent 公开 trace 事件类型
|
||||
* @param runId SuperAgent run ID
|
||||
* @param messageId AI message ID
|
||||
* @param toolCallId 工具调用 ID
|
||||
* @param toolName 工具名称
|
||||
* @param text 公开推理或进度摘要
|
||||
* @param inputSummary 工具入参摘要
|
||||
* @param outputSummary 工具输出摘要
|
||||
* @param status 运行或任务状态
|
||||
* @param ts SuperAgent 事件时间戳
|
||||
*/
|
||||
public record SuperAgentOpenApiTraceEvent(
|
||||
String event,
|
||||
@JsonProperty("run_id")
|
||||
String runId,
|
||||
@JsonProperty("message_id")
|
||||
String messageId,
|
||||
@JsonProperty("tool_call_id")
|
||||
String toolCallId,
|
||||
@JsonProperty("tool_name")
|
||||
String toolName,
|
||||
String text,
|
||||
@JsonProperty("input_summary")
|
||||
String inputSummary,
|
||||
@JsonProperty("output_summary")
|
||||
String outputSummary,
|
||||
String status,
|
||||
String ts
|
||||
) {
|
||||
}
|
||||
+13
-1
@@ -2,6 +2,8 @@ package cn.nianxx.thhotel.integrations.ai.superagent.service;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.request.SuperAgentMailDebugRequest;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* SuperAgent Open API 客户端端口。业务层只依赖该接口,不直接拼 HTTP 或解析 SSE。
|
||||
@@ -11,5 +13,15 @@ public interface SuperAgentOpenApiClient {
|
||||
/**
|
||||
* 创建 SuperAgent session 并发送邮件 Debug 消息,返回最终 AI 回答。
|
||||
*/
|
||||
SuperAgentOpenApiResult invokeMailDebug(SuperAgentMailDebugRequest request);
|
||||
default SuperAgentOpenApiResult invokeMailDebug(SuperAgentMailDebugRequest request) {
|
||||
return invokeMailDebug(request, traceEvent -> {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 创建 SuperAgent session 并发送邮件 Debug 消息,边解析边回调公开 Trace 事件。
|
||||
*/
|
||||
SuperAgentOpenApiResult invokeMailDebug(
|
||||
SuperAgentMailDebugRequest request,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer);
|
||||
}
|
||||
+37
-8
@@ -2,16 +2,22 @@ package cn.nianxx.thhotel.integrations.ai.superagent.service.impl;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.request.SuperAgentMailDebugRequest;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.SuperAgentOpenApiClient;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import java.io.InputStreamReader;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.UUID;
|
||||
import java.util.function.Consumer;
|
||||
import org.springframework.http.client.SimpleClientHttpRequestFactory;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.util.StreamUtils;
|
||||
import org.springframework.web.client.RestClient;
|
||||
import org.springframework.web.client.RestClientResponseException;
|
||||
|
||||
/**
|
||||
* SuperAgent Open API HTTP 客户端实现。负责创建 session、发送 SSE 消息和解析最终回答。
|
||||
@@ -39,7 +45,9 @@ public class SuperAgentOpenApiClientImpl implements SuperAgentOpenApiClient {
|
||||
* 调用 SuperAgent Open API 邮件 Debug 流程;配置缺失时直接失败,避免静默跳过真实调用。
|
||||
*/
|
||||
@Override
|
||||
public SuperAgentOpenApiResult invokeMailDebug(SuperAgentMailDebugRequest request) {
|
||||
public SuperAgentOpenApiResult invokeMailDebug(
|
||||
SuperAgentMailDebugRequest request,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
validateProperties();
|
||||
try {
|
||||
RestClient restClient = RestClient.builder()
|
||||
@@ -47,8 +55,7 @@ public class SuperAgentOpenApiClientImpl implements SuperAgentOpenApiClient {
|
||||
.requestFactory(requestFactory())
|
||||
.build();
|
||||
String sessionId = createSession(restClient, request);
|
||||
String sseBody = sendMessage(restClient, sessionId, request);
|
||||
return sseParser.parse(sessionId, sseBody);
|
||||
return sendMessage(restClient, sessionId, request, traceConsumer);
|
||||
} catch (SuperAgentOpenApiException exception) {
|
||||
throw exception;
|
||||
} catch (Exception exception) {
|
||||
@@ -86,24 +93,46 @@ public class SuperAgentOpenApiClientImpl implements SuperAgentOpenApiClient {
|
||||
}
|
||||
|
||||
/**
|
||||
* 发送 Debug 邮件消息并读取 SSE 文本响应。
|
||||
* 发送 Debug 邮件消息并流式解析 SSE 文本响应。
|
||||
*/
|
||||
private String sendMessage(RestClient restClient, String sessionId, SuperAgentMailDebugRequest request) {
|
||||
private SuperAgentOpenApiResult sendMessage(
|
||||
RestClient restClient,
|
||||
String sessionId,
|
||||
SuperAgentMailDebugRequest request,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
String csrfToken = UUID.randomUUID().toString();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("message", request.message());
|
||||
body.put("idempotency_key", request.idempotencyKey() + "-message");
|
||||
body.put("metadata", request.metadata());
|
||||
return restClient.post()
|
||||
.uri("/api/open/agent-sessions/{sessionId}/messages/stream", sessionId)
|
||||
.uri(uriBuilder -> uriBuilder
|
||||
.path("/api/open/agent-sessions/{sessionId}/messages/stream")
|
||||
.queryParam("include_trace", "true")
|
||||
.build(sessionId))
|
||||
.contentType(MediaType.APPLICATION_JSON)
|
||||
.accept(MediaType.TEXT_EVENT_STREAM)
|
||||
.header("Authorization", "Bearer " + properties.getApiKey())
|
||||
.header("X-CSRF-Token", csrfToken)
|
||||
.header("Cookie", "csrf_token=" + csrfToken)
|
||||
.body(body)
|
||||
.retrieve()
|
||||
.body(String.class);
|
||||
.exchange((clientRequest, clientResponse) -> {
|
||||
if (!clientResponse.getStatusCode().is2xxSuccessful()) {
|
||||
byte[] responseBody = StreamUtils.copyToByteArray(clientResponse.getBody());
|
||||
throw new RestClientResponseException(
|
||||
"SuperAgent Open API HTTP 调用失败。",
|
||||
clientResponse.getStatusCode(),
|
||||
clientResponse.getStatusText(),
|
||||
clientResponse.getHeaders(),
|
||||
responseBody,
|
||||
StandardCharsets.UTF_8);
|
||||
}
|
||||
try (InputStreamReader reader = new InputStreamReader(
|
||||
clientResponse.getBody(),
|
||||
StandardCharsets.UTF_8)) {
|
||||
return sseParser.parse(sessionId, reader, traceConsumer);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+215
-29
@@ -1,12 +1,19 @@
|
||||
package cn.nianxx.thhotel.integrations.ai.superagent.service.impl;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import java.io.BufferedReader;
|
||||
import java.io.IOException;
|
||||
import java.io.Reader;
|
||||
import java.io.StringReader;
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.regex.Pattern;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
@@ -15,6 +22,13 @@ import org.springframework.stereotype.Component;
|
||||
@Component
|
||||
public class SuperAgentOpenApiSseParser {
|
||||
|
||||
private static final Pattern JSON_SECRET_PATTERN = Pattern.compile(
|
||||
"(?i)(\"(?:api[_-]?key|token|secret|password|cookie|authorization)\"\\s*:\\s*\")[^\"]*(\")");
|
||||
private static final Pattern HEADER_SECRET_PATTERN = Pattern.compile(
|
||||
"(?i)((?:authorization|cookie)\\s*:\\s*)[^\\s,;]+");
|
||||
private static final Pattern TEXT_SECRET_PATTERN = Pattern.compile(
|
||||
"(?i)((?:api[_-]?key|token|secret|password)\\s*[=:]\\s*)[^\\s,;}]+");
|
||||
|
||||
private final ObjectMapper objectMapper;
|
||||
|
||||
/**
|
||||
@@ -35,18 +49,50 @@ public class SuperAgentOpenApiSseParser {
|
||||
* 解析 SSE 文本,返回最终 AI 回答和调用元数据。
|
||||
*/
|
||||
public SuperAgentOpenApiResult parse(String sessionId, String sseBody) {
|
||||
return parse(sessionId, sseBody, traceEvent -> {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 解析 SSE 文本,边解析边回调 SuperAgent 公开 Trace 事件。
|
||||
*/
|
||||
public SuperAgentOpenApiResult parse(
|
||||
String sessionId,
|
||||
String sseBody,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
if (sseBody == null || sseBody.isBlank()) {
|
||||
throw new SuperAgentOpenApiException("SuperAgent SSE 响应为空。");
|
||||
}
|
||||
return parse(sessionId, new StringReader(sseBody), traceConsumer);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从 Reader 流式解析 SSE,适用于 Open API 长连接响应。
|
||||
*/
|
||||
public SuperAgentOpenApiResult parse(
|
||||
String sessionId,
|
||||
Reader reader,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
if (reader == null) {
|
||||
throw new SuperAgentOpenApiException("SuperAgent SSE 响应为空。");
|
||||
}
|
||||
Set<String> eventTypes = new LinkedHashSet<>();
|
||||
ParsedState state = new ParsedState();
|
||||
for (SseEvent event : splitEvents(sseBody)) {
|
||||
eventTypes.add(event.eventType());
|
||||
consumeEvent(event, state);
|
||||
if (state.endSeen) {
|
||||
break;
|
||||
}
|
||||
try {
|
||||
parseLines(reader, event -> {
|
||||
eventTypes.add(event.eventType());
|
||||
consumeEvent(event, state, safeTraceConsumer(traceConsumer));
|
||||
}, state);
|
||||
} catch (IOException exception) {
|
||||
throw new SuperAgentOpenApiException("SuperAgent SSE 读取失败。", exception);
|
||||
}
|
||||
return buildResult(sessionId, eventTypes, state);
|
||||
}
|
||||
|
||||
/**
|
||||
* 汇总解析状态并生成结果对象。
|
||||
*/
|
||||
private SuperAgentOpenApiResult buildResult(String sessionId, Set<String> eventTypes, ParsedState state) {
|
||||
if (!state.endSeen) {
|
||||
throw new SuperAgentOpenApiException("SuperAgent SSE 未收到结束事件。");
|
||||
}
|
||||
@@ -70,13 +116,17 @@ public class SuperAgentOpenApiSseParser {
|
||||
state.inputTokens,
|
||||
state.outputTokens,
|
||||
state.totalTokens,
|
||||
List.copyOf(eventTypes));
|
||||
List.copyOf(eventTypes),
|
||||
List.copyOf(state.traceEvents));
|
||||
}
|
||||
|
||||
/**
|
||||
* 消费单个 SSE 事件。
|
||||
*/
|
||||
private void consumeEvent(SseEvent event, ParsedState state) {
|
||||
private void consumeEvent(
|
||||
SseEvent event,
|
||||
ParsedState state,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
try {
|
||||
if ("end".equals(event.eventType())) {
|
||||
state.endSeen = true;
|
||||
@@ -94,6 +144,9 @@ public class SuperAgentOpenApiSseParser {
|
||||
if ("values".equals(event.eventType())) {
|
||||
consumeValuesEvent(data, state);
|
||||
}
|
||||
if ("trace".equals(event.eventType())) {
|
||||
consumeTraceEvent(data, state, traceConsumer);
|
||||
}
|
||||
} catch (Exception exception) {
|
||||
throw new SuperAgentOpenApiException("SuperAgent SSE JSON 解析失败。", exception);
|
||||
}
|
||||
@@ -159,31 +212,96 @@ public class SuperAgentOpenApiSseParser {
|
||||
}
|
||||
|
||||
/**
|
||||
* 拆分 SSE 事件块,支持多行 data。
|
||||
* 消费公开 Trace 事件,只提取可展示摘要字段。
|
||||
*/
|
||||
private List<SseEvent> splitEvents(String sseBody) {
|
||||
String[] blocks = sseBody.split("\\R\\s*\\R");
|
||||
List<SseEvent> events = new ArrayList<>();
|
||||
for (String block : blocks) {
|
||||
String eventType = "message";
|
||||
StringBuilder data = new StringBuilder();
|
||||
for (String line : block.split("\\R")) {
|
||||
if (line.startsWith("event:")) {
|
||||
eventType = line.substring("event:".length()).trim();
|
||||
continue;
|
||||
}
|
||||
if (line.startsWith("data:")) {
|
||||
if (!data.isEmpty()) {
|
||||
data.append('\n');
|
||||
}
|
||||
data.append(line.substring("data:".length()).trim());
|
||||
}
|
||||
private void consumeTraceEvent(
|
||||
JsonNode data,
|
||||
ParsedState state,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
JsonNode trace = data.path("data").isObject() && data.path("event").isMissingNode()
|
||||
? data.path("data")
|
||||
: data;
|
||||
JsonNode toolCalls = trace.path("tool_calls");
|
||||
if (toolCalls.isArray() && !toolCalls.isEmpty()) {
|
||||
for (JsonNode toolCall : toolCalls) {
|
||||
emitTraceEvent(traceEvent(trace, toolCall), state, traceConsumer);
|
||||
}
|
||||
if (!data.isEmpty() || "end".equals(eventType)) {
|
||||
events.add(new SseEvent(eventType, data.toString()));
|
||||
return;
|
||||
}
|
||||
emitTraceEvent(traceEvent(trace, trace), state, traceConsumer);
|
||||
}
|
||||
|
||||
/**
|
||||
* 构造公开 Trace 事件对象。
|
||||
*/
|
||||
private SuperAgentOpenApiTraceEvent traceEvent(JsonNode trace, JsonNode detail) {
|
||||
return new SuperAgentOpenApiTraceEvent(
|
||||
safeTraceText(text(trace, "event", null)),
|
||||
safeTraceText(text(trace, "run_id", null)),
|
||||
safeTraceText(text(trace, "message_id", null)),
|
||||
safeTraceText(text(detail, "tool_call_id", text(trace, "tool_call_id", null))),
|
||||
safeTraceText(text(detail, "name", text(trace, "name", null))),
|
||||
safeTraceText(text(trace, "text", null)),
|
||||
safeTraceText(summary(detail.path("input"), summary(trace.path("input"), null))),
|
||||
safeTraceText(summary(detail.path("output"), summary(trace.path("output"), null))),
|
||||
safeTraceText(text(trace, "status", null)),
|
||||
safeTraceText(text(trace, "ts", null)));
|
||||
}
|
||||
|
||||
/**
|
||||
* 保存并回调 Trace 事件。
|
||||
*/
|
||||
private void emitTraceEvent(
|
||||
SuperAgentOpenApiTraceEvent traceEvent,
|
||||
ParsedState state,
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
if (traceEvent.event() == null || traceEvent.event().isBlank()) {
|
||||
return;
|
||||
}
|
||||
state.traceEvents.add(traceEvent);
|
||||
traceConsumer.accept(traceEvent);
|
||||
}
|
||||
|
||||
/**
|
||||
* 流式拆分 SSE 事件块,支持多行 data。
|
||||
*/
|
||||
private void parseLines(Reader reader, Consumer<SseEvent> consumer, ParsedState state) throws IOException {
|
||||
BufferedReader bufferedReader = reader instanceof BufferedReader existingReader
|
||||
? existingReader
|
||||
: new BufferedReader(reader);
|
||||
SseEventBuilder builder = new SseEventBuilder();
|
||||
String line;
|
||||
while ((line = bufferedReader.readLine()) != null) {
|
||||
if (line.isBlank()) {
|
||||
flushEvent(builder, consumer);
|
||||
if (state.endSeen) {
|
||||
return;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (line.startsWith(":")) {
|
||||
continue;
|
||||
}
|
||||
if (line.startsWith("event:")) {
|
||||
builder.eventType(line.substring("event:".length()).trim());
|
||||
continue;
|
||||
}
|
||||
if (line.startsWith("data:")) {
|
||||
builder.appendData(line.substring("data:".length()).trim());
|
||||
}
|
||||
}
|
||||
return events;
|
||||
flushEvent(builder, consumer);
|
||||
}
|
||||
|
||||
/**
|
||||
* 输出当前累积事件并重置 builder。
|
||||
*/
|
||||
private void flushEvent(SseEventBuilder builder, Consumer<SseEvent> consumer) {
|
||||
SseEvent event = builder.build();
|
||||
if (event != null) {
|
||||
consumer.accept(event);
|
||||
}
|
||||
builder.reset();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -205,6 +323,42 @@ public class SuperAgentOpenApiSseParser {
|
||||
return value.asInt();
|
||||
}
|
||||
|
||||
/**
|
||||
* 读取公开 summary 字段。
|
||||
*/
|
||||
private String summary(JsonNode node, String fallback) {
|
||||
if (node == null || node.isMissingNode() || node.isNull()) {
|
||||
return fallback;
|
||||
}
|
||||
String value = text(node, "summary", null);
|
||||
if (value != null) {
|
||||
return value;
|
||||
}
|
||||
return node.isTextual() ? node.asText() : fallback;
|
||||
}
|
||||
|
||||
/**
|
||||
* Trace 展示字段做短文本脱敏与截断。
|
||||
*/
|
||||
private String safeTraceText(String value) {
|
||||
if (value == null || value.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
String sanitized = JSON_SECRET_PATTERN.matcher(value).replaceAll("$1***$2");
|
||||
sanitized = HEADER_SECRET_PATTERN.matcher(sanitized).replaceAll("$1***");
|
||||
sanitized = TEXT_SECRET_PATTERN.matcher(sanitized).replaceAll("$1***");
|
||||
return sanitized.length() > 2048 ? sanitized.substring(0, 2048) : sanitized;
|
||||
}
|
||||
|
||||
/**
|
||||
* Trace 消费者为空时使用 no-op。
|
||||
*/
|
||||
private Consumer<SuperAgentOpenApiTraceEvent> safeTraceConsumer(
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer) {
|
||||
return traceConsumer == null ? traceEvent -> {
|
||||
} : traceConsumer;
|
||||
}
|
||||
|
||||
/**
|
||||
* 读取 AI content,兼容字符串、文本片段数组和简单文本对象。
|
||||
*/
|
||||
@@ -252,6 +406,37 @@ public class SuperAgentOpenApiSseParser {
|
||||
private record SseEvent(String eventType, String data) {
|
||||
}
|
||||
|
||||
/**
|
||||
* SSE 事件累积器。
|
||||
*/
|
||||
private static final class SseEventBuilder {
|
||||
private String eventType = "message";
|
||||
private final StringBuilder data = new StringBuilder();
|
||||
|
||||
private void eventType(String eventType) {
|
||||
this.eventType = eventType == null || eventType.isBlank() ? "message" : eventType;
|
||||
}
|
||||
|
||||
private void appendData(String value) {
|
||||
if (!data.isEmpty()) {
|
||||
data.append('\n');
|
||||
}
|
||||
data.append(value);
|
||||
}
|
||||
|
||||
private SseEvent build() {
|
||||
if (data.isEmpty() && !"end".equals(eventType)) {
|
||||
return null;
|
||||
}
|
||||
return new SseEvent(eventType, data.toString());
|
||||
}
|
||||
|
||||
private void reset() {
|
||||
eventType = "message";
|
||||
data.setLength(0);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* SSE 解析过程中的可变状态。
|
||||
*/
|
||||
@@ -270,5 +455,6 @@ public class SuperAgentOpenApiSseParser {
|
||||
private Integer fallbackInputTokens;
|
||||
private Integer fallbackOutputTokens;
|
||||
private Integer fallbackTotalTokens;
|
||||
private final List<SuperAgentOpenApiTraceEvent> traceEvents = new ArrayList<>();
|
||||
}
|
||||
}
|
||||
+4
@@ -1,5 +1,6 @@
|
||||
package cn.nianxx.thhotel.platform.debug.common.result;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import java.util.List;
|
||||
@@ -25,6 +26,7 @@ import java.util.Map;
|
||||
* @param superagentRunId SuperAgent run ID
|
||||
* @param superagentRawAnswer SuperAgent 原始最终回答
|
||||
* @param superagentParsedJson 后端解析出的 JSON
|
||||
* @param superagentTraceEvents SuperAgent 公开 Trace 事件列表
|
||||
* @param warnings 可展示的安全警告
|
||||
* @param status Debug 运行状态
|
||||
*/
|
||||
@@ -63,6 +65,8 @@ public record DebugEmlSuperAgentRunResult(
|
||||
String superagentRawAnswer,
|
||||
@JsonProperty("superagent_parsed_json")
|
||||
JsonNode superagentParsedJson,
|
||||
@JsonProperty("superagent_trace_events")
|
||||
List<SuperAgentOpenApiTraceEvent> superagentTraceEvents,
|
||||
List<String> warnings,
|
||||
String status
|
||||
) {
|
||||
|
||||
+21
@@ -3,6 +3,7 @@ package cn.nianxx.thhotel.platform.debug.control;
|
||||
import cn.nianxx.thhotel.platform.debug.common.result.DebugEmlSuperAgentRunResult;
|
||||
import cn.nianxx.thhotel.platform.debug.service.DebugEmlSuperAgentRunService;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.http.HttpHeaders;
|
||||
import org.springframework.http.HttpStatus;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.http.ResponseEntity;
|
||||
@@ -12,6 +13,7 @@ import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RequestParam;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
import org.springframework.web.servlet.mvc.method.annotation.StreamingResponseBody;
|
||||
|
||||
/**
|
||||
* Debug EML 上传 Controller。该入口只用于受控调试,不创建业务订单或任务。
|
||||
@@ -41,4 +43,23 @@ public class DebugEmlSuperAgentController {
|
||||
@RequestParam(name = "run_label", required = false) String runLabel) {
|
||||
return ResponseEntity.status(HttpStatus.CREATED).body(runService.uploadAndRun(accessKey, file, hotelId, runLabel));
|
||||
}
|
||||
|
||||
/**
|
||||
* 上传单封 .eml 邮件并以 SSE 实时返回本系统阶段、SuperAgent 公开 Trace 和最终回答。
|
||||
*/
|
||||
@PostMapping(path = "/stream", consumes = MediaType.MULTIPART_FORM_DATA_VALUE, produces = MediaType.TEXT_EVENT_STREAM_VALUE)
|
||||
public ResponseEntity<StreamingResponseBody> uploadStream(
|
||||
@RequestHeader(name = "X-TH-Hotel-Debug-Upload-Key", required = false) String accessKey,
|
||||
@RequestParam("file") MultipartFile file,
|
||||
@RequestParam(name = "hotel_id", required = false) String hotelId,
|
||||
@RequestParam(name = "run_label", required = false) String runLabel) {
|
||||
runService.validateUploadAccessKey(accessKey);
|
||||
StreamingResponseBody responseBody = outputStream ->
|
||||
runService.uploadAndRunStream(accessKey, file, hotelId, runLabel, outputStream);
|
||||
return ResponseEntity.ok()
|
||||
.header(HttpHeaders.CACHE_CONTROL, "no-cache")
|
||||
.header("X-Accel-Buffering", "no")
|
||||
.contentType(MediaType.TEXT_EVENT_STREAM)
|
||||
.body(responseBody);
|
||||
}
|
||||
}
|
||||
+16
@@ -1,6 +1,7 @@
|
||||
package cn.nianxx.thhotel.platform.debug.service;
|
||||
|
||||
import cn.nianxx.thhotel.platform.debug.common.result.DebugEmlSuperAgentRunResult;
|
||||
import java.io.OutputStream;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
/**
|
||||
@@ -8,6 +9,11 @@ import org.springframework.web.multipart.MultipartFile;
|
||||
*/
|
||||
public interface DebugEmlSuperAgentRunService {
|
||||
|
||||
/**
|
||||
* 校验 Debug EML 上传访问口令。流式入口需要在响应开始前完成校验,才能保留标准 HTTP 401。
|
||||
*/
|
||||
void validateUploadAccessKey(String accessKey);
|
||||
|
||||
/**
|
||||
* 上传并处理单封 EML,写入 SourceMessage 后调用 SuperAgent Open API。
|
||||
*/
|
||||
@@ -16,4 +22,14 @@ public interface DebugEmlSuperAgentRunService {
|
||||
MultipartFile file,
|
||||
String hotelId,
|
||||
String runLabel);
|
||||
|
||||
/**
|
||||
* 上传并处理单封 EML,按 text/event-stream 实时写出内部阶段、SuperAgent Trace 和最终结果。
|
||||
*/
|
||||
void uploadAndRunStream(
|
||||
String accessKey,
|
||||
MultipartFile file,
|
||||
String hotelId,
|
||||
String runLabel,
|
||||
OutputStream outputStream);
|
||||
}
|
||||
+235
-12
@@ -2,6 +2,7 @@ package cn.nianxx.thhotel.platform.debug.service.impl;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.request.SuperAgentMailDebugRequest;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.SuperAgentOpenApiClient;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.impl.SuperAgentOpenApiException;
|
||||
import cn.nianxx.thhotel.integrations.storage.aliyunoss.common.request.ObjectStoragePutRequest;
|
||||
@@ -34,6 +35,8 @@ import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.databind.node.ObjectNode;
|
||||
import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.net.SocketTimeoutException;
|
||||
import java.net.URLDecoder;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
@@ -48,6 +51,7 @@ import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import org.springframework.http.HttpStatus;
|
||||
@@ -106,6 +110,14 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
this.hotelContextService = hotelContextService;
|
||||
}
|
||||
|
||||
/**
|
||||
* 校验 Debug EML 上传访问口令。供 Controller 在流式响应开始前复用,避免错误被写成 SSE 事件后丢失 HTTP 状态。
|
||||
*/
|
||||
@Override
|
||||
public void validateUploadAccessKey(String accessKey) {
|
||||
validateAccessKey(accessKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理单封 Debug EML 上传;第一版只展示 SuperAgent 结果,不创建订单或任务。
|
||||
*/
|
||||
@@ -131,7 +143,7 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
DebugEmlSuperAgentRunStatus.CREATED.name(),
|
||||
now));
|
||||
try {
|
||||
return doUploadAndRun(runId, normalizedHotelId, trimToNull(runLabel), safeFileName, emlBytes, now);
|
||||
return doUploadAndRun(runId, normalizedHotelId, trimToNull(runLabel), safeFileName, emlBytes, now, null);
|
||||
} catch (DebugEmlSuperAgentException exception) {
|
||||
markFailed(runId, exception.getMessage(), statusForException(exception), nowUtc());
|
||||
throw exception;
|
||||
@@ -166,6 +178,80 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理单封 Debug EML 上传,并实时输出安全的调试 SSE 事件。
|
||||
*/
|
||||
@Override
|
||||
public void uploadAndRunStream(
|
||||
String accessKey,
|
||||
MultipartFile file,
|
||||
String hotelId,
|
||||
String runLabel,
|
||||
OutputStream outputStream) {
|
||||
OutputStreamDebugEmlStreamSink streamSink = new OutputStreamDebugEmlStreamSink(outputStream, objectMapper);
|
||||
Long runId = null;
|
||||
try {
|
||||
validateAccessKey(accessKey);
|
||||
String normalizedHotelId = normalizeHotelId(hotelId);
|
||||
validateFile(file);
|
||||
byte[] emlBytes = readFileBytes(file);
|
||||
String safeFileName = safeFileName(file.getOriginalFilename(), "debug-email.eml");
|
||||
if (!safeFileName.toLowerCase(Locale.ROOT).endsWith(".eml")) {
|
||||
throw new DebugEmlSuperAgentException(HttpStatus.BAD_REQUEST, "INVALID_FILE_TYPE", "只支持上传 .eml 邮件文件。");
|
||||
}
|
||||
|
||||
LocalDateTime now = nowUtc();
|
||||
runId = runRepository.insert(new DebugEmlSuperAgentRunDraft(
|
||||
normalizedHotelId,
|
||||
trimToNull(runLabel),
|
||||
DebugEmlSuperAgentRunStatus.CREATED.name(),
|
||||
now));
|
||||
DebugEmlSuperAgentRunResult result = doUploadAndRun(
|
||||
runId,
|
||||
normalizedHotelId,
|
||||
trimToNull(runLabel),
|
||||
safeFileName,
|
||||
emlBytes,
|
||||
now,
|
||||
streamSink);
|
||||
streamSink.result(result);
|
||||
streamSink.done();
|
||||
} catch (DebugEmlSuperAgentException exception) {
|
||||
if (runId != null) {
|
||||
markFailed(runId, exception.getMessage(), statusForException(exception), nowUtc());
|
||||
}
|
||||
streamSink.error(runId, exception.getErrorCode(), exception.getMessage(), statusForException(exception).name());
|
||||
} catch (EmlMessageParseException exception) {
|
||||
if (runId != null) {
|
||||
markFailed(runId, "EML 邮件解析失败。", DebugEmlSuperAgentRunStatus.FAILED, nowUtc());
|
||||
}
|
||||
streamSink.error(runId, "EML_PARSE_FAILED", "EML 邮件解析失败。", DebugEmlSuperAgentRunStatus.FAILED.name());
|
||||
} catch (ObjectStorageException exception) {
|
||||
if (runId != null) {
|
||||
markFailed(runId, "OSS 上传失败。", DebugEmlSuperAgentRunStatus.FAILED, nowUtc());
|
||||
}
|
||||
streamSink.error(runId, "OSS_UPLOAD_FAILED", "OSS 上传失败。", DebugEmlSuperAgentRunStatus.FAILED.name());
|
||||
} catch (SuperAgentOpenApiException exception) {
|
||||
if (runId != null) {
|
||||
markFailed(runId, superAgentFailureSummary(exception), DebugEmlSuperAgentRunStatus.SUPERAGENT_FAILED, nowUtc());
|
||||
}
|
||||
streamSink.error(
|
||||
runId,
|
||||
"SUPERAGENT_OPEN_API_FAILED",
|
||||
"SuperAgent 调用失败。",
|
||||
DebugEmlSuperAgentRunStatus.SUPERAGENT_FAILED.name());
|
||||
} catch (Exception exception) {
|
||||
if (runId != null) {
|
||||
markFailed(runId, "Debug EML 上传处理失败。", DebugEmlSuperAgentRunStatus.FAILED, nowUtc());
|
||||
}
|
||||
streamSink.error(
|
||||
runId,
|
||||
"DEBUG_EML_RUN_FAILED",
|
||||
"Debug EML 上传处理失败。",
|
||||
DebugEmlSuperAgentRunStatus.FAILED.name());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 解析 Debug 上传酒店上下文。第一版可不传 hotel_id,由单酒店系统上下文兜底。
|
||||
*/
|
||||
@@ -250,12 +336,14 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
String runLabel,
|
||||
String safeFileName,
|
||||
byte[] emlBytes,
|
||||
LocalDateTime createdAt) throws Exception {
|
||||
LocalDateTime createdAt,
|
||||
DebugEmlSuperAgentStreamSink streamSink) throws Exception {
|
||||
String sha256 = sha256(emlBytes);
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.PARSING_EML,
|
||||
"阶段:解析 EML 邮件,文件名:" + safeFileName + ",大小:" + emlBytes.length + " bytes。");
|
||||
"阶段:解析 EML 邮件,文件名:" + safeFileName + ",大小:" + emlBytes.length + " bytes。",
|
||||
streamSink);
|
||||
ParsedEmlMessage parsed = parseService.parse(emlBytes, safeFileName);
|
||||
String originalMessageId = parsed.messageId();
|
||||
String originalConversationId = parsed.conversationId();
|
||||
@@ -267,7 +355,8 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.UPLOADING_ORIGINAL_EML,
|
||||
"阶段:上传原始 EML 到 OSS,文件名:" + safeFileName + ",大小:" + emlBytes.length + " bytes。");
|
||||
"阶段:上传原始 EML 到 OSS,文件名:" + safeFileName + ",大小:" + emlBytes.length + " bytes。",
|
||||
streamSink);
|
||||
uploadedMedia.add(uploadOriginalEml(runId, safeFileName, emlBytes, createdAt));
|
||||
int inlineIndex = 1;
|
||||
int attachmentIndex = 1;
|
||||
@@ -276,13 +365,15 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.UPLOADING_MEDIA,
|
||||
mediaUploadStageSummary(mediaItem, "inline", inlineIndex));
|
||||
mediaUploadStageSummary(mediaItem, "inline", inlineIndex),
|
||||
streamSink);
|
||||
uploadedMedia.add(uploadParsedMedia(runId, createdAt, mediaItem, "inline", inlineIndex++));
|
||||
} else {
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.UPLOADING_MEDIA,
|
||||
mediaUploadStageSummary(mediaItem, "attachments", attachmentIndex));
|
||||
mediaUploadStageSummary(mediaItem, "attachments", attachmentIndex),
|
||||
streamSink);
|
||||
uploadedMedia.add(uploadParsedMedia(runId, createdAt, mediaItem, "attachments", attachmentIndex++));
|
||||
}
|
||||
}
|
||||
@@ -290,7 +381,8 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.BUILDING_SOURCE_MESSAGE,
|
||||
"阶段:构造 SourceMessage payload。");
|
||||
"阶段:构造 SourceMessage payload。",
|
||||
streamSink);
|
||||
String htmlWithOssUrls = replaceCidReferences(parsed.htmlBody(), uploadedMedia, warnings);
|
||||
String htmlBodySanitized = htmlSanitizerService.sanitizeHtml(htmlWithOssUrls);
|
||||
String htmlRenderMode = htmlSanitizerService.htmlRenderMode(htmlWithOssUrls);
|
||||
@@ -308,7 +400,8 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.CAPTURING_SOURCE_MESSAGE,
|
||||
"阶段:写入 SourceMessage Inbox。");
|
||||
"阶段:写入 SourceMessage Inbox。",
|
||||
streamSink);
|
||||
SourceMessageCaptureResult captureResult = sourceMessageCaptureService.capture(new CaptureSourceMessageCommand(
|
||||
hotelId,
|
||||
SOURCE_PROVIDER,
|
||||
@@ -338,15 +431,29 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
markStage(
|
||||
runId,
|
||||
DebugEmlSuperAgentRunStatus.CALLING_SUPERAGENT,
|
||||
"阶段:调用 SuperAgent Open API,source_message_id=" + captureResult.inboxId() + "。");
|
||||
SuperAgentOpenApiResult superAgentResult = superAgentOpenApiClient.invokeMailDebug(new SuperAgentMailDebugRequest(
|
||||
"阶段:调用 SuperAgent Open API,source_message_id=" + captureResult.inboxId() + "。",
|
||||
streamSink);
|
||||
SuperAgentMailDebugRequest superAgentRequest = new SuperAgentMailDebugRequest(
|
||||
buildSuperAgentMessage(payloadJson),
|
||||
"debug-eml-" + runId,
|
||||
Map.of(
|
||||
"source", "th-hotel-debug-eml-upload",
|
||||
"debug_run_id", runId.toString(),
|
||||
"source_message_id", captureResult.inboxId().toString(),
|
||||
"hotel_id", hotelId)));
|
||||
"hotel_id", hotelId));
|
||||
List<SuperAgentOpenApiTraceEvent> streamedTraceEvents = new ArrayList<>();
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer = traceEvent -> {
|
||||
streamedTraceEvents.add(traceEvent);
|
||||
if (streamSink != null) {
|
||||
streamSink.trace(traceEvent);
|
||||
}
|
||||
};
|
||||
SuperAgentOpenApiResult superAgentResult = streamSink == null
|
||||
? superAgentOpenApiClient.invokeMailDebug(superAgentRequest)
|
||||
: superAgentOpenApiClient.invokeMailDebug(superAgentRequest, traceConsumer);
|
||||
List<SuperAgentOpenApiTraceEvent> traceEvents = streamedTraceEvents.isEmpty()
|
||||
? superAgentResult.traceEvents()
|
||||
: List.copyOf(streamedTraceEvents);
|
||||
JsonNode parsedJson = parseSuperAgentJson(superAgentResult.rawAnswer(), warnings);
|
||||
String parsedJsonText = parsedJson == null ? null : objectMapper.writeValueAsString(parsedJson);
|
||||
DebugEmlSuperAgentRunStatus status = DebugEmlSuperAgentRunStatus.SUPERAGENT_SUCCEEDED;
|
||||
@@ -391,6 +498,7 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
superAgentResult.runId(),
|
||||
superAgentResult.rawAnswer(),
|
||||
parsedJson,
|
||||
traceEvents,
|
||||
List.copyOf(warnings),
|
||||
status.name());
|
||||
}
|
||||
@@ -693,14 +801,29 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
Long runId,
|
||||
DebugEmlSuperAgentRunStatus status,
|
||||
String safeSummary) {
|
||||
markStage(runId, status, safeSummary, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* 标记 Debug 运行中的当前阶段,并在流式调试入口实时输出阶段事件。
|
||||
*/
|
||||
private void markStage(
|
||||
Long runId,
|
||||
DebugEmlSuperAgentRunStatus status,
|
||||
String safeSummary,
|
||||
DebugEmlSuperAgentStreamSink streamSink) {
|
||||
if (runId == null) {
|
||||
return;
|
||||
}
|
||||
String truncatedSummary = truncate(safeSummary, 512);
|
||||
runRepository.updateStatus(new DebugEmlSuperAgentRunStatusUpdate(
|
||||
runId,
|
||||
status.name(),
|
||||
truncate(safeSummary, 512),
|
||||
truncatedSummary,
|
||||
nowUtc()));
|
||||
if (streamSink != null) {
|
||||
streamSink.stage(runId, status, truncatedSummary);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -912,6 +1035,106 @@ public class DebugEmlSuperAgentRunServiceImpl implements DebugEmlSuperAgentRunSe
|
||||
return LocalDateTime.now(ZoneOffset.UTC);
|
||||
}
|
||||
|
||||
/**
|
||||
* Debug EML 流式事件输出端口,避免主流程依赖 HTTP 细节。
|
||||
*/
|
||||
private interface DebugEmlSuperAgentStreamSink {
|
||||
|
||||
/**
|
||||
* 输出本系统处理阶段。
|
||||
*/
|
||||
void stage(Long runId, DebugEmlSuperAgentRunStatus status, String safeSummary);
|
||||
|
||||
/**
|
||||
* 输出 SuperAgent 公开 Trace 事件。
|
||||
*/
|
||||
void trace(SuperAgentOpenApiTraceEvent traceEvent);
|
||||
|
||||
/**
|
||||
* 输出最终 Debug 结果。
|
||||
*/
|
||||
void result(DebugEmlSuperAgentRunResult result);
|
||||
|
||||
/**
|
||||
* 输出安全错误事件。
|
||||
*/
|
||||
void error(Long runId, String errorCode, String message, String status);
|
||||
|
||||
/**
|
||||
* 输出流结束事件。
|
||||
*/
|
||||
void done();
|
||||
}
|
||||
|
||||
/**
|
||||
* 将 Debug EML 事件写成 text/event-stream。
|
||||
*/
|
||||
private static final class OutputStreamDebugEmlStreamSink implements DebugEmlSuperAgentStreamSink {
|
||||
|
||||
private final OutputStream outputStream;
|
||||
private final ObjectMapper objectMapper;
|
||||
private boolean active = true;
|
||||
|
||||
private OutputStreamDebugEmlStreamSink(OutputStream outputStream, ObjectMapper objectMapper) {
|
||||
this.outputStream = outputStream;
|
||||
this.objectMapper = objectMapper;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stage(Long runId, DebugEmlSuperAgentRunStatus status, String safeSummary) {
|
||||
Map<String, Object> data = new LinkedHashMap<>();
|
||||
data.put("debug_run_id", runId == null ? null : runId.toString());
|
||||
data.put("status", status.name());
|
||||
data.put("safe_summary", safeSummary);
|
||||
send("debug_stage", data);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void trace(SuperAgentOpenApiTraceEvent traceEvent) {
|
||||
send("superagent_trace", traceEvent);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void result(DebugEmlSuperAgentRunResult result) {
|
||||
send("superagent_result", result);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void error(Long runId, String errorCode, String message, String status) {
|
||||
Map<String, Object> data = new LinkedHashMap<>();
|
||||
data.put("debug_run_id", runId == null ? null : runId.toString());
|
||||
data.put("error_code", errorCode);
|
||||
data.put("message", message);
|
||||
data.put("status", status);
|
||||
send("debug_error", data);
|
||||
done();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void done() {
|
||||
send("done", Map.of("status", "done"));
|
||||
}
|
||||
|
||||
/**
|
||||
* 写出一个 SSE 事件块并立即 flush。写失败通常表示客户端断开,只关闭流输出,不影响后续业务落库。
|
||||
*/
|
||||
private void send(String eventName, Object data) {
|
||||
if (!active) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
String eventBlock = "event: " + eventName
|
||||
+ "\n"
|
||||
+ "data: " + objectMapper.writeValueAsString(data)
|
||||
+ "\n\n";
|
||||
outputStream.write(eventBlock.getBytes(StandardCharsets.UTF_8));
|
||||
outputStream.flush();
|
||||
} catch (IOException exception) {
|
||||
active = false;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 上传媒体内部组合对象。
|
||||
*/
|
||||
|
||||
+88
@@ -0,0 +1,88 @@
|
||||
package cn.nianxx.thhotel.integrations.ai.superagent.service.impl;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.request.SuperAgentMailDebugRequest;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.sun.net.httpserver.HttpServer;
|
||||
import java.io.OutputStream;
|
||||
import java.net.InetAddress;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.time.Duration;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
class SuperAgentOpenApiClientImplTest {
|
||||
|
||||
@Test
|
||||
void shouldRequestStreamingMessagesWithIncludeTraceAndEmitPublicTraceEvents() throws Exception {
|
||||
HttpServer server = HttpServer.create(new InetSocketAddress(InetAddress.getLoopbackAddress(), 0), 0);
|
||||
AtomicReference<String> streamQuery = new AtomicReference<>();
|
||||
server.createContext("/api/open/agent-sessions", exchange -> {
|
||||
byte[] response = "{\"session_id\":\"session-http-001\"}".getBytes(StandardCharsets.UTF_8);
|
||||
exchange.getResponseHeaders().add("Content-Type", "application/json");
|
||||
exchange.sendResponseHeaders(200, response.length);
|
||||
try (OutputStream outputStream = exchange.getResponseBody()) {
|
||||
outputStream.write(response);
|
||||
}
|
||||
});
|
||||
server.createContext("/api/open/agent-sessions/session-http-001/messages/stream", exchange -> {
|
||||
streamQuery.set(exchange.getRequestURI().getRawQuery());
|
||||
byte[] response = """
|
||||
event: metadata
|
||||
data: {"run_id":"run-http-001","resolved_profile_id":"profile-http","resolved_profile_version_id":"version-http"}
|
||||
|
||||
event: trace
|
||||
data: {"event":"reasoning.summary","run_id":"run-http-001","text":"正在分析 Debug 邮件。","ts":"2026-07-11T10:00:01Z"}
|
||||
|
||||
event: values
|
||||
data: {"messages":[{"type":"ai","content":"{\\"route_code\\":\\"S10\\"}","response_metadata":{"finish_reason":"stop","model_name":"debug-model"},"usage_metadata":{"input_tokens":9,"output_tokens":4,"total_tokens":13}}]}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
""".getBytes(StandardCharsets.UTF_8);
|
||||
exchange.getResponseHeaders().add("Content-Type", "text/event-stream");
|
||||
exchange.sendResponseHeaders(200, response.length);
|
||||
try (OutputStream outputStream = exchange.getResponseBody()) {
|
||||
outputStream.write(response);
|
||||
}
|
||||
});
|
||||
server.start();
|
||||
try {
|
||||
SuperAgentOpenApiProperties properties = new SuperAgentOpenApiProperties();
|
||||
properties.setEnabled(true);
|
||||
properties.setBaseUrl("http://127.0.0.1:" + server.getAddress().getPort());
|
||||
properties.setApiKey("df_open_test");
|
||||
properties.setExternalSubjectId("debug-subject");
|
||||
properties.setConnectTimeout(Duration.ofSeconds(5));
|
||||
properties.setReadTimeout(Duration.ofSeconds(5));
|
||||
SuperAgentOpenApiClientImpl client = new SuperAgentOpenApiClientImpl(
|
||||
properties,
|
||||
new SuperAgentOpenApiSseParser(new ObjectMapper()),
|
||||
new ObjectMapper());
|
||||
List<SuperAgentOpenApiTraceEvent> traceEvents = new ArrayList<>();
|
||||
|
||||
SuperAgentOpenApiResult result = client.invokeMailDebug(new SuperAgentMailDebugRequest(
|
||||
"debug message",
|
||||
"debug-idempotency",
|
||||
Map.of("source", "unit-test")), traceEvents::add);
|
||||
|
||||
assertThat(streamQuery.get()).isEqualTo("include_trace=true");
|
||||
assertThat(result.sessionId()).isEqualTo("session-http-001");
|
||||
assertThat(result.runId()).isEqualTo("run-http-001");
|
||||
assertThat(result.traceEvents()).hasSize(1);
|
||||
assertThat(traceEvents).hasSize(1);
|
||||
assertThat(traceEvents.get(0).event()).isEqualTo("reasoning.summary");
|
||||
assertThat(traceEvents.get(0).text()).isEqualTo("正在分析 Debug 邮件。");
|
||||
} finally {
|
||||
server.stop(0);
|
||||
}
|
||||
}
|
||||
}
|
||||
+40
@@ -4,6 +4,9 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
class SuperAgentOpenApiSseParserTest {
|
||||
@@ -164,4 +167,41 @@ class SuperAgentOpenApiSseParserTest {
|
||||
|
||||
assertThat(result.rawAnswer()).isEqualTo("{\"ai_task_results\":[{\"task_type\":\"Cancel Booking\"}]}");
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldCollectPublicTraceEventsAndKeepFinalAnswer() {
|
||||
String sse = """
|
||||
event: metadata
|
||||
data: {"run_id":"run-debug-trace","resolved_profile_id":"profile-debug"}
|
||||
|
||||
event: trace
|
||||
data: {"event":"reasoning.summary","run_id":"run-debug-trace","text":"正在分析邮件和附件。","ts":"2026-07-11T10:00:01Z"}
|
||||
|
||||
event: trace
|
||||
data: {"event":"tool.call.started","run_id":"run-debug-trace","message_id":"ai-1","tool_calls":[{"tool_call_id":"tool-1","name":"th_hotel_query_case_context","input":{"summary":"{\\"group_code\\":\\"G001\\"}"}}],"ts":"2026-07-11T10:00:02Z"}
|
||||
|
||||
event: trace
|
||||
data: {"event":"tool.call.completed","run_id":"run-debug-trace","tool_call_id":"tool-1","name":"th_hotel_query_case_context","output":{"summary":"matched 1 case"},"ts":"2026-07-11T10:00:03Z"}
|
||||
|
||||
event: values
|
||||
data: {"messages":[{"type":"ai","content":"{\\"route_code\\":\\"S10\\"}","response_metadata":{"finish_reason":"stop","model_name":"debug-model"},"usage_metadata":{"input_tokens":31,"output_tokens":5,"total_tokens":36}}]}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
""";
|
||||
List<SuperAgentOpenApiTraceEvent> emittedTraceEvents = new ArrayList<>();
|
||||
|
||||
SuperAgentOpenApiResult result = parser.parse("session-debug-trace", sse, emittedTraceEvents::add);
|
||||
|
||||
assertThat(result.rawAnswer()).isEqualTo("{\"route_code\":\"S10\"}");
|
||||
assertThat(result.eventTypes()).containsExactly("metadata", "trace", "values", "end");
|
||||
assertThat(result.traceEvents()).hasSize(3);
|
||||
assertThat(emittedTraceEvents).hasSize(3);
|
||||
assertThat(result.traceEvents().get(0).event()).isEqualTo("reasoning.summary");
|
||||
assertThat(result.traceEvents().get(0).text()).isEqualTo("正在分析邮件和附件。");
|
||||
assertThat(result.traceEvents().get(1).toolName()).isEqualTo("th_hotel_query_case_context");
|
||||
assertThat(result.traceEvents().get(1).inputSummary()).isEqualTo("{\"group_code\":\"G001\"}");
|
||||
assertThat(result.traceEvents().get(2).outputSummary()).isEqualTo("matched 1 case");
|
||||
}
|
||||
}
|
||||
+168
@@ -7,20 +7,27 @@ import static org.hamcrest.Matchers.not;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.reset;
|
||||
import static org.mockito.Mockito.when;
|
||||
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.asyncDispatch;
|
||||
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.multipart;
|
||||
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.content;
|
||||
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath;
|
||||
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.request;
|
||||
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
|
||||
|
||||
import cn.nianxx.thhotel.ThHotelApplication;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiTraceEvent;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.SuperAgentOpenApiClient;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.impl.SuperAgentOpenApiException;
|
||||
import cn.nianxx.thhotel.integrations.storage.aliyunoss.common.request.ObjectStoragePutRequest;
|
||||
import cn.nianxx.thhotel.integrations.storage.aliyunoss.common.result.ObjectStoragePutResult;
|
||||
import cn.nianxx.thhotel.integrations.storage.aliyunoss.service.ObjectStorageService;
|
||||
import cn.nianxx.thhotel.platform.debug.service.DebugEmlSuperAgentRunService;
|
||||
import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.List;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
@@ -32,6 +39,7 @@ import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.mock.web.MockMultipartFile;
|
||||
import org.springframework.test.context.ActiveProfiles;
|
||||
import org.springframework.test.web.servlet.MockMvc;
|
||||
import org.springframework.test.web.servlet.MvcResult;
|
||||
import org.springframework.web.client.RestClientResponseException;
|
||||
|
||||
@SpringBootTest(
|
||||
@@ -56,6 +64,9 @@ class DebugEmlSuperAgentControllerTest {
|
||||
@Autowired
|
||||
private JdbcTemplate jdbcTemplate;
|
||||
|
||||
@Autowired
|
||||
private DebugEmlSuperAgentRunService runService;
|
||||
|
||||
@MockBean
|
||||
private ObjectStorageService objectStorageService;
|
||||
|
||||
@@ -71,6 +82,17 @@ class DebugEmlSuperAgentControllerTest {
|
||||
.andExpect(content().string(not(containsString("test-debug-upload-key"))));
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldRejectStreamUploadWhenDebugKeyMissingBeforeStartingSse() throws Exception {
|
||||
mockMvc.perform(multipart(ENDPOINT + "/stream")
|
||||
.file(emlFile())
|
||||
.param("hotel_id", "HOTEL-TEST"))
|
||||
.andExpect(status().isUnauthorized())
|
||||
.andExpect(content().contentTypeCompatibleWith(MediaType.APPLICATION_JSON))
|
||||
.andExpect(jsonPath("$.error_code").value("DEBUG_UPLOAD_KEY_INVALID"))
|
||||
.andExpect(content().string(not(containsString("test-debug-upload-key"))));
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldUploadEmlToOssCaptureSourceMessageAndReturnSuperAgentResult() throws Exception {
|
||||
when(objectStorageService.putObject(any())).thenAnswer(invocation -> {
|
||||
@@ -158,6 +180,139 @@ class DebugEmlSuperAgentControllerTest {
|
||||
org.assertj.core.api.Assertions.assertThat(debugRunCount).isEqualTo(1L);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldStreamDebugStagesSuperAgentTraceAndFinalResult() throws Exception {
|
||||
when(objectStorageService.putObject(any())).thenAnswer(invocation -> {
|
||||
ObjectStoragePutRequest request = invocation.getArgument(0);
|
||||
return new ObjectStoragePutResult(
|
||||
request.objectKey(),
|
||||
"https://oss.example.test/" + request.objectKey(),
|
||||
request.contentType(),
|
||||
request.sizeBytes());
|
||||
});
|
||||
when(superAgentOpenApiClient.invokeMailDebug(any(), any())).thenAnswer(invocation -> {
|
||||
@SuppressWarnings("unchecked")
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer = invocation.getArgument(1);
|
||||
traceConsumer.accept(new SuperAgentOpenApiTraceEvent(
|
||||
"reasoning.summary",
|
||||
"run-stream-001",
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
"正在分析 Debug 邮件。",
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
"2026-07-11T10:00:01Z"));
|
||||
traceConsumer.accept(new SuperAgentOpenApiTraceEvent(
|
||||
"tool.call.started",
|
||||
"run-stream-001",
|
||||
"ai-1",
|
||||
"tool-1",
|
||||
"th_hotel_query_case_context",
|
||||
null,
|
||||
"{\"group_code\":\"G001\"}",
|
||||
null,
|
||||
null,
|
||||
"2026-07-11T10:00:02Z"));
|
||||
return new SuperAgentOpenApiResult(
|
||||
"session-stream-001",
|
||||
"run-stream-001",
|
||||
"profile-debug",
|
||||
"profile-version-debug",
|
||||
"debug-model",
|
||||
"{\"route_code\":\"S10\"}",
|
||||
11,
|
||||
7,
|
||||
18,
|
||||
List.of("metadata", "trace", "values", "end"),
|
||||
List.of());
|
||||
});
|
||||
|
||||
MvcResult mvcResult = mockMvc.perform(multipart(ENDPOINT + "/stream")
|
||||
.file(emlFile())
|
||||
.param("hotel_id", "HOTEL-TEST")
|
||||
.param("run_label", "stream-debug-upload")
|
||||
.header("X-TH-Hotel-Debug-Upload-Key", "test-debug-upload-key"))
|
||||
.andExpect(request().asyncStarted())
|
||||
.andReturn();
|
||||
|
||||
mockMvc.perform(asyncDispatch(mvcResult))
|
||||
.andExpect(status().isOk())
|
||||
.andExpect(content().contentTypeCompatibleWith(MediaType.TEXT_EVENT_STREAM))
|
||||
.andExpect(content().string(containsString("event: debug_stage")))
|
||||
.andExpect(content().string(containsString("\"status\":\"PARSING_EML\"")))
|
||||
.andExpect(content().string(containsString("\"status\":\"CALLING_SUPERAGENT\"")))
|
||||
.andExpect(content().string(containsString("event: superagent_trace")))
|
||||
.andExpect(content().string(containsString("\"event\":\"reasoning.summary\"")))
|
||||
.andExpect(content().string(containsString("\"event\":\"tool.call.started\"")))
|
||||
.andExpect(content().string(containsString("\"tool_name\":\"th_hotel_query_case_context\"")))
|
||||
.andExpect(content().string(containsString("event: superagent_result")))
|
||||
.andExpect(content().string(containsString("\"superagent_run_id\":\"run-stream-001\"")))
|
||||
.andExpect(content().string(containsString("\"superagent_parsed_json\":{\"route_code\":\"S10\"}")))
|
||||
.andExpect(content().string(containsString("event: done")))
|
||||
.andExpect(content().string(not(containsString("test-debug-upload-key"))));
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldKeepBusinessRunStatusWhenSseClientDisconnects() {
|
||||
when(objectStorageService.putObject(any())).thenAnswer(invocation -> {
|
||||
ObjectStoragePutRequest request = invocation.getArgument(0);
|
||||
return new ObjectStoragePutResult(
|
||||
request.objectKey(),
|
||||
"https://oss.example.test/" + request.objectKey(),
|
||||
request.contentType(),
|
||||
request.sizeBytes());
|
||||
});
|
||||
when(superAgentOpenApiClient.invokeMailDebug(any(), any())).thenAnswer(invocation -> {
|
||||
@SuppressWarnings("unchecked")
|
||||
Consumer<SuperAgentOpenApiTraceEvent> traceConsumer = invocation.getArgument(1);
|
||||
traceConsumer.accept(new SuperAgentOpenApiTraceEvent(
|
||||
"reasoning.summary",
|
||||
"run-stream-disconnect",
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
"正在分析 Debug 邮件。",
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
"2026-07-11T10:00:01Z"));
|
||||
return new SuperAgentOpenApiResult(
|
||||
"session-stream-disconnect",
|
||||
"run-stream-disconnect",
|
||||
"profile-debug",
|
||||
"profile-version-debug",
|
||||
"debug-model",
|
||||
"{\"route_code\":\"S10\"}",
|
||||
11,
|
||||
7,
|
||||
18,
|
||||
List.of("metadata", "trace", "values", "end"),
|
||||
List.of());
|
||||
});
|
||||
|
||||
org.assertj.core.api.Assertions.assertThatCode(() -> runService.uploadAndRunStream(
|
||||
"test-debug-upload-key",
|
||||
emlFile(),
|
||||
"HOTEL-TEST",
|
||||
"stream-client-disconnected",
|
||||
new FailingOutputStream()))
|
||||
.doesNotThrowAnyException();
|
||||
|
||||
Long debugRunCount = jdbcTemplate.queryForObject("""
|
||||
SELECT COUNT(*)
|
||||
FROM platform_debug_eml_superagent_run
|
||||
WHERE hotel_id = 'HOTEL-TEST'
|
||||
AND run_status = 'SUPERAGENT_SUCCEEDED'
|
||||
AND superagent_session_id = 'session-stream-disconnect'
|
||||
AND superagent_run_id = 'run-stream-disconnect'
|
||||
AND run_label = 'stream-client-disconnected'
|
||||
AND safe_error_summary IS NULL
|
||||
""", Long.class);
|
||||
org.assertj.core.api.Assertions.assertThat(debugRunCount).isEqualTo(1L);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldRecordPhaseBeforeUploadingOriginalEmlToOss() throws Exception {
|
||||
AtomicInteger uploadIndex = new AtomicInteger();
|
||||
@@ -408,6 +563,19 @@ class DebugEmlSuperAgentControllerTest {
|
||||
unsafeHtmlEmlBytes());
|
||||
}
|
||||
|
||||
private static final class FailingOutputStream extends OutputStream {
|
||||
|
||||
@Override
|
||||
public void write(int b) throws IOException {
|
||||
throw new IOException("client disconnected");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void write(byte[] b, int off, int len) throws IOException {
|
||||
throw new IOException("client disconnected");
|
||||
}
|
||||
}
|
||||
|
||||
private byte[] emlBytes() {
|
||||
return """
|
||||
From: Guest <guest@example.test>
|
||||
|
||||
Reference in new issue
Block a user