feat(booking): add safe retries and operational metrics
verify / booking-verify (push) Has been cancelled

This commit is contained in:
鲨鱼辣椒 committed 2026-08-08 16:25:17 +08:00
1 parent c826b6e574
commit 7330ac853b
26 files changed
+749 -31

No files matched your search

+4
View File
@@ -680,6 +680,10 @@ export default {
refreshing: 'Refreshing…',
source: 'Source message',
revision: 'Revision',
retryCount: 'Retry count',
retry: 'Retry processing',
retrying: 'Retrying…',
retrySuccess: 'The run was returned to the processing queue and will refresh automatically.',
disposition: 'Disposition',
status: 'Current status',
timeline: 'Processing progress',
+4
View File
@@ -680,6 +680,10 @@ export default {
refreshing: 'กำลังรีเฟรช…',
source: 'อีเมลต้นทาง',
revision: 'ฉบับแก้ไข',
retryCount: 'จำนวนครั้งที่ลองใหม่',
retry: 'ประมวลผลอีกครั้ง',
retrying: 'กำลังประมวลผลอีกครั้ง…',
retrySuccess: 'ส่งกลับเข้าคิวประมวลผลแล้ว และสถานะจะรีเฟรชอัตโนมัติ',
disposition: 'ผลการจัดการ',
status: 'สถานะปัจจุบัน',
timeline: 'ความคืบหน้าการประมวลผล',
+4
View File
@@ -679,6 +679,10 @@ export default {
refreshing: '刷新中…',
source: '来源邮件',
revision: '修订',
retryCount: '重试次数',
retry: '重新处理',
retrying: '正在重新处理…',
retrySuccess: '已重新进入处理队列,将自动刷新状态。',
disposition: '处理结果',
status: '当前状态',
timeline: '处理进度',
@@ -63,6 +63,13 @@ export async function fetchBookingProcessingRun(runId: string): Promise<BookingP
)
}
export async function retryBookingProcessingRun(runId: string): Promise<BookingProcessingRunView> {
return sendJson<BookingProcessingRunView>(
`/api/reservation/booking-processing-runs/${encodeURIComponent(runId)}/retry`,
'POST',
)
}
export async function confirmBookingProcessingRun(
runId: string,
request: BookingConfirmationRequest,
@@ -13,6 +13,7 @@ vi.mock('@/services/reservationService', async (importOriginal) => {
...actual,
fetchBookingProcessingRun: vi.fn(),
confirmBookingProcessingRun: vi.fn(),
retryBookingProcessingRun: vi.fn(),
}
})
@@ -22,6 +23,7 @@ describe('ReservationBookingProcessingRunView', () => {
beforeEach(() => {
vi.mocked(service.fetchBookingProcessingRun).mockReset()
vi.mocked(service.confirmBookingProcessingRun).mockReset()
vi.mocked(service.retryBookingProcessingRun).mockReset()
})
it('shows the safe confirmation projection and freezes corrected values without a PMS action', async () => {
@@ -64,6 +66,38 @@ describe('ReservationBookingProcessingRunView', () => {
expect(wrapper.text()).toContain('参数已冻结')
expect(wrapper.text()).toContain('已确认')
})
it('retries only a failed processing run and resumes safe polling', async () => {
const failed = {
...awaitingRun(),
status: 'FAILED' as const,
result_disposition: 'RISK_NOTIFICATION' as const,
retry_count: 0,
error_code: 'BOOKING_ORCHESTRATION_FAILED',
error_summary: 'Booking processing failed before confirmation.',
confirmation_projection: null,
}
const retried = {
...failed,
status: 'CONTEXT_READY' as const,
retry_count: 1,
error_code: null,
error_summary: null,
}
vi.mocked(service.fetchBookingProcessingRun).mockResolvedValue(failed)
vi.mocked(service.retryBookingProcessingRun).mockResolvedValue(retried)
const wrapper = await mountView()
await flushPromises()
const retryButton = wrapper.findAll('button').find((button) => button.text().includes('重新处理'))
expect(retryButton).toBeDefined()
await retryButton!.trigger('click')
await flushPromises()
expect(service.retryBookingProcessingRun).toHaveBeenCalledWith('booking-001')
expect(wrapper.text()).toContain('上下文已组装')
expect(wrapper.text()).toContain('已重新进入处理队列')
})
})
async function mountView() {
@@ -64,6 +64,10 @@
<dt>{{ t('bookingProcessingRun.revision') }}</dt>
<dd>{{ run.revision_number }}</dd>
</div>
<div>
<dt>{{ t('bookingProcessingRun.retryCount') }}</dt>
<dd>{{ run.retry_count }}</dd>
</div>
<div>
<dt>{{ t('bookingProcessingRun.disposition') }}</dt>
<dd>{{ run.result_disposition ?? '—' }}</dd>
@@ -105,6 +109,15 @@
>
{{ run.error_code ?? 'BOOKING_PROCESSING_ERROR' }} · {{ run.error_summary }}
</p>
<button
v-if="canRetry"
type="button"
class="secondary-button"
:disabled="retrying"
@click="retryRun"
>
{{ retrying ? t('bookingProcessingRun.retrying') : t('bookingProcessingRun.retry') }}
</button>
</section>
<section
@@ -258,6 +271,7 @@ import { ApiError } from '@/services/httpClient'
import {
confirmBookingProcessingRun,
fetchBookingProcessingRun,
retryBookingProcessingRun,
} from '@/services/reservationService'
import type {
BookingConfirmationField,
@@ -282,6 +296,7 @@ const route = useRoute()
const run = ref<BookingProcessingRunView | null>(null)
const loading = ref(false)
const confirming = ref(false)
const retrying = ref(false)
const errorMessage = ref('')
const successMessage = ref('')
const draftValues = ref<Record<string, Record<string, string>>>({})
@@ -299,6 +314,7 @@ const projectionIssues = computed(() => projection.value?.issues ?? [])
const evidenceCount = computed(() => projection.value?.evidence_refs.length ?? 0)
const canSubmit = computed(() => run.value?.status === 'AWAITING_CONFIRMATION'
&& confirmationItems.value.length > 0)
const canRetry = computed(() => run.value?.status === 'FAILED')
onMounted(() => {
void loadRun(true)
@@ -401,6 +417,26 @@ async function confirmRun(): Promise<void> {
}
}
async function retryRun(): Promise<void> {
if (!run.value || !canRetry.value || retrying.value) {
return
}
retrying.value = true
errorMessage.value = ''
successMessage.value = ''
try {
const nextRun = await retryBookingProcessingRun(run.value.processing_run_id)
run.value = nextRun
hydrate(nextRun)
schedulePolling(nextRun)
successMessage.value = t('bookingProcessingRun.retrySuccess')
} catch (error) {
errorMessage.value = apiErrorMessage(error)
} finally {
retrying.value = false
}
}
function itemFields(item: BookingConfirmationProjection['confirmation_items'][number]): BookingConfirmationField[] {
const byCode = new Map<string, BookingConfirmationField>()
for (const field of item.display_fields) {
@@ -18,3 +18,16 @@
- 本地验证使用 H2,不使用远程数据库默认值。
- 远程测试环境只允许通过 Secret 注入连接信息。
- 真实邮件样本、密码和 `.planning` 过程文件不进入产品提交。
## 测试环境部署与回滚 Runbook
1. 先由数据库管理员完成历史凭据轮换,并在 Secret 系统注入新的 `TH_HOTEL_BOOKING_PG_URL`、`TH_HOTEL_BOOKING_PG_USERNAME`、`TH_HOTEL_BOOKING_PG_PASSWORD`;仓库不保存其值。
2. 在应用启动前,以同一 Secret 在受控环境执行 `scripts/verify-booking-postgresql-migration.sh`。脚本必须先完成项目 schema preflight,再串行验证 V1/V2;没有 URL 时会以 exit 64 失败关闭。
3. 仅在确认 schema 仅为 `th_hotel_booking` 且 pre/postflight 通过后,设置 `TH_HOTEL_BOOKING_PG_ENABLED=true`。旧 MySQL 仍只读兼容,不允许双写或跨库事务。
4. 如需回滚应用行为,先停止 AgentBus 入口/worker,再将 `TH_HOTEL_BOOKING_PG_ENABLED=false` 并重启;不要执行 Flyway clean、不要删除 schema。已保存的 contracts artifact 与 attempt 供排障和三个月 retention 清理使用。
5. 排障从受控查询 `GET /api/reservation/booking-processing-runs/{runId}` 开始;只有 `FAILED` run 可调用 `POST /api/reservation/booking-processing-runs/{runId}/retry`。重试只读取已脱敏的 contracts artifact,绝不重拉原件。
6. 系统管理员用 `GET /api/reservation/booking-processing-metrics` 观察最近 24 小时的失败、Agent fallback、Risk、积压和重试趋势。该 API 不返回原始邮件、附件、团号或 Secret。
## 当前外部阻塞
- 远程 PostgreSQL 测试库在本轮的非写入 JDBC 预检中,于认证/SSL 协商前读取超时;没有执行 SQL、没有创建 schema。待网络或服务端恢复后,按上述 Runbook 的 preflight/migrate/postflight 路径继续,不能以重复重试替代环境修复。
@@ -96,6 +96,13 @@ Excel、图片和正文的原始材料体积大、结构多变且含历史内容
- 同团同封的 standalone Trace 在信息系统内合并为一张候选卡,保留每条服务项和证据;跨封 Trace 绝不并入旧卡,Context 只以无自由文本的时间序摘要提供历史顺序。
- Trace 的最终确认至少选择 `FO` 或 `HSK`,可同时选择二者;本期确认不触发部门流转。
### 5.4 运行重试与可观测性
- `FAILED` run 只能从已保存的 `PARSED_FACT_SET`、`CONTEXT_PACKAGE`、`MATERIAL_PACKAGE` 重放;不会重新拉取邮件、附件或历史,也不会创建新的 source revision。重试在 worker 中继续,processing attempt 与 `retry_count` 形成可审计的尝试链。
- `FAILED` 之外的状态拒绝重试;旧的确认投影在重试期间不返回,避免用户确认过期候选。
- 管理员只读指标固定统计最近 24 小时:失败率、Agent fallback 比率、Risk 比率、活跃积压、已重试 run 与 retry attempt。指标只来自 `th_hotel_booking` 聚合,不含个人数据、邮件内容、附件、团号或跨库查询。
- 任何监控、重试或 retention 任务都以用户确认前为终点,不能触发 PMS/Opera、付款、库存扣减或部门流转。
## 6. 数据与事实源
- 新 Booking 主线的事实源是 PostgreSQL schema `th_hotel_booking`。处理运行、版本化契约、证据引用、候选、校验和确认投影均落入该 schema。
@@ -148,7 +148,7 @@ Validator 必须校验:版本一致性、证据存在、目标解析、字段
| `blocked_fields[]` | 缺失/冲突字段及需用户选择的候选;不暴露内部 payload。 |
| `linked_actions[]` | 如 Allotment source/actual、Trace 合并关系;只显示本期可确认参数。 |
| `safe_evidence[]` | 文件/Sheet/行/图片等安全引用与脱敏摘要。 |
| `processing_run` | run ID、当前状态、可重试状态与安全错误摘要。 |
| `processing_run` | run ID、当前状态、重试次数、可重试状态与安全错误摘要。 |
确认 API 只冻结用户确认的参数和审计,必须携带 `processing_run_id` 与并发版本。确认成功不调用 PMS/Opera,也不执行付款、库存扣减或部门流转。
@@ -156,9 +156,10 @@ Validator 必须校验:版本一致性、证据存在、目标解析、字段
1. AgentBus 与手工 EML 都调用同一个 `BookingMessageOrchestrator`,并产生一致的 `SourceMessageEnvelope` 幂等语义。
2. 同步完成且无需异步 Agent 时返回 `201`;进入 Agent 或异步处理时返回 `202` 与 `processing_run_id`。
3. 状态查询为 `GET /api/reservation/booking-processing-runs/{runId}`;已有 `POST /api/reservation/booking-email-intakes` 逐步适配为统一入口。
3. 状态查询为 `GET /api/reservation/booking-processing-runs/{runId}`;已有 `POST /api/reservation/booking-email-intakes` 逐步适配为统一入口。仅 `FAILED` run 可经 `POST /api/reservation/booking-processing-runs/{runId}/retry` 以 `202` 重放;它只能重用已持久化的脱敏 contracts artifact,新增 processing attempt,不重新下载邮件/附件、不创建 revision、更不调用 PMS/Opera。
4. PostgreSQL `th_hotel_booking` 保存 run、attempt、版本、证据引用、候选、校验和确认投影;`PARSED_FACT_SET` 中只允许保存受 policy 约束的 Agent 摘要,所有这些项目 schema 数据按 `retention_until` 三个月清理。旧 MySQL 只能被 Context 兼容读取,不能参与新主线写入或跨库事务。
5. 所有持久化/接口代码必须使用 `contract_version` 做版本门禁;未知未来版本 fail closed 并形成安全 Risk/技术错误记录。
6. 仅拥有 `SYSTEM_ADMIN_CONSOLE_ACCESS` 的用户可读取 `GET /api/reservation/booking-processing-metrics`。该接口固定返回最近 24 小时的总 run、失败率、Agent fallback 比率、Risk 比率、积压、已重试 run 与 retry attempt 聚合,不返回任何邮件、附件、团号或用户信息。
## 10. 兼容与测试要求
@@ -0,0 +1,23 @@
package cn.nianxx.thhotel.workflows.reservation.booking.common.dto;
/** 项目 schema 内的聚合运行指标;不包含邮件、附件、团号或任何个人信息。 */
public record BookingOperationalMetricsSnapshot(
int windowHours,
long totalRuns,
long failedRuns,
long agentFallbackRuns,
long riskRuns,
long backlogRuns,
long retriedRuns,
long retryAttempts) {
public BookingOperationalMetricsSnapshot {
if (windowHours <= 0) {
throw new IllegalArgumentException("windowHours 必须大于 0");
}
if (totalRuns < 0 || failedRuns < 0 || agentFallbackRuns < 0 || riskRuns < 0
|| backlogRuns < 0 || retriedRuns < 0 || retryAttempts < 0) {
throw new IllegalArgumentException("Booking 运行指标不能为负数");
}
}
}
@@ -0,0 +1,29 @@
package cn.nianxx.thhotel.workflows.reservation.booking.common.result;
import com.fasterxml.jackson.annotation.JsonProperty;
/** 面向系统管理员的安全 Booking 运行指标视图。 */
public record BookingOperationalMetricsView(
@JsonProperty("window_hours")
int windowHours,
@JsonProperty("total_runs")
long totalRuns,
@JsonProperty("failed_runs")
long failedRuns,
@JsonProperty("failure_rate_percent")
double failureRatePercent,
@JsonProperty("agent_fallback_runs")
long agentFallbackRuns,
@JsonProperty("agent_fallback_rate_percent")
double agentFallbackRatePercent,
@JsonProperty("risk_runs")
long riskRuns,
@JsonProperty("risk_rate_percent")
double riskRatePercent,
@JsonProperty("backlog_runs")
long backlogRuns,
@JsonProperty("retried_runs")
long retriedRuns,
@JsonProperty("retry_attempts")
long retryAttempts) {
}
@@ -0,0 +1,9 @@
package cn.nianxx.thhotel.workflows.reservation.booking.repository;
import cn.nianxx.thhotel.workflows.reservation.booking.common.dto.BookingOperationalMetricsSnapshot;
/** 只读聚合项目 schema 的运行指标。 */
public interface BookingOperationalMetricsRepository {
BookingOperationalMetricsSnapshot snapshot(int windowHours);
}
@@ -0,0 +1,59 @@
package cn.nianxx.thhotel.workflows.reservation.booking.repository;
import cn.nianxx.thhotel.workflows.reservation.booking.common.dto.BookingOperationalMetricsSnapshot;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional;
/** PostgreSQL 项目 schema 的聚合监控查询;不会扫描旧 MySQL 或返回单条邮件数据。 */
@Repository
@ConditionalOnProperty(prefix = "booking.postgres", name = "enabled", havingValue = "true")
public class BookingPostgresOperationalMetricsRepository implements BookingOperationalMetricsRepository {
private static final String SNAPSHOT_SQL = "SELECT "
+ "COUNT(*) AS total_runs, "
+ "COUNT(*) FILTER (WHERE run.run_status = 'FAILED') AS failed_runs, "
+ "COUNT(*) FILTER (WHERE run.result_disposition = 'RISK_NOTIFICATION') AS risk_runs, "
+ "COUNT(*) FILTER (WHERE run.run_status IN ("
+ "'RECEIVED', 'MATERIAL_READY', 'PARSER_COMPLETE', 'PARSER_NOT_APPLICABLE', "
+ "'PARSER_FAILED', 'CONTEXT_READY', 'AGENT_COMPLETE')) AS backlog_runs, "
+ "COUNT(*) FILTER (WHERE run.retry_count > 0) AS retried_runs, "
+ "COALESCE(SUM(run.retry_count), 0) AS retry_attempts, "
+ "COUNT(*) FILTER (WHERE EXISTS ("
+ "SELECT 1 FROM th_hotel_booking.booking_contract_artifact artifact "
+ "WHERE artifact.processing_run_id = run.id "
+ "AND artifact.artifact_type = 'CANDIDATE_DECISION' "
+ "AND artifact.payload ->> 'decision_origin' IN ('AGENT', 'MIXED')"
+ ")) AS agent_fallback_runs "
+ "FROM th_hotel_booking.booking_contract_processing_run run "
+ "WHERE run.started_at >= CURRENT_TIMESTAMP - (? * INTERVAL '1 hour')";
private final JdbcTemplate jdbcTemplate;
public BookingPostgresOperationalMetricsRepository(
@Qualifier("bookingPostgresJdbcTemplate") JdbcTemplate jdbcTemplate) {
this.jdbcTemplate = jdbcTemplate;
}
@Override
@Transactional(readOnly = true, transactionManager = "bookingPostgresTransactionManager")
public BookingOperationalMetricsSnapshot snapshot(int windowHours) {
if (windowHours <= 0 || windowHours > 24 * 31) {
throw new IllegalArgumentException("Booking 监控窗口必须在 1 到 744 小时之间。");
}
return jdbcTemplate.queryForObject(
SNAPSHOT_SQL,
(resultSet, rowNumber) -> new BookingOperationalMetricsSnapshot(
windowHours,
resultSet.getLong("total_runs"),
resultSet.getLong("failed_runs"),
resultSet.getLong("agent_fallback_runs"),
resultSet.getLong("risk_runs"),
resultSet.getLong("backlog_runs"),
resultSet.getLong("retried_runs"),
resultSet.getLong("retry_attempts")),
windowHours);
}
}
@@ -404,6 +404,40 @@ public class BookingPostgresProcessingRunRepository implements BookingProcessing
}
}
/** 失败 run 的显式重试起点;不新建来源邮件、不改变 revision,也不删除前一 attempt 的审计。 */
@Override
@Transactional(transactionManager = "bookingPostgresTransactionManager")
public BookingProcessingRunSnapshot beginRetry(String runId) {
requireText(runId, "runId");
Long processingRunDatabaseId = jdbcTemplate.query(
"UPDATE th_hotel_booking.booking_contract_processing_run SET "
+ "retry_count = retry_count + 1, run_status = 'CONTEXT_READY', "
+ "result_disposition = NULL, error_code = NULL, error_summary = NULL, "
+ "completed_at = NULL, updated_at = CURRENT_TIMESTAMP "
+ "WHERE run_id = ? AND run_status = 'FAILED' RETURNING id",
(resultSet, rowNumber) -> resultSet.getLong("id"),
runId)
.stream()
.findFirst()
.orElse(null);
if (processingRunDatabaseId == null) {
throw new IllegalStateException("只有 FAILED 状态的 Booking processing run 可以重试。");
}
Integer nextAttemptNumber = jdbcTemplate.queryForObject(
"SELECT COALESCE(MAX(attempt_number), 0) + 1 "
+ "FROM th_hotel_booking.booking_contract_processing_attempt WHERE processing_run_id = ?",
Integer.class,
processingRunDatabaseId);
jdbcTemplate.update(
"INSERT INTO th_hotel_booking.booking_contract_processing_attempt ("
+ "processing_run_id, attempt_number, attempt_status"
+ ") VALUES (?, ?, 'STARTED')",
processingRunDatabaseId,
nextAttemptNumber == null ? 1 : nextAttemptNumber);
return findByRunId(runId)
.orElseThrow(() -> new IllegalStateException("Booking processing run 重试后无法读取。"));
}
/** 按 run ID 查询安全摘要。 */
@Override
@Transactional(readOnly = true, transactionManager = "bookingPostgresTransactionManager")
@@ -52,6 +52,12 @@ public interface BookingProcessingRunRepository {
String errorCode,
String errorSummary);
/**
* 仅将失败 run 置回 Context-ready 并新增一条 processing attempt;调用方随后必须只从已持久化的安全
* contracts artifact 重放,不能重新下载原始邮件或附件。
*/
BookingProcessingRunSnapshot beginRetry(String runId);
/** 读取某个安全 contracts artifact 的 JSON;只用于后端确认/查询,不返回原始邮件材料。 */
Optional<String> findArtifactPayload(String runId, String artifactType);
@@ -0,0 +1,9 @@
package cn.nianxx.thhotel.workflows.reservation.booking.service;
import cn.nianxx.thhotel.workflows.reservation.booking.common.result.BookingOperationalMetricsView;
/** Booking 运行状态的安全聚合监控服务。 */
public interface BookingOperationalMetricsService {
BookingOperationalMetricsView last24Hours();
}
@@ -9,5 +9,8 @@ public interface BookingProcessingRunService {
BookingProcessingRunView get(String runId);
/** 对失败 run 仅使用已持久化的安全 contracts artifact 发起异步重试。 */
BookingProcessingRunView retry(String runId);
BookingConfirmationResult confirm(String runId, BookingConfirmationRequest request, String actorId);
}
@@ -0,0 +1,44 @@
package cn.nianxx.thhotel.workflows.reservation.booking.service.impl;
import cn.nianxx.thhotel.workflows.reservation.common.dto.booking.BookingContracts;
import java.util.List;
import java.util.Locale;
/** 统一判断 current message 是否具有 Booking 业务信号;不读取 quoted history 或附件。 */
final class BookingMessageSignalPolicy {
private static final List<String> BUSINESS_SIGNAL_KEYWORDS = List.of(
"BOOKING",
"AMEND",
"UPDATE",
"CANCEL",
"CXL",
"ROOMING",
"PAYMENT",
"VOUCHER",
"TRACE",
"ALLOTMENT",
"TOUR CODE",
"ยกเลิก",
"จอง",
"ห้อง");
private BookingMessageSignalPolicy() {
}
static boolean containsBusinessSignal(String safeSubject, String currentBody) {
String source = String.join("\n",
safeSubject == null ? "" : safeSubject,
currentBody == null ? "" : currentBody)
.toUpperCase(Locale.ROOT);
return BUSINESS_SIGNAL_KEYWORDS.stream()
.map(keyword -> keyword.toUpperCase(Locale.ROOT))
.anyMatch(source::contains);
}
/** retry 只可使用已持久化的受控 current-message 摘要重建业务信号。 */
static boolean containsBusinessSignal(BookingContracts.AgentInputMaterial agentInputMaterial) {
return agentInputMaterial != null
&& containsBusinessSignal(null, agentInputMaterial.currentMessageText());
}
}
@@ -0,0 +1,46 @@
package cn.nianxx.thhotel.workflows.reservation.booking.service.impl;
import cn.nianxx.thhotel.workflows.reservation.booking.common.dto.BookingOperationalMetricsSnapshot;
import cn.nianxx.thhotel.workflows.reservation.booking.common.result.BookingOperationalMetricsView;
import cn.nianxx.thhotel.workflows.reservation.booking.repository.BookingOperationalMetricsRepository;
import cn.nianxx.thhotel.workflows.reservation.booking.service.BookingOperationalMetricsService;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service;
/** 将项目 schema 的计数聚合为可比较的百分比;不读取单封邮件内容。 */
@Service
@ConditionalOnProperty(prefix = "booking.postgres", name = "enabled", havingValue = "true")
public class BookingOperationalMetricsServiceImpl implements BookingOperationalMetricsService {
private static final int DEFAULT_WINDOW_HOURS = 24;
private final BookingOperationalMetricsRepository repository;
public BookingOperationalMetricsServiceImpl(BookingOperationalMetricsRepository repository) {
this.repository = repository;
}
@Override
public BookingOperationalMetricsView last24Hours() {
BookingOperationalMetricsSnapshot snapshot = repository.snapshot(DEFAULT_WINDOW_HOURS);
return new BookingOperationalMetricsView(
snapshot.windowHours(),
snapshot.totalRuns(),
snapshot.failedRuns(),
percentage(snapshot.failedRuns(), snapshot.totalRuns()),
snapshot.agentFallbackRuns(),
percentage(snapshot.agentFallbackRuns(), snapshot.totalRuns()),
snapshot.riskRuns(),
percentage(snapshot.riskRuns(), snapshot.totalRuns()),
snapshot.backlogRuns(),
snapshot.retriedRuns(),
snapshot.retryAttempts());
}
private double percentage(long numerator, long denominator) {
if (denominator <= 0) {
return 0D;
}
return Math.round((numerator * 10_000D / denominator)) / 100D;
}
}
@@ -19,7 +19,6 @@ import java.security.NoSuchAlgorithmException;
import java.util.HexFormat;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.UUID;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
@@ -36,22 +35,6 @@ import org.springframework.stereotype.Service;
@ConditionalOnProperty(prefix = "booking.postgres", name = "enabled", havingValue = "true")
public class BookingPostgresMessageOrchestrator implements BookingMessageOrchestrator {
private static final List<String> BUSINESS_SIGNAL_KEYWORDS = List.of(
"BOOKING",
"AMEND",
"UPDATE",
"CANCEL",
"CXL",
"ROOMING",
"PAYMENT",
"VOUCHER",
"TRACE",
"ALLOTMENT",
"TOUR CODE",
"ยกเลิก",
"จอง",
"ห้อง");
private final BookingContractAssembler contractAssembler;
private final BookingProcessingRunRepository processingRunRepository;
private final BookingContextAssembler contextAssembler;
@@ -160,7 +143,8 @@ public class BookingPostgresMessageOrchestrator implements BookingMessageOrchest
mergeWarningCodes(masterData.warningCodes(), completion.warningCodes()));
}
boolean businessSignalPresent = containsBusinessSignal(input);
boolean businessSignalPresent = BookingMessageSignalPolicy.containsBusinessSignal(
input.safeSubject(), input.currentBody());
boolean attachmentPresent = !input.attachments().isEmpty();
if (shouldDispatchAgentAsynchronously(parsedFactSet)) {
agentContinuationDispatcher.completeAsync(
@@ -277,16 +261,6 @@ public class BookingPostgresMessageOrchestrator implements BookingMessageOrchest
}
}
private boolean containsBusinessSignal(BookingMessageInput input) {
String source = String.join("\n",
input.safeSubject() == null ? "" : input.safeSubject(),
input.currentBody() == null ? "" : input.currentBody())
.toUpperCase(Locale.ROOT);
return BUSINESS_SIGNAL_KEYWORDS.stream()
.map(keyword -> keyword.toUpperCase(Locale.ROOT))
.anyMatch(source::contains);
}
private String sha256(String source) {
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
@@ -35,16 +35,19 @@ public class BookingProcessingRunServiceImpl implements BookingProcessingRunServ
private final BookingProcessingRunRepository processingRunRepository;
private final BookingDecisionValidator decisionValidator;
private final BookingConfirmationProjectionFactory projectionFactory;
private final BookingAgentContinuationDispatcher agentContinuationDispatcher;
private final ObjectMapper objectMapper;
public BookingProcessingRunServiceImpl(
BookingProcessingRunRepository processingRunRepository,
BookingDecisionValidator decisionValidator,
BookingConfirmationProjectionFactory projectionFactory,
BookingAgentContinuationDispatcher agentContinuationDispatcher,
ObjectMapper objectMapper) {
this.processingRunRepository = processingRunRepository;
this.decisionValidator = decisionValidator;
this.projectionFactory = projectionFactory;
this.agentContinuationDispatcher = agentContinuationDispatcher;
this.objectMapper = objectMapper;
}
@@ -52,7 +55,54 @@ public class BookingProcessingRunServiceImpl implements BookingProcessingRunServ
@Override
public BookingProcessingRunView get(String runId) {
BookingProcessingRunSnapshot snapshot = snapshot(runId);
return view(snapshot, confirmationProjection(runId));
return view(snapshot, projectionFor(snapshot));
}
/**
* 重试只允许 FAILED run,且只会重用已持久化、脱敏的 contracts artifact。为避免用户请求线程等待真实
* Agent,所有 retry 都经既有 worker 异步执行并立即返回 Context-ready 状态。
*/
@Override
public BookingProcessingRunView retry(String runId) {
BookingProcessingRunSnapshot before = snapshot(runId);
if (before.status() != BookingContracts.ProcessingStatus.FAILED) {
throw workflowError(
HttpStatus.CONFLICT,
"BOOKING_PROCESSING_RUN_NOT_RETRYABLE",
"只有失败的 Booking 邮件可以重试。");
}
BookingContracts.ParsedFactSet parsedFactSet = artifact(
runId,
"PARSED_FACT_SET",
BookingContracts.ParsedFactSet.class);
BookingContracts.ContextPackage contextPackage = artifact(
runId,
"CONTEXT_PACKAGE",
BookingContracts.ContextPackage.class);
BookingContracts.MaterialPackage materialPackage = artifact(
runId,
"MATERIAL_PACKAGE",
BookingContracts.MaterialPackage.class);
BookingProcessingRunSnapshot retryStarted = processingRunRepository.beginRetry(runId);
try {
agentContinuationDispatcher.completeAsync(
runId,
parsedFactSet,
contextPackage,
BookingMessageSignalPolicy.containsBusinessSignal(parsedFactSet.agentInputMaterial()),
materialPackage.currentMaterials().stream()
.anyMatch(material -> material.currentMaterial()
&& material.kind() != BookingContracts.EvidenceKind.CURRENT_BODY));
} catch (RuntimeException exception) {
processingRunRepository.updateRunStatus(
runId,
BookingContracts.ProcessingStatus.FAILED,
BookingContracts.ResultDisposition.RISK_NOTIFICATION,
"BOOKING_RETRY_DISPATCH_FAILED",
"Booking retry dispatch failed before processing.");
throw exception;
}
return view(retryStarted, null);
}
/**
@@ -308,6 +358,14 @@ public class BookingProcessingRunServiceImpl implements BookingProcessingRunServ
return payload == null ? null : read(payload, BookingContracts.ConfirmationProjection.class);
}
private BookingContracts.ConfirmationProjection projectionFor(BookingProcessingRunSnapshot snapshot) {
return switch (snapshot.status()) {
case VALIDATED, AWAITING_CONFIRMATION, CONFIRMED -> confirmationProjection(snapshot.runId());
case RECEIVED, MATERIAL_READY, PARSER_COMPLETE, PARSER_NOT_APPLICABLE,
PARSER_FAILED, CONTEXT_READY, AGENT_COMPLETE, FAILED -> null;
};
}
private <T> T artifact(String runId, String artifactType, Class<T> type) {
String payload = processingRunRepository.findArtifactPayload(runId, artifactType)
.orElseThrow(() -> workflowError(
@@ -0,0 +1,42 @@
package cn.nianxx.thhotel.workflows.reservation.control;
import cn.nianxx.thhotel.platform.access.common.enums.PlatformPermissionCode;
import cn.nianxx.thhotel.platform.security.service.FrontendAuthorizationService;
import cn.nianxx.thhotel.workflows.reservation.booking.common.result.BookingOperationalMetricsView;
import cn.nianxx.thhotel.workflows.reservation.booking.service.BookingOperationalMetricsService;
import cn.nianxx.thhotel.workflows.reservation.service.impl.ReservationTaskWorkflowException;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/** 系统管理员可读取的 Booking 聚合监控 API;不返回任一邮件、附件或团号。 */
@RestController
@RequestMapping("/api/reservation/booking-processing-metrics")
public class ReservationBookingOperationalMetricsController {
private final ObjectProvider<BookingOperationalMetricsService> metricsServiceProvider;
private final FrontendAuthorizationService authorizationService;
public ReservationBookingOperationalMetricsController(
ObjectProvider<BookingOperationalMetricsService> metricsServiceProvider,
FrontendAuthorizationService authorizationService) {
this.metricsServiceProvider = metricsServiceProvider;
this.authorizationService = authorizationService;
}
@GetMapping(produces = MediaType.APPLICATION_JSON_VALUE)
public BookingOperationalMetricsView getLast24Hours() {
authorizationService.requirePermission(PlatformPermissionCode.SYSTEM_ADMIN_CONSOLE_ACCESS.name());
BookingOperationalMetricsService service = metricsServiceProvider.getIfAvailable();
if (service == null) {
throw new ReservationTaskWorkflowException(
HttpStatus.SERVICE_UNAVAILABLE,
"BOOKING_POSTGRES_RUNTIME_DISABLED",
"新 Booking PostgreSQL 运行时尚未启用。");
}
return service.last24Hours();
}
}
@@ -16,6 +16,7 @@ import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.ResponseStatus;
import org.springframework.web.bind.annotation.RestController;
/** contracts-v1 processing run 查询和用户确认 API;没有 PMS/Opera 调用。 */
@@ -40,6 +41,14 @@ public class ReservationBookingProcessingRunController {
return service().get(runId);
}
/** 仅重放失败 run 的安全 contracts artifact;异步 worker 完成后由前端继续轮询。 */
@PostMapping(value = "/{runId}/retry", produces = MediaType.APPLICATION_JSON_VALUE)
@ResponseStatus(HttpStatus.ACCEPTED)
public BookingProcessingRunView retry(@PathVariable String runId) {
authorizationService.requirePermission(PlatformPermissionCode.RESERVATION_TASK_EDIT.name());
return service().retry(runId);
}
/** 冻结经校验的参数并写确认审计;不会创建 PMS/Opera 操作。 */
@PostMapping(
value = "/{runId}/confirmations",
@@ -0,0 +1,29 @@
package cn.nianxx.thhotel.workflows.reservation.booking.service.impl;
import static org.assertj.core.api.Assertions.assertThat;
import cn.nianxx.thhotel.workflows.reservation.common.dto.booking.BookingContracts;
import org.junit.jupiter.api.Test;
class BookingMessageSignalPolicyTest {
@Test
void shouldRecognizeCurrentMessageSignalsWithoutUsingHistory() {
assertThat(BookingMessageSignalPolicy.containsBusinessSignal("Update booking", "ordinary text")).isTrue();
assertThat(BookingMessageSignalPolicy.containsBusinessSignal(null, "ยกเลิกห้องพัก")).isTrue();
assertThat(BookingMessageSignalPolicy.containsBusinessSignal("FYI", "ordinary notice")).isFalse();
}
@Test
void shouldUseOnlyPersistedSafeAgentMaterialForRetry() {
BookingContracts.AgentInputMaterial material = new BookingContracts.AgentInputMaterial(
"current-evidence-001",
"Subject: booking update\nPlease amend one room.",
52,
false,
"booking-agent-current-message-v1");
assertThat(BookingMessageSignalPolicy.containsBusinessSignal(material)).isTrue();
assertThat(BookingMessageSignalPolicy.containsBusinessSignal((BookingContracts.AgentInputMaterial) null)).isFalse();
}
}
@@ -0,0 +1,42 @@
package cn.nianxx.thhotel.workflows.reservation.booking.service.impl;
import static org.assertj.core.api.Assertions.assertThat;
import cn.nianxx.thhotel.workflows.reservation.booking.common.dto.BookingOperationalMetricsSnapshot;
import cn.nianxx.thhotel.workflows.reservation.booking.common.result.BookingOperationalMetricsView;
import cn.nianxx.thhotel.workflows.reservation.booking.service.BookingOperationalMetricsService;
import org.junit.jupiter.api.Test;
class BookingOperationalMetricsServiceImplTest {
@Test
void shouldExposeSafeCountsAndRoundedRatesForTheFixedWindow() {
BookingOperationalMetricsService service = new BookingOperationalMetricsServiceImpl(hours -> {
assertThat(hours).isEqualTo(24);
return new BookingOperationalMetricsSnapshot(24, 8, 2, 3, 1, 4, 2, 5);
});
BookingOperationalMetricsView view = service.last24Hours();
assertThat(view.windowHours()).isEqualTo(24);
assertThat(view.totalRuns()).isEqualTo(8);
assertThat(view.failureRatePercent()).isEqualTo(25D);
assertThat(view.agentFallbackRatePercent()).isEqualTo(37.5D);
assertThat(view.riskRatePercent()).isEqualTo(12.5D);
assertThat(view.backlogRuns()).isEqualTo(4);
assertThat(view.retriedRuns()).isEqualTo(2);
assertThat(view.retryAttempts()).isEqualTo(5);
}
@Test
void shouldAvoidDivisionByZeroForAnEmptyWindow() {
BookingOperationalMetricsService service = new BookingOperationalMetricsServiceImpl(
hours -> new BookingOperationalMetricsSnapshot(24, 0, 0, 0, 0, 0, 0, 0));
BookingOperationalMetricsView view = service.last24Hours();
assertThat(view.failureRatePercent()).isZero();
assertThat(view.agentFallbackRatePercent()).isZero();
assertThat(view.riskRatePercent()).isZero();
}
}
@@ -0,0 +1,192 @@
package cn.nianxx.thhotel.workflows.reservation.booking.service.impl;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.eq;
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.workflows.reservation.booking.common.dto.BookingProcessingRunSnapshot;
import cn.nianxx.thhotel.workflows.reservation.booking.common.result.BookingProcessingRunView;
import cn.nianxx.thhotel.workflows.reservation.booking.repository.BookingProcessingRunRepository;
import cn.nianxx.thhotel.workflows.reservation.booking.service.BookingConfirmationProjectionFactory;
import cn.nianxx.thhotel.workflows.reservation.booking.service.BookingDecisionValidator;
import cn.nianxx.thhotel.workflows.reservation.common.dto.booking.BookingContracts;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.time.OffsetDateTime;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import org.junit.jupiter.api.Test;
class BookingProcessingRunServiceImplTest {
@Test
void shouldRetryOnlyFromSafePersistedContractsAndReturnAsyncContextReadyView() throws Exception {
BookingProcessingRunRepository repository = mock(BookingProcessingRunRepository.class);
BookingAgentContinuationDispatcher dispatcher = mock(BookingAgentContinuationDispatcher.class);
ObjectMapper objectMapper = new ObjectMapper().findAndRegisterModules();
BookingProcessingRunServiceImpl service = new BookingProcessingRunServiceImpl(
repository,
mock(BookingDecisionValidator.class),
mock(BookingConfirmationProjectionFactory.class),
dispatcher,
objectMapper);
BookingContracts.ContractMeta meta = meta();
BookingContracts.EvidenceReference evidence = evidence(meta);
BookingContracts.ParsedFactSet parsedFactSet = new BookingContracts.ParsedFactSet(
meta,
List.of(evidence),
BookingContracts.ParserApplicability.NOT_APPLICABLE,
null,
List.of(),
List.of(),
List.of("NO_FIXED_CHANNEL"),
List.of(),
new BookingContracts.AgentInputMaterial(
evidence.evidenceId(),
"Subject: Booking update\nPlease amend this reservation.",
54,
false,
"booking-agent-current-message-v1"));
BookingContracts.ContextPackage contextPackage = new BookingContracts.ContextPackage(
meta,
List.of(evidence),
"MESSAGE",
List.of(),
Map.of(),
Map.of(),
Map.of(),
Map.of(),
List.of());
BookingContracts.MaterialPackage materialPackage = materialPackage(meta, evidence);
when(repository.findByRunId("run-retry-001")).thenReturn(Optional.of(snapshot(
BookingContracts.ProcessingStatus.FAILED,
0)));
when(repository.findArtifactPayload("run-retry-001", "PARSED_FACT_SET"))
.thenReturn(Optional.of(objectMapper.writeValueAsString(parsedFactSet)));
when(repository.findArtifactPayload("run-retry-001", "CONTEXT_PACKAGE"))
.thenReturn(Optional.of(objectMapper.writeValueAsString(contextPackage)));
when(repository.findArtifactPayload("run-retry-001", "MATERIAL_PACKAGE"))
.thenReturn(Optional.of(objectMapper.writeValueAsString(materialPackage)));
when(repository.beginRetry("run-retry-001")).thenReturn(snapshot(
BookingContracts.ProcessingStatus.CONTEXT_READY,
1));
BookingProcessingRunView view = service.retry("run-retry-001");
assertThat(view.status()).isEqualTo("CONTEXT_READY");
assertThat(view.retryCount()).isEqualTo(1);
assertThat(view.confirmationProjection()).isNull();
verify(repository).beginRetry("run-retry-001");
verify(dispatcher).completeAsync(
eq("run-retry-001"),
eq(parsedFactSet),
eq(contextPackage),
eq(true),
eq(true));
}
@Test
void shouldRejectRetryWhenRunHasNotFailed() {
BookingProcessingRunRepository repository = mock(BookingProcessingRunRepository.class);
BookingProcessingRunServiceImpl service = new BookingProcessingRunServiceImpl(
repository,
mock(BookingDecisionValidator.class),
mock(BookingConfirmationProjectionFactory.class),
mock(BookingAgentContinuationDispatcher.class),
new ObjectMapper().findAndRegisterModules());
when(repository.findByRunId("run-retry-002")).thenReturn(Optional.of(snapshot(
BookingContracts.ProcessingStatus.AWAITING_CONFIRMATION,
0)));
assertThatThrownBy(() -> service.retry("run-retry-002"))
.hasMessageContaining("只有失败的 Booking 邮件可以重试");
verify(repository, never()).beginRetry("run-retry-002");
}
private BookingContracts.MaterialPackage materialPackage(
BookingContracts.ContractMeta meta,
BookingContracts.EvidenceReference evidence) {
BookingContracts.SourceMessageEnvelope envelope = new BookingContracts.SourceMessageEnvelope(
meta,
List.of(evidence),
BookingContracts.MessageOrigin.MANUAL_EML,
"external-retry-001",
"conversation-retry-001",
null,
"idempotency-retry-001",
OffsetDateTime.parse("2026-08-08T00:00:00Z"),
"Booking update",
evidence.evidenceId(),
null,
List.of(),
null);
return new BookingContracts.MaterialPackage(
meta,
List.of(evidence),
envelope,
List.of(
new BookingContracts.MaterialReference(
evidence.evidenceId(),
BookingContracts.EvidenceKind.CURRENT_BODY,
true,
Map.of("content_available", true)),
new BookingContracts.MaterialReference(
"attachment-evidence-retry-001",
BookingContracts.EvidenceKind.ATTACHMENT,
true,
Map.of("attachment_id", "attachment-retry-001"))),
List.of(),
List.of(),
List.of(),
BookingContracts.MaterialSelectionStatus.READY,
List.of());
}
private BookingProcessingRunSnapshot snapshot(
BookingContracts.ProcessingStatus status,
int retryCount) {
return new BookingProcessingRunSnapshot(
1L,
"run-retry-001",
2L,
"source-retry-001",
1,
status,
status == BookingContracts.ProcessingStatus.FAILED
? BookingContracts.ResultDisposition.RISK_NOTIFICATION
: null,
false,
retryCount,
null,
null,
OffsetDateTime.parse("2026-08-08T00:00:00Z"),
null);
}
private BookingContracts.ContractMeta meta() {
return new BookingContracts.ContractMeta(
BookingContracts.CONTRACT_VERSION,
"catalog-v1",
"parser-v1",
BookingContracts.AGENT_NOT_INVOKED,
"run-retry-001",
"source-retry-001",
0L);
}
private BookingContracts.EvidenceReference evidence(BookingContracts.ContractMeta meta) {
return new BookingContracts.EvidenceReference(
"current-evidence-retry-001",
BookingContracts.EvidenceKind.CURRENT_BODY,
meta.sourceMessageId(),
"current-body",
"a".repeat(64),
true);
}
}