实现 SourceMessage Inbox 与 AgentBus 入站闭环
This commit is contained in:
1 parent
5a498635d8
commit
7e276469ef
75 files changed
+4608
-27
No files matched your search
+126
@@ -0,0 +1,126 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
import java.time.OffsetDateTime;
|
||||
import java.time.ZoneOffset;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* AgentBus 连接与入站处理状态。该组件只保存安全计数器和错误代码,不保存 raw frame。
|
||||
*/
|
||||
@Component
|
||||
public class AgentBusConnectionStatus {
|
||||
|
||||
private final AtomicLong receivedFrameCount = new AtomicLong();
|
||||
private final AtomicLong capturedFrameCount = new AtomicLong();
|
||||
private final AtomicLong ignoredFrameCount = new AtomicLong();
|
||||
private final AtomicLong rejectedFrameCount = new AtomicLong();
|
||||
private final AtomicLong failedFrameCount = new AtomicLong();
|
||||
|
||||
private volatile boolean connected;
|
||||
private volatile boolean sessionReady;
|
||||
private volatile String lastErrorCode;
|
||||
private volatile String lastErrorSummary;
|
||||
private volatile OffsetDateTime connectedAt;
|
||||
private volatile OffsetDateTime disconnectedAt;
|
||||
private volatile OffsetDateTime lastFrameReceivedAt;
|
||||
|
||||
/**
|
||||
* 标记 WebSocket 已连接。
|
||||
*/
|
||||
public void markConnected() {
|
||||
connected = true;
|
||||
connectedAt = nowUtc();
|
||||
lastErrorCode = null;
|
||||
lastErrorSummary = null;
|
||||
}
|
||||
|
||||
/**
|
||||
* 标记 WebSocket 已断开。
|
||||
*/
|
||||
public void markDisconnected() {
|
||||
connected = false;
|
||||
sessionReady = false;
|
||||
disconnectedAt = nowUtc();
|
||||
}
|
||||
|
||||
/**
|
||||
* 标记已经收到 AgentBus session.ready 控制事件。
|
||||
*/
|
||||
public void markSessionReady() {
|
||||
sessionReady = true;
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录收到一条入站 frame。
|
||||
*/
|
||||
public void markFrameReceived() {
|
||||
receivedFrameCount.incrementAndGet();
|
||||
lastFrameReceivedAt = nowUtc();
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录一条 frame 已写入 SourceMessage Inbox。
|
||||
*/
|
||||
public void markFrameCaptured() {
|
||||
capturedFrameCount.incrementAndGet();
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录一条控制事件或配置关闭时被忽略的 frame。
|
||||
*/
|
||||
public void markFrameIgnored() {
|
||||
ignoredFrameCount.incrementAndGet();
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录一条因为安全边界被拒绝的 frame。
|
||||
*/
|
||||
public void markFrameRejected(String errorCode, String errorSummary) {
|
||||
rejectedFrameCount.incrementAndGet();
|
||||
recordError(errorCode, errorSummary);
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录一条处理失败的 frame。
|
||||
*/
|
||||
public void markFrameFailed(String errorCode, String errorSummary) {
|
||||
failedFrameCount.incrementAndGet();
|
||||
recordError(errorCode, errorSummary);
|
||||
}
|
||||
|
||||
/**
|
||||
* 记录最近一次安全错误代码和摘要,不包含 raw frame、Token 或正文。
|
||||
*/
|
||||
public void recordError(String errorCode, String errorSummary) {
|
||||
lastErrorCode = errorCode;
|
||||
lastErrorSummary = errorSummary;
|
||||
}
|
||||
|
||||
/**
|
||||
* 返回当前安全状态快照,供系统探针接口展示。
|
||||
*/
|
||||
public AgentBusStatusSnapshot snapshot() {
|
||||
return new AgentBusStatusSnapshot(
|
||||
connected,
|
||||
sessionReady,
|
||||
receivedFrameCount.get(),
|
||||
capturedFrameCount.get(),
|
||||
ignoredFrameCount.get(),
|
||||
rejectedFrameCount.get(),
|
||||
failedFrameCount.get(),
|
||||
lastErrorCode,
|
||||
lastErrorSummary,
|
||||
connectedAt,
|
||||
disconnectedAt,
|
||||
lastFrameReceivedAt
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 生成 UTC 时间,保证状态接口时间字段稳定。
|
||||
*/
|
||||
private OffsetDateTime nowUtc() {
|
||||
return OffsetDateTime.now(ZoneOffset.UTC);
|
||||
}
|
||||
}
|
||||
+40
@@ -0,0 +1,40 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
/**
|
||||
* AgentBus 单条 frame 处理结果。结果中不包含原始 payload 或客户正文。
|
||||
*/
|
||||
public record AgentBusFrameProcessResult(
|
||||
String outcome,
|
||||
Long inboxId,
|
||||
String captureStatus,
|
||||
String errorCode
|
||||
) {
|
||||
|
||||
/**
|
||||
* 构造忽略结果,用于控制事件或捕获关闭场景。
|
||||
*/
|
||||
public static AgentBusFrameProcessResult ignored() {
|
||||
return new AgentBusFrameProcessResult("IGNORED", null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* 构造捕获成功结果,只返回内部 Inbox ID 和捕获状态。
|
||||
*/
|
||||
public static AgentBusFrameProcessResult captured(Long inboxId, String captureStatus) {
|
||||
return new AgentBusFrameProcessResult("CAPTURED", inboxId, captureStatus, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* 构造拒绝结果,用于超大 frame 等安全边界。
|
||||
*/
|
||||
public static AgentBusFrameProcessResult rejected(String errorCode) {
|
||||
return new AgentBusFrameProcessResult("REJECTED", null, null, errorCode);
|
||||
}
|
||||
|
||||
/**
|
||||
* 构造失败结果,用于 JSON 解析或捕获异常。
|
||||
*/
|
||||
public static AgentBusFrameProcessResult failed(String errorCode) {
|
||||
return new AgentBusFrameProcessResult("FAILED", null, null, errorCode);
|
||||
}
|
||||
}
|
||||
+118
@@ -0,0 +1,118 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
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;
|
||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.Set;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* AgentBus 入站 frame 处理器。只负责控制事件识别和 SourceMessage 捕获,不做业务写操作或客户回复。
|
||||
*/
|
||||
@Component
|
||||
public class AgentBusFrameProcessor {
|
||||
|
||||
private static final String FRAME_TOO_LARGE = "FRAME_TOO_LARGE";
|
||||
private static final String INVALID_JSON = "INVALID_JSON";
|
||||
private static final String CAPTURE_FAILED = "CAPTURE_FAILED";
|
||||
private static final Set<String> CONTROL_EVENTS = Set.of("session.ready", "task.progress", "task.result");
|
||||
|
||||
private final ObjectMapper objectMapper;
|
||||
private final AgentBusSourceMessageAdapter sourceMessageAdapter;
|
||||
private final SourceMessageCaptureService captureService;
|
||||
private final AgentBusConnectionStatus status;
|
||||
private final AgentBusProperties properties;
|
||||
|
||||
/**
|
||||
* 注入 AgentBus frame 处理依赖。外部协议转换与平台捕获服务通过稳定命令隔离。
|
||||
*/
|
||||
public AgentBusFrameProcessor(
|
||||
ObjectMapper objectMapper,
|
||||
AgentBusSourceMessageAdapter sourceMessageAdapter,
|
||||
SourceMessageCaptureService captureService,
|
||||
AgentBusConnectionStatus status,
|
||||
AgentBusProperties properties) {
|
||||
this.objectMapper = objectMapper;
|
||||
this.sourceMessageAdapter = sourceMessageAdapter;
|
||||
this.captureService = captureService;
|
||||
this.status = status;
|
||||
this.properties = properties;
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理一条 AgentBus raw frame。该方法不保存 raw frame,也不会发送 ACK 或客户回复。
|
||||
*/
|
||||
public AgentBusFrameProcessResult process(String rawFrame) {
|
||||
status.markFrameReceived();
|
||||
if (rawFrame == null || rawFrame.getBytes(StandardCharsets.UTF_8).length > properties.getMaxFrameBytes()) {
|
||||
status.markFrameRejected(FRAME_TOO_LARGE, "AgentBus frame exceeds max allowed bytes.");
|
||||
return AgentBusFrameProcessResult.rejected(FRAME_TOO_LARGE);
|
||||
}
|
||||
JsonNode frame = parse(rawFrame);
|
||||
if (frame == null) {
|
||||
return AgentBusFrameProcessResult.failed(INVALID_JSON);
|
||||
}
|
||||
String eventType = firstText(frame, "type", frame, "event");
|
||||
if (eventType != null && CONTROL_EVENTS.contains(eventType)) {
|
||||
if ("session.ready".equals(eventType)) {
|
||||
status.markSessionReady();
|
||||
}
|
||||
status.markFrameIgnored();
|
||||
return AgentBusFrameProcessResult.ignored();
|
||||
}
|
||||
if (!properties.getCapture().isEnabled()) {
|
||||
status.markFrameIgnored();
|
||||
return AgentBusFrameProcessResult.ignored();
|
||||
}
|
||||
try {
|
||||
CaptureSourceMessageCommand command = sourceMessageAdapter.toCaptureCommand(
|
||||
properties.getCapture().getDefaultHotelId(),
|
||||
frame);
|
||||
SourceMessageCaptureResult result = captureService.capture(command);
|
||||
status.markFrameCaptured();
|
||||
return AgentBusFrameProcessResult.captured(result.inboxId(), result.captureStatus());
|
||||
} catch (RuntimeException exception) {
|
||||
status.markFrameFailed(CAPTURE_FAILED, "AgentBus frame capture failed.");
|
||||
return AgentBusFrameProcessResult.failed(CAPTURE_FAILED);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 解析 AgentBus JSON frame。解析失败只记录安全错误代码,不回显原始报文。
|
||||
*/
|
||||
private JsonNode parse(String rawFrame) {
|
||||
try {
|
||||
return objectMapper.readTree(rawFrame);
|
||||
} catch (JsonProcessingException exception) {
|
||||
status.markFrameFailed(INVALID_JSON, "AgentBus frame is not valid JSON.");
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 从两个候选字段中读取第一个非空文本。
|
||||
*/
|
||||
private String firstText(JsonNode firstNode, String firstField, JsonNode secondNode, String secondField) {
|
||||
String firstValue = text(firstNode, firstField);
|
||||
if (firstValue != null) {
|
||||
return firstValue;
|
||||
}
|
||||
return text(secondNode, secondField);
|
||||
}
|
||||
|
||||
/**
|
||||
* 读取并清理 JSON 文本字段,空白字符串按缺失处理。
|
||||
*/
|
||||
private String text(JsonNode node, String fieldName) {
|
||||
if (node == null || node.path(fieldName).isMissingNode() || node.path(fieldName).isNull()) {
|
||||
return null;
|
||||
}
|
||||
String value = node.path(fieldName).asText();
|
||||
String trimmed = value == null ? "" : value.trim();
|
||||
return trimmed.isEmpty() ? null : trimmed;
|
||||
}
|
||||
}
|
||||
+145
@@ -0,0 +1,145 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
import java.time.Duration;
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* AgentBus 入站连接配置。Secret 只从环境变量注入,不在响应、日志或测试夹具中输出。
|
||||
*/
|
||||
@Component
|
||||
@ConfigurationProperties(prefix = "agentbus")
|
||||
public class AgentBusProperties {
|
||||
|
||||
/** AgentBus 探针和长连接开关,默认关闭。 */
|
||||
private Probe probe = new Probe();
|
||||
/** AgentBus WebSocket 连接参数。 */
|
||||
private Ws ws = new Ws();
|
||||
/** AgentBus 入站捕获配置。 */
|
||||
private Capture capture = new Capture();
|
||||
/** 单个入站 frame 最大字节数,超过后拒绝处理。 */
|
||||
private int maxFrameBytes = 1048576;
|
||||
|
||||
public Probe getProbe() {
|
||||
return probe;
|
||||
}
|
||||
|
||||
public void setProbe(Probe probe) {
|
||||
this.probe = probe;
|
||||
}
|
||||
|
||||
public Ws getWs() {
|
||||
return ws;
|
||||
}
|
||||
|
||||
public void setWs(Ws ws) {
|
||||
this.ws = ws;
|
||||
}
|
||||
|
||||
public Capture getCapture() {
|
||||
return capture;
|
||||
}
|
||||
|
||||
public void setCapture(Capture capture) {
|
||||
this.capture = capture;
|
||||
}
|
||||
|
||||
public int getMaxFrameBytes() {
|
||||
return maxFrameBytes;
|
||||
}
|
||||
|
||||
public void setMaxFrameBytes(int maxFrameBytes) {
|
||||
this.maxFrameBytes = maxFrameBytes;
|
||||
}
|
||||
|
||||
/**
|
||||
* AgentBus WebSocket 启停配置。
|
||||
*/
|
||||
public static class Probe {
|
||||
|
||||
/** 是否启用 AgentBus WebSocket 长连接。 */
|
||||
private boolean enabled = false;
|
||||
|
||||
public boolean isEnabled() {
|
||||
return enabled;
|
||||
}
|
||||
|
||||
public void setEnabled(boolean enabled) {
|
||||
this.enabled = enabled;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* AgentBus WebSocket 连接参数。
|
||||
*/
|
||||
public static class Ws {
|
||||
|
||||
/** AgentBus WebSocket 地址。 */
|
||||
private String url = "wss://mesh.nianxx.cn/ws";
|
||||
/** AgentBus WebSocket 鉴权 Token。 */
|
||||
private String token;
|
||||
/** AgentBus 断线重连等待时间。 */
|
||||
private Duration reconnectDelay = Duration.ofSeconds(5);
|
||||
/** AgentBus 连接超时时间。 */
|
||||
private Duration connectTimeout = Duration.ofSeconds(15);
|
||||
|
||||
public String getUrl() {
|
||||
return url;
|
||||
}
|
||||
|
||||
public void setUrl(String url) {
|
||||
this.url = url;
|
||||
}
|
||||
|
||||
public String getToken() {
|
||||
return token;
|
||||
}
|
||||
|
||||
public void setToken(String token) {
|
||||
this.token = token;
|
||||
}
|
||||
|
||||
public Duration getReconnectDelay() {
|
||||
return reconnectDelay;
|
||||
}
|
||||
|
||||
public void setReconnectDelay(Duration reconnectDelay) {
|
||||
this.reconnectDelay = reconnectDelay;
|
||||
}
|
||||
|
||||
public Duration getConnectTimeout() {
|
||||
return connectTimeout;
|
||||
}
|
||||
|
||||
public void setConnectTimeout(Duration connectTimeout) {
|
||||
this.connectTimeout = connectTimeout;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* AgentBus 入站消息捕获配置。
|
||||
*/
|
||||
public static class Capture {
|
||||
|
||||
/** 是否把业务 frame 写入 SourceMessage Inbox。 */
|
||||
private boolean enabled = true;
|
||||
/** AgentBus 未提供酒店上下文时使用的默认酒店 ID。 */
|
||||
private String defaultHotelId = "HOTEL-TEST";
|
||||
|
||||
public boolean isEnabled() {
|
||||
return enabled;
|
||||
}
|
||||
|
||||
public void setEnabled(boolean enabled) {
|
||||
this.enabled = enabled;
|
||||
}
|
||||
|
||||
public String getDefaultHotelId() {
|
||||
return defaultHotelId;
|
||||
}
|
||||
|
||||
public void setDefaultHotelId(String defaultHotelId) {
|
||||
this.defaultHotelId = defaultHotelId;
|
||||
}
|
||||
}
|
||||
}
|
||||
+191
@@ -0,0 +1,191 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageCommand;
|
||||
import cn.nianxx.thhotel.platform.message.common.request.CaptureSourceMessageMedia;
|
||||
import cn.nianxx.thhotel.platform.message.common.enums.SourceMessageMediaType;
|
||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import java.time.Instant;
|
||||
import java.time.format.DateTimeParseException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* AgentBus Outlook 邮件 payload 适配器。外部协议字段只在适配层解析,不泄漏到平台消息领域。
|
||||
*/
|
||||
@Component
|
||||
public class AgentBusSourceMessageAdapter {
|
||||
|
||||
private static final String PROVIDER_AGENTBUS = "AGENTBUS";
|
||||
private static final String SCHEMA_VERSION = "agentbus-outlook-v1";
|
||||
|
||||
private final ObjectMapper objectMapper;
|
||||
|
||||
public AgentBusSourceMessageAdapter(ObjectMapper objectMapper) {
|
||||
this.objectMapper = objectMapper;
|
||||
}
|
||||
|
||||
/**
|
||||
* 将 AgentBus frame 转换为平台 SourceMessage 捕获命令。
|
||||
*/
|
||||
public CaptureSourceMessageCommand toCaptureCommand(String hotelId, JsonNode frame) {
|
||||
JsonNode payload = frame.path("payload").isMissingNode() ? frame : frame.path("payload");
|
||||
JsonNode source = payload.path("source");
|
||||
JsonNode body = payload.path("body");
|
||||
String channel = uppercaseOrDefault(text(source, "channel"), "EMAIL");
|
||||
|
||||
return new CaptureSourceMessageCommand(
|
||||
hotelId,
|
||||
PROVIDER_AGENTBUS,
|
||||
channel,
|
||||
text(source, "external_message_id"),
|
||||
text(source, "external_conversation_id"),
|
||||
text(frame, "id"),
|
||||
text(frame, "session_id"),
|
||||
parseInstant(text(source, "sent_at")),
|
||||
text(source, "sender"),
|
||||
text(source, "subject"),
|
||||
firstText(body, "text", payload, "text"),
|
||||
text(body, "html"),
|
||||
toPayloadJson(payload),
|
||||
SCHEMA_VERSION,
|
||||
mediaItems(payload)
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从 AgentBus payload 中提取内嵌图片和附件引用。
|
||||
*/
|
||||
private List<CaptureSourceMessageMedia> mediaItems(JsonNode payload) {
|
||||
List<CaptureSourceMessageMedia> result = new ArrayList<>();
|
||||
addMedia(result, SourceMessageMediaType.INLINE_IMAGE.code(), payload.path("inline_images"));
|
||||
addMedia(result, SourceMessageMediaType.ATTACHMENT.code(), payload.path("attachments"));
|
||||
return List.copyOf(result);
|
||||
}
|
||||
|
||||
/**
|
||||
* 将指定媒体数组追加为平台统一媒体请求项。
|
||||
*/
|
||||
private void addMedia(List<CaptureSourceMessageMedia> result, String mediaType, JsonNode items) {
|
||||
if (!items.isArray()) {
|
||||
return;
|
||||
}
|
||||
for (JsonNode item : items) {
|
||||
result.add(new CaptureSourceMessageMedia(
|
||||
mediaType,
|
||||
firstText(item, "file_name", item, "filename", item, "name"),
|
||||
firstText(item, "content_type", item, "mime_type"),
|
||||
firstLong(item, "size_bytes", item, "size"),
|
||||
firstText(item, "external_url", item, "url"),
|
||||
text(item, "id")
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 将规范化后的 AgentBus payload 序列化,供 SourceMessage 原始载荷追溯。
|
||||
*/
|
||||
private String toPayloadJson(JsonNode payload) {
|
||||
try {
|
||||
return objectMapper.writeValueAsString(payload);
|
||||
} catch (JsonProcessingException exception) {
|
||||
throw new IllegalArgumentException("AgentBus payload 无法序列化为规范 JSON", exception);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 解析 AgentBus 来源时间,格式非法时保持为空,避免阻断消息入站。
|
||||
*/
|
||||
private Instant parseInstant(String value) {
|
||||
if (value == null) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
return Instant.parse(value);
|
||||
} catch (DateTimeParseException exception) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 将渠道等稳定代码转为大写,缺失时使用服务端默认值。
|
||||
*/
|
||||
private String uppercaseOrDefault(String value, String defaultValue) {
|
||||
if (value == null) {
|
||||
return defaultValue;
|
||||
}
|
||||
return value.toUpperCase(Locale.ROOT);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从两个候选字段中读取第一个非空文本。
|
||||
*/
|
||||
private String firstText(JsonNode firstNode, String firstField, JsonNode secondNode, String secondField) {
|
||||
String firstValue = text(firstNode, firstField);
|
||||
if (firstValue != null) {
|
||||
return firstValue;
|
||||
}
|
||||
return text(secondNode, secondField);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从三个候选字段中读取第一个非空文本。
|
||||
*/
|
||||
private String firstText(
|
||||
JsonNode firstNode,
|
||||
String firstField,
|
||||
JsonNode secondNode,
|
||||
String secondField,
|
||||
JsonNode thirdNode,
|
||||
String thirdField) {
|
||||
String firstValue = text(firstNode, firstField);
|
||||
if (firstValue != null) {
|
||||
return firstValue;
|
||||
}
|
||||
String secondValue = text(secondNode, secondField);
|
||||
if (secondValue != null) {
|
||||
return secondValue;
|
||||
}
|
||||
return text(thirdNode, thirdField);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从两个候选字段中读取第一个可转换为 Long 的值。
|
||||
*/
|
||||
private Long firstLong(JsonNode firstNode, String firstField, JsonNode secondNode, String secondField) {
|
||||
Long firstValue = longValue(firstNode, firstField);
|
||||
if (firstValue != null) {
|
||||
return firstValue;
|
||||
}
|
||||
return longValue(secondNode, secondField);
|
||||
}
|
||||
|
||||
/**
|
||||
* 读取并清理 AgentBus JSON 文本字段,空白字符串按缺失处理。
|
||||
*/
|
||||
private String text(JsonNode node, String fieldName) {
|
||||
if (node == null || node.isMissingNode() || node.path(fieldName).isMissingNode() || node.path(fieldName).isNull()) {
|
||||
return null;
|
||||
}
|
||||
String value = node.path(fieldName).asText();
|
||||
String trimmed = value == null ? "" : value.trim();
|
||||
return trimmed.isEmpty() ? null : trimmed;
|
||||
}
|
||||
|
||||
/**
|
||||
* 读取 AgentBus JSON 数字字段,非数字或缺失时返回空。
|
||||
*/
|
||||
private Long longValue(JsonNode node, String fieldName) {
|
||||
if (node == null || node.isMissingNode() || node.path(fieldName).isMissingNode() || node.path(fieldName).isNull()) {
|
||||
return null;
|
||||
}
|
||||
JsonNode value = node.path(fieldName);
|
||||
if (value.canConvertToLong()) {
|
||||
return value.asLong();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
+22
@@ -0,0 +1,22 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
import java.time.OffsetDateTime;
|
||||
|
||||
/**
|
||||
* AgentBus 连接状态快照。只包含安全状态和计数器,不包含 Token 或原始 frame。
|
||||
*/
|
||||
public record AgentBusStatusSnapshot(
|
||||
boolean connected,
|
||||
boolean sessionReady,
|
||||
long receivedFrameCount,
|
||||
long capturedFrameCount,
|
||||
long ignoredFrameCount,
|
||||
long rejectedFrameCount,
|
||||
long failedFrameCount,
|
||||
String lastErrorCode,
|
||||
String lastErrorSummary,
|
||||
OffsetDateTime connectedAt,
|
||||
OffsetDateTime disconnectedAt,
|
||||
OffsetDateTime lastFrameReceivedAt
|
||||
) {
|
||||
}
|
||||
+235
@@ -0,0 +1,235 @@
|
||||
package cn.nianxx.thhotel.integrations.messaging.agentbus.adapter;
|
||||
|
||||
import java.net.URI;
|
||||
import java.net.http.HttpClient;
|
||||
import java.net.http.WebSocket;
|
||||
import java.time.Duration;
|
||||
import java.util.concurrent.CompletionStage;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import jakarta.annotation.PreDestroy;
|
||||
import org.springframework.context.SmartLifecycle;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* AgentBus WebSocket 长连接客户端。默认关闭,开启后只接收入站 frame 并交给 SourceMessage 捕获链路。
|
||||
*/
|
||||
@Component
|
||||
public class AgentBusWebSocketClient implements SmartLifecycle {
|
||||
|
||||
private static final String MISSING_TOKEN = "MISSING_TOKEN";
|
||||
private static final String INVALID_WS_URL = "INVALID_WS_URL";
|
||||
private static final String CONNECT_FAILED = "CONNECT_FAILED";
|
||||
private static final String WEBSOCKET_ERROR = "WEBSOCKET_ERROR";
|
||||
|
||||
private final AgentBusProperties properties;
|
||||
private final AgentBusFrameProcessor frameProcessor;
|
||||
private final AgentBusConnectionStatus status;
|
||||
private final ScheduledExecutorService reconnectExecutor;
|
||||
private final AtomicBoolean running = new AtomicBoolean(false);
|
||||
|
||||
private volatile WebSocket webSocket;
|
||||
|
||||
/**
|
||||
* 注入 AgentBus 连接配置、frame 处理器和安全状态组件。
|
||||
*/
|
||||
public AgentBusWebSocketClient(
|
||||
AgentBusProperties properties,
|
||||
AgentBusFrameProcessor frameProcessor,
|
||||
AgentBusConnectionStatus status) {
|
||||
this.properties = properties;
|
||||
this.frameProcessor = frameProcessor;
|
||||
this.status = status;
|
||||
this.reconnectExecutor = Executors.newSingleThreadScheduledExecutor(runnable -> {
|
||||
Thread thread = new Thread(runnable, "agentbus-websocket-reconnect");
|
||||
thread.setDaemon(true);
|
||||
return thread;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Spring 容器启动后按配置决定是否连接 AgentBus WebSocket。
|
||||
*/
|
||||
@Override
|
||||
public void start() {
|
||||
if (!properties.getProbe().isEnabled()) {
|
||||
return;
|
||||
}
|
||||
if (running.compareAndSet(false, true)) {
|
||||
connect();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 停止 WebSocket 长连接并关闭重连循环。
|
||||
*/
|
||||
@Override
|
||||
public void stop() {
|
||||
running.set(false);
|
||||
WebSocket current = webSocket;
|
||||
if (current != null) {
|
||||
current.sendClose(WebSocket.NORMAL_CLOSURE, "application stopping");
|
||||
}
|
||||
status.markDisconnected();
|
||||
}
|
||||
|
||||
/**
|
||||
* 返回 WebSocket 客户端是否处于运行状态。
|
||||
*/
|
||||
@Override
|
||||
public boolean isRunning() {
|
||||
return running.get();
|
||||
}
|
||||
|
||||
/**
|
||||
* 当前客户端默认随应用自动启动,但只有配置开启时才真正建立连接。
|
||||
*/
|
||||
@Override
|
||||
public boolean isAutoStartup() {
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Spring 容器销毁时关闭重连线程,避免本地测试或重启时残留后台任务。
|
||||
*/
|
||||
@PreDestroy
|
||||
public void shutdown() {
|
||||
stop();
|
||||
reconnectExecutor.shutdownNow();
|
||||
}
|
||||
|
||||
/**
|
||||
* 建立 AgentBus WebSocket 连接。Token 缺失时只记录错误代码,不尝试连接。
|
||||
*/
|
||||
private void connect() {
|
||||
if (!running.get()) {
|
||||
return;
|
||||
}
|
||||
String token = trimToNull(properties.getWs().getToken());
|
||||
if (token == null) {
|
||||
status.recordError(MISSING_TOKEN, "AgentBus WebSocket token is not configured.");
|
||||
return;
|
||||
}
|
||||
URI uri = buildReadyUri(properties.getWs().getUrl());
|
||||
if (uri == null) {
|
||||
status.recordError(INVALID_WS_URL, "AgentBus WebSocket URL is invalid.");
|
||||
return;
|
||||
}
|
||||
HttpClient.newBuilder()
|
||||
.connectTimeout(normalizeDuration(properties.getWs().getConnectTimeout(), Duration.ofSeconds(15)))
|
||||
.build()
|
||||
.newWebSocketBuilder()
|
||||
.header("Authorization", "Bearer " + token)
|
||||
.buildAsync(uri, new Listener())
|
||||
.whenComplete((socket, exception) -> {
|
||||
if (exception != null) {
|
||||
status.recordError(CONNECT_FAILED, "AgentBus WebSocket connection failed.");
|
||||
scheduleReconnect();
|
||||
} else {
|
||||
webSocket = socket;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 安排断线后重连。停止状态下不会继续重连。
|
||||
*/
|
||||
private void scheduleReconnect() {
|
||||
if (!running.get()) {
|
||||
return;
|
||||
}
|
||||
long delayMillis = normalizeDuration(properties.getWs().getReconnectDelay(), Duration.ofSeconds(5)).toMillis();
|
||||
reconnectExecutor.schedule(this::connect, delayMillis, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
/**
|
||||
* 将 WebSocket URL 补上 ready=1 查询参数,用于请求 AgentBus session.ready。
|
||||
*/
|
||||
private URI buildReadyUri(String rawUrl) {
|
||||
String url = trimToNull(rawUrl);
|
||||
if (url == null) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
String separator = url.contains("?") ? "&" : "?";
|
||||
return URI.create(url + separator + "ready=1");
|
||||
} catch (IllegalArgumentException exception) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 标准化 Duration 配置,缺失或非正数时使用默认值。
|
||||
*/
|
||||
private Duration normalizeDuration(Duration value, Duration defaultValue) {
|
||||
if (value == null || value.isZero() || value.isNegative()) {
|
||||
return defaultValue;
|
||||
}
|
||||
return value;
|
||||
}
|
||||
|
||||
/**
|
||||
* 将空白字符串统一视为未配置。
|
||||
*/
|
||||
private String trimToNull(String value) {
|
||||
if (value == null) {
|
||||
return null;
|
||||
}
|
||||
String trimmed = value.trim();
|
||||
return trimmed.isEmpty() ? null : trimmed;
|
||||
}
|
||||
|
||||
/**
|
||||
* JDK WebSocket Listener。只接收文本 frame,不发送 ACK、task.result 或客户回复。
|
||||
*/
|
||||
private class Listener implements WebSocket.Listener {
|
||||
|
||||
private final StringBuilder textBuffer = new StringBuilder();
|
||||
|
||||
/**
|
||||
* 标记 WebSocket 已连接,并请求第一条消息。
|
||||
*/
|
||||
@Override
|
||||
public void onOpen(WebSocket webSocket) {
|
||||
status.markConnected();
|
||||
WebSocket.Listener.super.onOpen(webSocket);
|
||||
}
|
||||
|
||||
/**
|
||||
* 收到完整文本 frame 后交给 AgentBusFrameProcessor 处理。
|
||||
*/
|
||||
@Override
|
||||
public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) {
|
||||
textBuffer.append(data);
|
||||
if (last) {
|
||||
frameProcessor.process(textBuffer.toString());
|
||||
textBuffer.setLength(0);
|
||||
}
|
||||
webSocket.request(1);
|
||||
return WebSocket.Listener.super.onText(webSocket, data, last);
|
||||
}
|
||||
|
||||
/**
|
||||
* 连接关闭时记录断开并按配置重连。
|
||||
*/
|
||||
@Override
|
||||
public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) {
|
||||
status.markDisconnected();
|
||||
scheduleReconnect();
|
||||
return WebSocket.Listener.super.onClose(webSocket, statusCode, reason);
|
||||
}
|
||||
|
||||
/**
|
||||
* WebSocket 异常只记录安全错误代码,不输出 Token 或原始 frame。
|
||||
*/
|
||||
@Override
|
||||
public void onError(WebSocket webSocket, Throwable error) {
|
||||
status.markDisconnected();
|
||||
status.recordError(WEBSOCKET_ERROR, "AgentBus WebSocket error.");
|
||||
scheduleReconnect();
|
||||
WebSocket.Listener.super.onError(webSocket, error);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user