diff --git a/client/src/i18n/locales/en-US.ts b/client/src/i18n/locales/en-US.ts index a45b386..4e8d9c2 100644 --- a/client/src/i18n/locales/en-US.ts +++ b/client/src/i18n/locales/en-US.ts @@ -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', diff --git a/client/src/i18n/locales/th-TH.ts b/client/src/i18n/locales/th-TH.ts index 749478c..202d46b 100644 --- a/client/src/i18n/locales/th-TH.ts +++ b/client/src/i18n/locales/th-TH.ts @@ -680,6 +680,10 @@ export default { refreshing: 'กำลังรีเฟรช…', source: 'อีเมลต้นทาง', revision: 'ฉบับแก้ไข', + retryCount: 'จำนวนครั้งที่ลองใหม่', + retry: 'ประมวลผลอีกครั้ง', + retrying: 'กำลังประมวลผลอีกครั้ง…', + retrySuccess: 'ส่งกลับเข้าคิวประมวลผลแล้ว และสถานะจะรีเฟรชอัตโนมัติ', disposition: 'ผลการจัดการ', status: 'สถานะปัจจุบัน', timeline: 'ความคืบหน้าการประมวลผล', diff --git a/client/src/i18n/locales/zh-CN.ts b/client/src/i18n/locales/zh-CN.ts index 473210e..0a3e8e0 100644 --- a/client/src/i18n/locales/zh-CN.ts +++ b/client/src/i18n/locales/zh-CN.ts @@ -679,6 +679,10 @@ export default { refreshing: '刷新中…', source: '来源邮件', revision: '修订', + retryCount: '重试次数', + retry: '重新处理', + retrying: '正在重新处理…', + retrySuccess: '已重新进入处理队列,将自动刷新状态。', disposition: '处理结果', status: '当前状态', timeline: '处理进度', diff --git a/client/src/services/reservationService.ts b/client/src/services/reservationService.ts index 0e20dee..89698b2 100644 --- a/client/src/services/reservationService.ts +++ b/client/src/services/reservationService.ts @@ -63,6 +63,13 @@ export async function fetchBookingProcessingRun(runId: string): Promise { + return sendJson( + `/api/reservation/booking-processing-runs/${encodeURIComponent(runId)}/retry`, + 'POST', + ) +} + export async function confirmBookingProcessingRun( runId: string, request: BookingConfirmationRequest, diff --git a/client/src/tests/reservationBookingProcessingRunView.spec.ts b/client/src/tests/reservationBookingProcessingRunView.spec.ts index a016ed6..c12ab40 100644 --- a/client/src/tests/reservationBookingProcessingRunView.spec.ts +++ b/client/src/tests/reservationBookingProcessingRunView.spec.ts @@ -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() { diff --git a/client/src/views/reservation/ReservationBookingProcessingRunView.vue b/client/src/views/reservation/ReservationBookingProcessingRunView.vue index 72a049f..7d87037 100644 --- a/client/src/views/reservation/ReservationBookingProcessingRunView.vue +++ b/client/src/views/reservation/ReservationBookingProcessingRunView.vue @@ -64,6 +64,10 @@
{{ t('bookingProcessingRun.revision') }}
{{ run.revision_number }}
+
+
{{ t('bookingProcessingRun.retryCount') }}
+
{{ run.retry_count }}
+
{{ t('bookingProcessingRun.disposition') }}
{{ run.result_disposition ?? '—' }}
@@ -105,6 +109,15 @@ > {{ run.error_code ?? 'BOOKING_PROCESSING_ERROR' }} · {{ run.error_summary }}

+
(null) const loading = ref(false) const confirming = ref(false) +const retrying = ref(false) const errorMessage = ref('') const successMessage = ref('') const draftValues = ref>>({}) @@ -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 { } } +async function retryRun(): Promise { + 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() for (const field of item.display_fields) { diff --git a/docs/project/operations/booking-v01-security-checkpoint.md b/docs/project/operations/booking-v01-security-checkpoint.md index f733c34..45e84c9 100644 --- a/docs/project/operations/booking-v01-security-checkpoint.md +++ b/docs/project/operations/booking-v01-security-checkpoint.md @@ -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 路径继续,不能以重复重试替代环境修复。 diff --git a/docs/project/requirements/booking-email-architecture-v0.4.md b/docs/project/requirements/booking-email-architecture-v0.4.md index c2838f0..054112a 100644 --- a/docs/project/requirements/booking-email-architecture-v0.4.md +++ b/docs/project/requirements/booking-email-architecture-v0.4.md @@ -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。 diff --git a/docs/project/requirements/booking-email-contracts-v1.md b/docs/project/requirements/booking-email-contracts-v1.md index 513bdca..077328c 100644 --- a/docs/project/requirements/booking-email-contracts-v1.md +++ b/docs/project/requirements/booking-email-contracts-v1.md @@ -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. 兼容与测试要求 diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/common/dto/BookingOperationalMetricsSnapshot.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/common/dto/BookingOperationalMetricsSnapshot.java new file mode 100644 index 0000000..0235e82 --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/common/dto/BookingOperationalMetricsSnapshot.java @@ -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 运行指标不能为负数"); + } + } +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/common/result/BookingOperationalMetricsView.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/common/result/BookingOperationalMetricsView.java new file mode 100644 index 0000000..bdd214b --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/common/result/BookingOperationalMetricsView.java @@ -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) { +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingOperationalMetricsRepository.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingOperationalMetricsRepository.java new file mode 100644 index 0000000..bd1cc2f --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingOperationalMetricsRepository.java @@ -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); +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresOperationalMetricsRepository.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresOperationalMetricsRepository.java new file mode 100644 index 0000000..0199c91 --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresOperationalMetricsRepository.java @@ -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); + } +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresProcessingRunRepository.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresProcessingRunRepository.java index 3b9b39a..83ed89a 100644 --- a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresProcessingRunRepository.java +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingPostgresProcessingRunRepository.java @@ -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") diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingProcessingRunRepository.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingProcessingRunRepository.java index f6d3189..1235ba4 100644 --- a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingProcessingRunRepository.java +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/repository/BookingProcessingRunRepository.java @@ -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 findArtifactPayload(String runId, String artifactType); diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingOperationalMetricsService.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingOperationalMetricsService.java new file mode 100644 index 0000000..f2c705d --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingOperationalMetricsService.java @@ -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(); +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingProcessingRunService.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingProcessingRunService.java index a824b30..5905784 100644 --- a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingProcessingRunService.java +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/BookingProcessingRunService.java @@ -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); } diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingMessageSignalPolicy.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingMessageSignalPolicy.java new file mode 100644 index 0000000..7ec50bb --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingMessageSignalPolicy.java @@ -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 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()); + } +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingOperationalMetricsServiceImpl.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingOperationalMetricsServiceImpl.java new file mode 100644 index 0000000..12bdf17 --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingOperationalMetricsServiceImpl.java @@ -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; + } +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingPostgresMessageOrchestrator.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingPostgresMessageOrchestrator.java index 651102c..ccca603 100644 --- a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingPostgresMessageOrchestrator.java +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingPostgresMessageOrchestrator.java @@ -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 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"); diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImpl.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImpl.java index 84c970e..560c413 100644 --- a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImpl.java +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImpl.java @@ -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 artifact(String runId, String artifactType, Class type) { String payload = processingRunRepository.findArtifactPayload(runId, artifactType) .orElseThrow(() -> workflowError( diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingOperationalMetricsController.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingOperationalMetricsController.java new file mode 100644 index 0000000..816fb87 --- /dev/null +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingOperationalMetricsController.java @@ -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 metricsServiceProvider; + private final FrontendAuthorizationService authorizationService; + + public ReservationBookingOperationalMetricsController( + ObjectProvider 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(); + } +} diff --git a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingProcessingRunController.java b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingProcessingRunController.java index 648ebe9..fc0369f 100644 --- a/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingProcessingRunController.java +++ b/server/src/main/java/cn/nianxx/thhotel/workflows/reservation/control/ReservationBookingProcessingRunController.java @@ -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", diff --git a/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingMessageSignalPolicyTest.java b/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingMessageSignalPolicyTest.java new file mode 100644 index 0000000..d632a6f --- /dev/null +++ b/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingMessageSignalPolicyTest.java @@ -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(); + } +} diff --git a/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingOperationalMetricsServiceImplTest.java b/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingOperationalMetricsServiceImplTest.java new file mode 100644 index 0000000..edda4a6 --- /dev/null +++ b/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingOperationalMetricsServiceImplTest.java @@ -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(); + } +} diff --git a/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImplTest.java b/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImplTest.java new file mode 100644 index 0000000..6569a96 --- /dev/null +++ b/server/src/test/java/cn/nianxx/thhotel/workflows/reservation/booking/service/impl/BookingProcessingRunServiceImplTest.java @@ -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); + } +}