实现AgentBus自动分发SuperAgent链路
This commit is contained in:
@@ -0,0 +1,193 @@
|
||||
package cn.nianxx.thhotel.integrations.ai.superagent.service.impl;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.never;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.dto.SuperAgentDispatchRunDraft;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.dto.SuperAgentDispatchRunSnapshot;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.dto.SuperAgentDispatchRunSuccessUpdate;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.enums.SuperAgentDispatchSource;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.enums.SuperAgentDispatchStatus;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.request.SuperAgentOpenApiMessageRequest;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.common.result.SuperAgentOpenApiResult;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.repository.SuperAgentDispatchRunRepository;
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.SuperAgentOpenApiClient;
|
||||
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageCommand;
|
||||
import cn.nianxx.thhotel.platform.message.common.result.SourceMessageCaptureResult;
|
||||
import cn.nianxx.thhotel.platform.message.repository.SourceMessageInboxRepository;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
|
||||
class SuperAgentDispatchServiceImplTest {
|
||||
|
||||
@Test
|
||||
void shouldCreatePendingDispatchForReceivedAgentBusSourceMessage() {
|
||||
SuperAgentDispatchRunRepository runRepository = mock(SuperAgentDispatchRunRepository.class);
|
||||
when(runRepository.insertIfAbsent(any(SuperAgentDispatchRunDraft.class))).thenReturn(99001L);
|
||||
SuperAgentDispatchServiceImpl service = service(runRepository, mock(SourceMessageInboxRepository.class),
|
||||
mock(SuperAgentOpenApiClient.class), properties(true, false));
|
||||
|
||||
service.enqueueAgentBusRealtime(agentBusCommand("AGENTBUS"), new SourceMessageCaptureResult(
|
||||
88001L,
|
||||
true,
|
||||
false,
|
||||
"RECEIVED"));
|
||||
|
||||
ArgumentCaptor<SuperAgentDispatchRunDraft> captor = ArgumentCaptor.forClass(SuperAgentDispatchRunDraft.class);
|
||||
verify(runRepository).insertIfAbsent(captor.capture());
|
||||
assertThat(captor.getValue().sourceMessageId()).isEqualTo(88001L);
|
||||
assertThat(captor.getValue().dispatchSource()).isEqualTo(SuperAgentDispatchSource.AGENTBUS_REALTIME.code());
|
||||
assertThat(captor.getValue().dispatchStatus()).isEqualTo(SuperAgentDispatchStatus.PENDING.code());
|
||||
assertThat(captor.getValue().idempotencyKey()).isEqualTo("agentbus-source-message-88001");
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldIgnoreNonAgentBusOrFailedCaptureWhenEnqueueDispatch() {
|
||||
SuperAgentDispatchRunRepository runRepository = mock(SuperAgentDispatchRunRepository.class);
|
||||
SuperAgentDispatchServiceImpl service = service(runRepository, mock(SourceMessageInboxRepository.class),
|
||||
mock(SuperAgentOpenApiClient.class), properties(true, false));
|
||||
|
||||
service.enqueueAgentBusRealtime(agentBusCommand("DEBUG_EML_UPLOAD"), new SourceMessageCaptureResult(
|
||||
88002L,
|
||||
true,
|
||||
false,
|
||||
"RECEIVED"));
|
||||
service.enqueueAgentBusRealtime(agentBusCommand("AGENTBUS"), new SourceMessageCaptureResult(
|
||||
88003L,
|
||||
true,
|
||||
false,
|
||||
"FAILED"));
|
||||
service.enqueueAgentBusRealtime(agentBusCommand("AGENTBUS"), new SourceMessageCaptureResult(
|
||||
88004L,
|
||||
false,
|
||||
false,
|
||||
"RECEIVED"));
|
||||
|
||||
verify(runRepository, never()).insertIfAbsent(any(SuperAgentDispatchRunDraft.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldProcessDueDispatchAndPersistSuperAgentSuccess() {
|
||||
SuperAgentDispatchRunRepository runRepository = mock(SuperAgentDispatchRunRepository.class);
|
||||
SourceMessageInboxRepository sourceMessageRepository = mock(SourceMessageInboxRepository.class);
|
||||
SuperAgentOpenApiClient openApiClient = mock(SuperAgentOpenApiClient.class);
|
||||
SuperAgentDispatchServiceImpl service = service(runRepository, sourceMessageRepository, openApiClient,
|
||||
properties(true, true));
|
||||
SuperAgentDispatchRunSnapshot snapshot = snapshot();
|
||||
when(runRepository.claimDue(any(), any(), any(), any(), any(Integer.class))).thenReturn(List.of(snapshot));
|
||||
when(sourceMessageRepository.findPayloadJson(88001L)).thenReturn(Optional.of("""
|
||||
{"source":{"external_message_id":"mail-agentbus-001"},"body":{"text":"Please book a room"}}
|
||||
"""));
|
||||
when(openApiClient.invokeMessage(any(SuperAgentOpenApiMessageRequest.class))).thenReturn(new SuperAgentOpenApiResult(
|
||||
"session-dispatch-001",
|
||||
"run-dispatch-001",
|
||||
"profile-dispatch",
|
||||
"version-dispatch",
|
||||
"model-dispatch",
|
||||
"{\"route_code\":\"S10\"}",
|
||||
10,
|
||||
5,
|
||||
15,
|
||||
List.of("trace", "end"),
|
||||
List.of(),
|
||||
"/api/open/agent-sessions/session-dispatch-001/runs/run-dispatch-001",
|
||||
"4"));
|
||||
|
||||
int processed = service.processDueDispatches();
|
||||
|
||||
assertThat(processed).isEqualTo(1);
|
||||
ArgumentCaptor<SuperAgentOpenApiMessageRequest> requestCaptor =
|
||||
ArgumentCaptor.forClass(SuperAgentOpenApiMessageRequest.class);
|
||||
verify(openApiClient).invokeMessage(requestCaptor.capture());
|
||||
assertThat(requestCaptor.getValue().idempotencyKey()).isEqualTo("agentbus-source-message-88001");
|
||||
assertThat(requestCaptor.getValue().metadata()).containsEntry("external_message_id", "mail-agentbus-001");
|
||||
ArgumentCaptor<SuperAgentDispatchRunSuccessUpdate> updateCaptor =
|
||||
ArgumentCaptor.forClass(SuperAgentDispatchRunSuccessUpdate.class);
|
||||
verify(runRepository).markSucceeded(updateCaptor.capture());
|
||||
assertThat(updateCaptor.getValue().id()).isEqualTo(99001L);
|
||||
assertThat(updateCaptor.getValue().superagentRunId()).isEqualTo("run-dispatch-001");
|
||||
assertThat(updateCaptor.getValue().lastEventId()).isEqualTo("4");
|
||||
assertThat(updateCaptor.getValue().superagentParsedJson()).isEqualTo("{\"route_code\":\"S10\"}");
|
||||
}
|
||||
|
||||
private SuperAgentDispatchServiceImpl service(
|
||||
SuperAgentDispatchRunRepository runRepository,
|
||||
SourceMessageInboxRepository sourceMessageRepository,
|
||||
SuperAgentOpenApiClient openApiClient,
|
||||
SuperAgentDispatchProperties properties) {
|
||||
return new SuperAgentDispatchServiceImpl(
|
||||
properties,
|
||||
new SuperAgentOpenApiProperties(),
|
||||
runRepository,
|
||||
sourceMessageRepository,
|
||||
openApiClient,
|
||||
new ObjectMapper());
|
||||
}
|
||||
|
||||
private SuperAgentDispatchProperties properties(boolean enabled, boolean workerEnabled) {
|
||||
SuperAgentDispatchProperties properties = new SuperAgentDispatchProperties();
|
||||
properties.setEnabled(enabled);
|
||||
properties.setWorkerEnabled(workerEnabled);
|
||||
return properties;
|
||||
}
|
||||
|
||||
private CaptureSourceMessageCommand agentBusCommand(String provider) {
|
||||
return new CaptureSourceMessageCommand(
|
||||
"HOTEL-TEST",
|
||||
provider,
|
||||
"OUTLOOK",
|
||||
"mail-agentbus-001",
|
||||
"conversation-agentbus-001",
|
||||
"frame-agentbus-001",
|
||||
"session-agentbus-001",
|
||||
null,
|
||||
null,
|
||||
"guest@example.test",
|
||||
"Booking",
|
||||
"Please book a room",
|
||||
"<html>Please book a room</html>",
|
||||
"{\"source\":{\"external_message_id\":\"mail-agentbus-001\"}}",
|
||||
"agentbus-outlook-v1",
|
||||
List.of());
|
||||
}
|
||||
|
||||
private SuperAgentDispatchRunSnapshot snapshot() {
|
||||
LocalDateTime now = LocalDateTime.of(2026, 7, 12, 12, 0);
|
||||
return new SuperAgentDispatchRunSnapshot(
|
||||
99001L,
|
||||
"HOTEL-TEST",
|
||||
88001L,
|
||||
"AGENTBUS",
|
||||
"OUTLOOK",
|
||||
"mail-agentbus-001",
|
||||
"conversation-agentbus-001",
|
||||
SuperAgentDispatchSource.AGENTBUS_REALTIME.code(),
|
||||
"agentbus-source-message-88001",
|
||||
"agentbus-dispatch-source-88001",
|
||||
SuperAgentDispatchStatus.RUNNING.code(),
|
||||
1,
|
||||
3,
|
||||
now,
|
||||
now.plusMinutes(5),
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
now,
|
||||
now);
|
||||
}
|
||||
}
|
||||
@@ -15,6 +15,7 @@ import java.time.Duration;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
@@ -35,20 +36,25 @@ class SuperAgentOpenApiClientImplTest {
|
||||
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"}
|
||||
|
||||
id: 1
|
||||
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}}]}
|
||||
id: 2
|
||||
event: trace
|
||||
data: {"event":"message.final","run_id":"run-http-001","text":"{\\"route_code\\":\\"S10\\"}"}
|
||||
|
||||
id: 3
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-http-001","status":"success"}
|
||||
|
||||
id: 4
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
""".getBytes(StandardCharsets.UTF_8);
|
||||
exchange.getResponseHeaders().add("Content-Type", "text/event-stream");
|
||||
exchange.getResponseHeaders().add("Content-Location", "/api/open/agent-sessions/session-http-001/runs/run-http-001");
|
||||
exchange.sendResponseHeaders(200, response.length);
|
||||
try (OutputStream outputStream = exchange.getResponseBody()) {
|
||||
outputStream.write(response);
|
||||
@@ -77,12 +83,101 @@ class SuperAgentOpenApiClientImplTest {
|
||||
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(result.lastEventId()).isEqualTo("4");
|
||||
assertThat(result.traceEvents()).hasSize(3);
|
||||
assertThat(traceEvents).hasSize(3);
|
||||
assertThat(traceEvents.get(0).event()).isEqualTo("reasoning.summary");
|
||||
assertThat(traceEvents.get(0).text()).isEqualTo("正在分析 Debug 邮件。");
|
||||
} finally {
|
||||
server.stop(0);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldRecoverPrematureEofBySubscribingRunEventsWithoutRepostingInitialMessage() throws Exception {
|
||||
HttpServer server = HttpServer.create(new InetSocketAddress(InetAddress.getLoopbackAddress(), 0), 0);
|
||||
AtomicInteger initialMessagePostCount = new AtomicInteger();
|
||||
AtomicReference<String> lastEventId = new AtomicReference<>();
|
||||
server.createContext("/api/open/agent-sessions", exchange -> {
|
||||
byte[] response = "{\"session_id\":\"session-recover-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-recover-001/messages/stream", exchange -> {
|
||||
initialMessagePostCount.incrementAndGet();
|
||||
byte[] response = """
|
||||
id: 1
|
||||
event: trace
|
||||
data: {"event":"message.delta","run_id":"run-recover-001","text":"{\\"route"}
|
||||
|
||||
""".getBytes(StandardCharsets.UTF_8);
|
||||
exchange.getResponseHeaders().add("Content-Type", "text/event-stream");
|
||||
exchange.getResponseHeaders().add("Content-Location", "/api/open/agent-sessions/session-recover-001/runs/run-recover-001");
|
||||
exchange.sendResponseHeaders(200, response.length);
|
||||
try (OutputStream outputStream = exchange.getResponseBody()) {
|
||||
outputStream.write(response);
|
||||
}
|
||||
});
|
||||
server.createContext("/api/open/agent-sessions/session-recover-001/runs/run-recover-001", exchange -> {
|
||||
byte[] response = "{\"status\":\"running\"}".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-recover-001/runs/run-recover-001/events", exchange -> {
|
||||
lastEventId.set(exchange.getRequestHeaders().getFirst("Last-Event-ID"));
|
||||
byte[] response = """
|
||||
id: 2
|
||||
event: trace
|
||||
data: {"event":"message.final","run_id":"run-recover-001","text":"{\\"route_code\\":\\"S10\\"}"}
|
||||
|
||||
id: 3
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-recover-001","status":"success"}
|
||||
|
||||
id: 4
|
||||
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));
|
||||
properties.setSseRecoveryMaxAttempts(2);
|
||||
SuperAgentOpenApiClientImpl client = new SuperAgentOpenApiClientImpl(
|
||||
properties,
|
||||
new SuperAgentOpenApiSseParser(new ObjectMapper()),
|
||||
new ObjectMapper());
|
||||
|
||||
SuperAgentOpenApiResult result = client.invokeMailDebug(new SuperAgentMailDebugRequest(
|
||||
"debug message",
|
||||
"debug-recover-idempotency",
|
||||
Map.of("source", "unit-test")));
|
||||
|
||||
assertThat(initialMessagePostCount).hasValue(1);
|
||||
assertThat(lastEventId).hasValue("1");
|
||||
assertThat(result.rawAnswer()).isEqualTo("{\"route_code\":\"S10\"}");
|
||||
assertThat(result.runId()).isEqualTo("run-recover-001");
|
||||
assertThat(result.lastEventId()).isEqualTo("4");
|
||||
} finally {
|
||||
server.stop(0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -25,6 +25,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
event: values
|
||||
data: {"messages":[{"type":"human","content":"input"},{"type":"ai","content":"{\\"ai_task_results\\":[{\\"task_type\\":\\"New Booking\\"}]}","response_metadata":{"finish_reason":"stop","model_name":"debug-model"},"usage_metadata":{"input_tokens":11,"output_tokens":7,"total_tokens":18}}]}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-001","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -40,7 +43,7 @@ class SuperAgentOpenApiSseParserTest {
|
||||
assertThat(result.inputTokens()).isEqualTo(11);
|
||||
assertThat(result.outputTokens()).isEqualTo(7);
|
||||
assertThat(result.totalTokens()).isEqualTo(18);
|
||||
assertThat(result.eventTypes()).containsExactly("metadata", "messages", "values", "end");
|
||||
assertThat(result.eventTypes()).containsExactly("metadata", "messages", "values", "trace", "end");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -52,6 +55,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
event: messages
|
||||
data: {"type":"ai","content":"{\\"ai_task_results\\":[{\\"task_type\\":\\"Update Booking\\"}]}","response_metadata":{"finish_reason":"stop","model_name":"debug-model-from-message"},"usage_metadata":{"input_tokens":21,"output_tokens":13,"total_tokens":34}}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-002","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -65,7 +71,7 @@ class SuperAgentOpenApiSseParserTest {
|
||||
assertThat(result.inputTokens()).isEqualTo(21);
|
||||
assertThat(result.outputTokens()).isEqualTo(13);
|
||||
assertThat(result.totalTokens()).isEqualTo(34);
|
||||
assertThat(result.eventTypes()).containsExactly("metadata", "messages", "end");
|
||||
assertThat(result.eventTypes()).containsExactly("metadata", "messages", "trace", "end");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -77,6 +83,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
event: messages
|
||||
data: {"type":"ai","content":"{\\"ai_task_results\\":[{\\"task_type\\":\\"Rooming List\\"}]}","response_metadata":{"model_name":"debug-model-fallback"}}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-003","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -98,21 +107,23 @@ class SuperAgentOpenApiSseParserTest {
|
||||
|
||||
assertThatThrownBy(() -> parser.parse("session-debug-partial", sse))
|
||||
.isInstanceOf(SuperAgentOpenApiException.class)
|
||||
.hasMessageContaining("SuperAgent SSE 未收到结束事件。");
|
||||
.hasMessageContaining("SuperAgent SSE 提前结束,尚未收到 end 事件。");
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldAcceptMissingEndEventWhenFinalAiContentExists() {
|
||||
void shouldFailWhenEndEventMissingEvenIfFinalAiContentExists() {
|
||||
String sse = """
|
||||
event: messages
|
||||
data: {"type":"ai","content":"{\\"ai_task_results\\":[]}","response_metadata":{"finish_reason":"stop"}}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-004","status":"success"}
|
||||
|
||||
""";
|
||||
|
||||
SuperAgentOpenApiResult result = parser.parse("session-debug-004", sse);
|
||||
|
||||
assertThat(result.rawAnswer()).isEqualTo("{\"ai_task_results\":[]}");
|
||||
assertThat(result.eventTypes()).containsExactly("messages");
|
||||
assertThatThrownBy(() -> parser.parse("session-debug-004", sse))
|
||||
.isInstanceOf(SuperAgentOpenApiException.class)
|
||||
.hasMessageContaining("SuperAgent SSE 提前结束,尚未收到 end 事件。");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -124,6 +135,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
event: values
|
||||
data: {"messages":[{"type":"human","content":"input"}]}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-005","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -137,6 +151,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
@Test
|
||||
void shouldNotUseMessagesAfterEndAsFinalAnswer() {
|
||||
String sse = """
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-006","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -156,6 +173,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
event: messages
|
||||
data: {"type":"ai","content":"{\\"ai_task_results\\":[{\\"task_type\\":\\"Voucher Payment\\"}]}","response_metadata":{"finish_reason":"stop"}}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-007","status":"success"}
|
||||
|
||||
event: end
|
||||
|
||||
""";
|
||||
@@ -163,7 +183,7 @@ class SuperAgentOpenApiSseParserTest {
|
||||
SuperAgentOpenApiResult result = parser.parse("session-debug-007", sse);
|
||||
|
||||
assertThat(result.rawAnswer()).isEqualTo("{\"ai_task_results\":[{\"task_type\":\"Voucher Payment\"}]}");
|
||||
assertThat(result.eventTypes()).containsExactly("messages", "end");
|
||||
assertThat(result.eventTypes()).containsExactly("messages", "trace", "end");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -172,6 +192,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
event: messages
|
||||
data: {"type":"ai","content":{"ai_task_results":[{"task_type":"Cancel Booking"}]},"response_metadata":{"finish_reason":"stop"}}
|
||||
|
||||
event: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-008","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -200,6 +223,9 @@ class SuperAgentOpenApiSseParserTest {
|
||||
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: trace
|
||||
data: {"event":"run.completed","run_id":"run-debug-trace","status":"success"}
|
||||
|
||||
event: end
|
||||
data: {}
|
||||
|
||||
@@ -210,8 +236,8 @@ class SuperAgentOpenApiSseParserTest {
|
||||
|
||||
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()).hasSize(4);
|
||||
assertThat(emittedTraceEvents).hasSize(4);
|
||||
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");
|
||||
|
||||
@@ -7,6 +7,7 @@ import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.verifyNoInteractions;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import cn.nianxx.thhotel.integrations.ai.superagent.service.SuperAgentDispatchService;
|
||||
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageCommand;
|
||||
import cn.nianxx.thhotel.platform.message.common.result.SourceMessageCaptureResult;
|
||||
import cn.nianxx.thhotel.platform.message.service.SourceMessageCaptureService;
|
||||
@@ -38,10 +39,12 @@ class AgentBusFrameProcessorTest {
|
||||
@Test
|
||||
void shouldCaptureBusinessFrameIntoSourceMessageInbox() {
|
||||
SourceMessageCaptureService captureService = mock(SourceMessageCaptureService.class);
|
||||
SourceMessageCaptureResult captureResult = new SourceMessageCaptureResult(88001L, true, false, "RECEIVED");
|
||||
when(captureService.capture(any(CaptureSourceMessageCommand.class)))
|
||||
.thenReturn(new SourceMessageCaptureResult(88001L, true, false, "RECEIVED"));
|
||||
.thenReturn(captureResult);
|
||||
SuperAgentDispatchService dispatchService = mock(SuperAgentDispatchService.class);
|
||||
AgentBusConnectionStatus status = new AgentBusConnectionStatus();
|
||||
AgentBusFrameProcessor processor = processor(captureService, status, properties(true, 4096));
|
||||
AgentBusFrameProcessor processor = processor(captureService, dispatchService, status, properties(true, 4096));
|
||||
|
||||
AgentBusFrameProcessResult result = processor.process("""
|
||||
{
|
||||
@@ -68,6 +71,31 @@ class AgentBusFrameProcessorTest {
|
||||
assertThat(captor.getValue().hotelId()).isEqualTo("HOTEL-SYSTEM");
|
||||
assertThat(captor.getValue().externalMessageId()).isEqualTo("mail-agentbus-capture-001");
|
||||
assertThat(captor.getValue().providerFrameId()).isEqualTo("frame-agentbus-capture-001");
|
||||
verify(dispatchService).enqueueAgentBusRealtime(captor.getValue(), captureResult);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldNotEnqueueSuperAgentDispatchWhenCaptureFailed() {
|
||||
SourceMessageCaptureService captureService = mock(SourceMessageCaptureService.class);
|
||||
when(captureService.capture(any(CaptureSourceMessageCommand.class)))
|
||||
.thenReturn(new SourceMessageCaptureResult(88002L, true, false, "FAILED"));
|
||||
SuperAgentDispatchService dispatchService = mock(SuperAgentDispatchService.class);
|
||||
AgentBusConnectionStatus status = new AgentBusConnectionStatus();
|
||||
AgentBusFrameProcessor processor = processor(captureService, dispatchService, status, properties(true, 4096));
|
||||
|
||||
AgentBusFrameProcessResult result = processor.process("""
|
||||
{
|
||||
"id": "frame-agentbus-failed-001",
|
||||
"payload": {
|
||||
"body": {"text": "Missing source message id"},
|
||||
"source": {"channel": "email"}
|
||||
}
|
||||
}
|
||||
""");
|
||||
|
||||
assertThat(result.outcome()).isEqualTo("CAPTURED");
|
||||
assertThat(result.captureStatus()).isEqualTo("FAILED");
|
||||
verifyNoInteractions(dispatchService);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -108,12 +136,21 @@ class AgentBusFrameProcessorTest {
|
||||
SourceMessageCaptureService captureService,
|
||||
AgentBusConnectionStatus status,
|
||||
AgentBusProperties properties) {
|
||||
return processor(captureService, mock(SuperAgentDispatchService.class), status, properties);
|
||||
}
|
||||
|
||||
private AgentBusFrameProcessor processor(
|
||||
SourceMessageCaptureService captureService,
|
||||
SuperAgentDispatchService dispatchService,
|
||||
AgentBusConnectionStatus status,
|
||||
AgentBusProperties properties) {
|
||||
HotelContextService hotelContextService = mock(HotelContextService.class);
|
||||
when(hotelContextService.resolveSystemHotelId()).thenReturn("HOTEL-SYSTEM");
|
||||
return new AgentBusFrameProcessor(
|
||||
objectMapper,
|
||||
new AgentBusSourceMessageAdapter(objectMapper),
|
||||
captureService,
|
||||
dispatchService,
|
||||
status,
|
||||
properties,
|
||||
hotelContextService);
|
||||
|
||||
Reference in New Issue
Block a user