From 40230ce33d6df77b1c5eb5b52a4556f7feb5e1dc Mon Sep 17 00:00:00 2001 From: andy Date: Sun, 20 Sep 2026 14:59:19 +0800 Subject: [PATCH] =?UTF-8?q?=E5=9B=9E=E6=89=A7=E5=A2=9E=E5=8A=A0=E5=86=85?= =?UTF-8?q?=E5=AE=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...0260920-team-progress-receipts-01a0858b.md | 33 ++ LianSyn-platform/app.js | 2 +- agent设计规范/agentbus-reply-contract.md | 15 + agent设计规范/business-adaptation-registry.md | 4 + .../team-progress-receipts-01a0858b/README.md | 33 ++ .../migrations/026_team_progress.sql | 32 ++ control-plane/src/agentbus-delivery.ts | 10 + control-plane/src/agentbus.ts | 5 +- control-plane/src/db.ts | 2 +- control-plane/src/reply-contract.ts | 4 +- control-plane/src/task-service.ts | 23 +- control-plane/src/team-progress.ts | 291 ++++++++++++++++++ control-plane/test/arrangement-ledger.test.ts | 3 + control-plane/test/control-plane.test.ts | 2 +- .../test/support/team-progress-memory.ts | 30 ++ control-plane/test/team-progress.test.ts | 192 ++++++++++++ tools/arrangement-ledger-postgres-smoke.mjs | 57 +++- 17 files changed, 729 insertions(+), 9 deletions(-) create mode 100644 .project-docs/30-worklog/tasks/20260920-team-progress-receipts-01a0858b.md create mode 100644 archive/evidence/2026-09-20/team-progress-receipts-01a0858b/README.md create mode 100644 control-plane/migrations/026_team_progress.sql create mode 100644 control-plane/src/team-progress.ts create mode 100644 control-plane/test/support/team-progress-memory.ts create mode 100644 control-plane/test/team-progress.test.ts diff --git a/.project-docs/30-worklog/tasks/20260920-team-progress-receipts-01a0858b.md b/.project-docs/30-worklog/tasks/20260920-team-progress-receipts-01a0858b.md new file mode 100644 index 0000000..0e18fe7 --- /dev/null +++ b/.project-docs/30-worklog/tasks/20260920-team-progress-receipts-01a0858b.md @@ -0,0 +1,33 @@ +# 全业务团队进度与安排明细回执 + +## Scope / 授权 + +- 用户批准讨论方案 B:各业务成功回执增加团队办理进度、当前有效安排明细;名单只显示人数,不输出旅客个人资料。 +- 按当前账号和完整团号隔离;散拼子单关联所属母团,批量按团分别展示。未能确认团号时不能推断全部未安排。 +- 当前有效安排逐条显示;变更/清除后保留历史。名单使用最近成功结果,不累计重复导入;修改、导出显示最近记录。取消/恢复不清空安排。 +- 强制删除任务保留这些独立业务记录;历史回执冻结生成时快照。 +- 不部署、不重启、不执行真实 ERP 写入或外部发送。单线程实施。 + +## Worktree Gate + +- Branch: main;Worktree: /Users/andy/IdeaProjects/LWLT-AIBOT。 +- Base: ea9c4b34c6ebc0c41b89d05d4d02e1b306c60135;仅一个 worktree。 +- 既有未跟踪 .idea/、dist/ltjt-order-assistant-0.5.188.zip 保留;源码无未提交改动。 +- Overlap: Clear;Feature 模式,只写本任务项目记忆。 + +## Progress + +- 已完成。新增 `026_team_progress.sql` 和 `team-progress.ts`,加密保存按组织/账号/完整团号隔离的独立业务事件;精确执行编号补录仍存在的历史凭据。 +- 统一成功回执追加下单/建计划、具体子单、最近名单人数、最近修改、文件类型与导出时间、取消/恢复及五类当前有效安排明细。子单与母团严格区分,身份冲突/无法确认团号时不推断未安排。 +- 名单不累计重复导入,不保留游客正文;搜索条件与实际名称区分。批量任务按明确对应关系归团,中断批次只保存各条独立成功凭据,父任务结果不变。 +- 任务结果、业务记录和回执快照同事务写入,幂等;强制删除前补录并标记独立历史,导出事实保留但文件附件仍删除;旧回执冻结。 +- AgentBus 超过既有 20,000 字符上限明确提示平台查看完整回执。平台仅调整强制删除说明;现有指令、ERP 写入逻辑和插件版本不变。 +- 已同步业务登记与回执契约。未部署、重启、访问 ERP、重放真实任务或外发消息。 + +## Validation / Promotion Candidates + +- `check`、`build` 通过;控制面 293/293,旧功能 431/431。直接工作树 `check:repo` 的两项失败仅为既有 `.idea`、`.DS_Store`、被忽略的历史日志;保留原文件,当前源码独立快照检查 10/10。 +- 独立临时 PostgreSQL 执行至迁移 026(共 25 个迁移文件):实际 TaskService 入库、账号/组织隔离、幂等、事务回滚、历史补录、删除后保留、冻结回执、名单最近值、母团恢复不恢复子单、中断批次独立成功记录均通过。测试无真实 ERP。 +- 证据:[验证记录](../../../archive/evidence/2026-09-20/team-progress-receipts-01a0858b/README.md)。 +- 上线项:执行迁移 026;正常 `npm run dev` / `npm start` 已自动迁移。插件保持 0.5.188,无需重新安装。升级前已经删除且未进入独立台账的业务不可恢复。 +- Promotion:后续 Integration 将“全业务账号内团队进度、独立事件保留、不可变回执快照”纳入 `20-architecture/system-overview.md` 和 `30-worklog/current-state.md`,并记录迁移 026。产品语义已获用户批准,无待决产品冲突。 diff --git a/LianSyn-platform/app.js b/LianSyn-platform/app.js index 4320e3f..96a2881 100644 --- a/LianSyn-platform/app.js +++ b/LianSyn-platform/app.js @@ -5265,7 +5265,7 @@ function confirmTaskHardDelete(tasks) { ? `任务 ${selectedTasks[0].task_id}` : `所选 ${selectedTasks.length} 个任务`; return window.confirm( - `确认永久强制删除${scope}?\n\n任务、原始输入、生命周期、执行记录、附件和任务回执都会被物理删除,无法恢复。独立保存的团队安排明细、变更历史及成功凭据会保留,安排状态不受影响。此操作不受“正在处理”或“等待 ERP 执行”状态限制。\n\n如果 ERP 已经开始写入,删除平台记录不会撤销 ERP 中已经发生的操作;系统只会向任务所属账号的在线插件发送停止与清理指令。已经由外部 Webhook 受理并可能到达组长微信群的任务摘要也无法撤回。` + `确认永久强制删除${scope}?\n\n任务、原始输入、生命周期、执行记录、附件和任务回执都会被物理删除,无法恢复。独立保存的团队办理进度、安排明细、变更历史及成功凭据会保留;导出文件记录保留,文件附件仍会删除。此操作不受“正在处理”或“等待 ERP 执行”状态限制。\n\n如果 ERP 已经开始写入,删除平台记录不会撤销 ERP 中已经发生的操作;系统只会向任务所属账号的在线插件发送停止与清理指令。已经由外部 Webhook 受理并可能到达组长微信群的任务摘要也无法撤回。` ); } diff --git a/agent设计规范/agentbus-reply-contract.md b/agent设计规范/agentbus-reply-contract.md index 6c109f0..4cc4fef 100644 --- a/agent设计规范/agentbus-reply-contract.md +++ b/agent设计规范/agentbus-reply-contract.md @@ -60,6 +60,7 @@ | `team_order_batch_create` | 批量下单成功 | 全部团号,不只第一条 | | `shared_plan_create` | 散拼计划创建成功 | 全部母团团号 | | `shared_child_order_create` | 新子单成功 | 子单号;有母团号时可同时返回母团号 | +| `shared_child_order_batch_create` | 批量新增子单成功 | 每个已确认子单号及对应母团号 | | `arrangement_guide` | 导游安排成功/清除成功 | 团号、导游或清除结果 | | `arrangement_vehicle` | 用车安排成功/清除成功 | 团号、用车日期、数量、车辆/项目 | | `arrangement_hotel` | 客户酒店模板或清除成功 | 团号、入住/退房、晚数、房型/房数 | @@ -73,6 +74,20 @@ 成功回执中的字段必须来自已验证的 ERP 回执或写后回查;不能使用“已提交并完成回查”这种技术过程作为业务正文。 +## 全业务团队进度与安排明细 + +各业务成功回执在本次操作结果之后追加“本系统团队办理进度(当前账号)”。范围严格限定为同一组织、任务归属账号及完整团号的本系统成功记录,不读取 ERP 外部操作补全进度,不跨账号合并。团号无法确认或身份凭据冲突时只提示暂无法关联,不推断五类安排全部未安排。 + +- 独立团展示下单、最近名单导入人数与时间、最近修改、各类文件最近导出和取消/恢复状态。名单按最近成功导入统计,不累加重复导入,不输出游客姓名、证件、电话或名单原文。 +- 散拼母团展示建计划、系统已记录子单数量和母团修改/文件;各子单分别展示新增、名单、修改、文件及状态。已记录子单不代表母团全部子单;母团恢复不代表子单已恢复。只有明确的子单—母团成功凭据或当前账号唯一历史关联才能归团。 +- 批量结果按团展示,日期须来自每条对应的成功凭据,不按数组顺序猜测。中断批次只保留有独立完成凭据的条目,父任务仍使用原有失败/待回查回执。 +- 导游、用车、酒店、大交通、其他/备案显示当前有效的每一条安排及已记录资源、项目、日期、数量、电话、备注等业务明细。变更更新当前值,清除移除对应有效安排,历史继续保留;缺失字段不编造。未安排只表示本账号无本系统有效记录,已安排表示至少一条,不代表全行程已安排齐全。 +- 文件只汇总导出类型和时间。此后有相关业务变更时提示可能需要重新导出,不删除旧记录。成功但类型缺失时标注未记录;强制删除任务保留导出事实,文件附件仍按原流程删除。 + +控制面将办理事件与安排明细加密保存为独立记录,不随任务强制删除级联消失。历史可证明的完成记录按原执行编号补录;已经删除且没有独立记录的数据不能凭空恢复。每次成功回执保存当时快照,后续操作不修改旧回执。 + +平台保存完整回执。AgentBus 最终文本超过现有 20,000 字符限制时,按行截取并明确提示“完整内容请在平台任务回执中查看”,不静默截断;阶段消息和附件交付规则不变。 + ## 酒店安排客户模板 酒店新增成功且执行结果包含完整酒店字段时,使用以下格式。变量全部由真实回查结果填充;英文和中文日期使用同一组实际日期。 diff --git a/agent设计规范/business-adaptation-registry.md b/agent设计规范/business-adaptation-registry.md index 5e7912b..52d2c08 100644 --- a/agent设计规范/business-adaptation-registry.md +++ b/agent设计规范/business-adaptation-registry.md @@ -30,6 +30,10 @@ 测试对象的 `order_delete` 回收能力仍保留在 Schema、插件和生命周期测试夹具中,但不属于生产运营业务,不登记为用户输入入口。 +## 跨业务成功回执 + +各业务成功回执追加当前账号、完整团号范围内的团队办理进度及五类当前有效安排明细;散拼子单分别记录,名单只展示最近成功导入人数,文件保留导出事实。该汇总只来源于本系统成功记录,强制删除任务不删除独立业务历史,旧回执保持生成时快照。具体规则见 [AgentBus 用户回复契约](agentbus-reply-contract.md#全业务团队进度与安排明细);实现位于 `control-plane/src/team-progress.ts`、`arrangement-ledger.ts` 与 `task-service.ts`,不改变 ERP 执行或用户指令。 + ## 列表检索与派发结果边界 `0.5.185` 起,插件在 ERP 只读目标查询发生主框架移除时,最多等待 8 秒确认当前主文档稳定后重新查询一次;仅同一写入前运行任务可重查,保存及业务候选不唯一不重试。异常保留实际阶段和当时自动恢复诊断,不能仅凭稍后心跳结果推断刷新原因。见[当前任务](../.project-docs/30-worklog/tasks/20260915-resolution-navigation-01a0858b.md)。 diff --git a/archive/evidence/2026-09-20/team-progress-receipts-01a0858b/README.md b/archive/evidence/2026-09-20/team-progress-receipts-01a0858b/README.md new file mode 100644 index 0000000..eaf4cd7 --- /dev/null +++ b/archive/evidence/2026-09-20/team-progress-receipts-01a0858b/README.md @@ -0,0 +1,33 @@ +# 团队办理进度与安排明细:本地验证 + +日期:2026-09-20。此目录为不可变验证证据,不是当前业务规则;当前规则见 [回执契约](../../../../agent设计规范/agentbus-reply-contract.md)。仅使用合成测试数据,未连接 ERP、生产数据库或外部发送渠道。 + +## 检查结果 + +| 检查 | 结果 | +|---|---| +| `node --run check` | 通过 | +| `node --run build` | 通过 | +| `node --run test:control-plane` | 293/293 通过 | +| `node --run test:legacy` | 431/431 通过 | +| `node --run check:repo`(当前源码独立快照) | 10/10 通过 | +| `git diff --check` | 通过 | +| `tools/arrangement-ledger-postgres-smoke.mjs` | 全部通过,迁移执行至 026(25 个 SQL 文件) | + +当前工作树仓库治理检查 8/10;其两项失败来自预先存在的 `.idea/`、根目录及插件目录 `.DS_Store`、2026-09-14 历史目录中的被忽略日志。未移动、删除或认领这些文件。验证快照从当前 Git 跟踪文件及本次新增文件复制,包含现有版本化扩展 ZIP,排除 IDE 元数据、忽略产物及秘密配置,并复用本地依赖。旧功能测试首次受沙箱本机 HTTP 监听限制,获工具审批后仅在本机测试环境补跑通过。 + +## 实际数据库覆盖 + +使用独立无数据卷 PostgreSQL 16 临时容器和 `ledger_test` 数据库。脚本新建随机 schema,结束时删除 schema;没有使用应用 `.env` 或应用数据库。通过真实迁移和 TaskService 事务接口验证: + +- 成功结果入库、公开任务与 AgentBus 共用完整回执;重复回传幂等。 +- 组织/账号隔离,管理员归属与跨组织外键阻断。 +- 非安排操作也生成团队进度;具体子单单独统计;重复名单导入使用最近人数,不出现名单内容。 +- 后续操作不修改旧回执;母团恢复不恢复仍取消的子单。 +- 强制删除前补录已归档历史;删除任务及执行记录后独立业务与安排明细仍在。 +- 安排精确清除只移除对应行;异常时业务记录与任务结果共同回滚。 +- 中断批次仅记录有独立成功凭据的子单,父任务继续待回查,不生成成功回执。 + +## 验证边界 + +未执行真实 ERP 写入、部署、业务服务重启或插件重载。插件与发布包未修改。已有且可证明的历史记录可以补录;升级前已删除且无独立记录的数据不能恢复。真实长消息的渠道显示仍由 AgentBus 渠道限制决定,平台保留完整文本,程序在既有长度上限处明确提示。 diff --git a/control-plane/migrations/026_team_progress.sql b/control-plane/migrations/026_team_progress.sql new file mode 100644 index 0000000..831271f --- /dev/null +++ b/control-plane/migrations/026_team_progress.sql @@ -0,0 +1,32 @@ +-- Durable business history survives task/attempt/attachment deletion. +CREATE TABLE team_progress_owners ( + organization_id uuid NOT NULL REFERENCES organizations(id), + owner_user_id uuid NOT NULL, + backfilled_at timestamptz, + PRIMARY KEY (organization_id, owner_user_id), + FOREIGN KEY (organization_id, owner_user_id) REFERENCES users(organization_id, id) +); +CREATE TABLE team_progress_events ( + id bigserial PRIMARY KEY, + organization_id uuid NOT NULL REFERENCES organizations(id), + owner_user_id uuid NOT NULL, + group_key text NOT NULL, + subject_key text NOT NULL, + subject_kind text NOT NULL CHECK (subject_kind IN ('team', 'child')), + event_key text NOT NULL, + source_task_id text NOT NULL, + execution_id text NOT NULL, + effective_at timestamptz NOT NULL, + recorded_at timestamptz NOT NULL, + source_task_deleted_at timestamptz, + detail_ciphertext text NOT NULL, + UNIQUE (organization_id, owner_user_id, execution_id, event_key), + FOREIGN KEY (organization_id, owner_user_id) REFERENCES users(organization_id, id) +); +CREATE INDEX team_progress_group_idx ON team_progress_events (organization_id, owner_user_id, group_key, effective_at, id); +CREATE INDEX team_progress_subject_idx ON team_progress_events (organization_id, owner_user_id, subject_key); +CREATE INDEX team_progress_task_idx ON team_progress_events (organization_id, owner_user_id, source_task_id); +CREATE TRIGGER team_progress_owner_not_admin BEFORE INSERT OR UPDATE OF organization_id, owner_user_id ON team_progress_owners +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal('owner_user_id', 'team_progress_owner_not_admin'); +CREATE TRIGGER team_progress_event_not_admin BEFORE INSERT OR UPDATE OF organization_id, owner_user_id ON team_progress_events +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal('owner_user_id', 'team_progress_event_not_admin'); diff --git a/control-plane/src/agentbus-delivery.ts b/control-plane/src/agentbus-delivery.ts index c24f071..735181a 100644 --- a/control-plane/src/agentbus-delivery.ts +++ b/control-plane/src/agentbus-delivery.ts @@ -3,6 +3,16 @@ export const AGENTBUS_ROSTER_WAITING_TEXT = '已识别名单业务,等待 .xls export const AGENTBUS_ROSTER_ATTACHMENT_RECEIVED_TEXT = '名单附件已收到,正在校验并处理。'; export const AGENTBUS_FINAL_REPLY_ELIGIBLE_FIELD = '_final_reply_eligible'; +// Preserve the complete receipt on the platform; never silently cut a long +// multi-team/detail receipt in the middle of a field on the channel. +export function boundedAgentBusReplyText(value: string): string { + if (value.length <= 20_000) return value; + const suffix = '\n\n回执明细较长,此处仅展示部分;完整内容请在平台任务回执中查看。'; + const prefix = value.slice(0, 20_000 - suffix.length); + const line = prefix.lastIndexOf('\n'); + return (line > 0 ? prefix.slice(0, line) : prefix.replace(/[\uD800-\uDBFF]$/, '')) + suffix; +} + export interface AgentBusAcceptedDeliveryOptions { text?: string; finalReplyEligible?: boolean; diff --git a/control-plane/src/agentbus.ts b/control-plane/src/agentbus.ts index 9614ba3..10031ca 100644 --- a/control-plane/src/agentbus.ts +++ b/control-plane/src/agentbus.ts @@ -21,6 +21,7 @@ import { } from './input-attachment.js'; import { diagnosticDurationMs, diagnosticError } from './diagnostics.js'; import { + boundedAgentBusReplyText, AGENTBUS_GENERIC_ACCEPTED_TEXT, AGENTBUS_ROSTER_ATTACHMENT_RECEIVED_TEXT, publicAgentBusDeliveryPayload @@ -482,7 +483,7 @@ export function createTaskResultFrame( payload: { event: 'task.result', status, - text: resultText.slice(0, 20_000), + text: boundedAgentBusReplyText(resultText), ...(attachments.length ? { attachments } : {}), ...(importantMessage ? { important_message: channelImportantMessage(importantMessage) } : {}) } @@ -1149,7 +1150,7 @@ export class AgentBusListener { event: 'task.result', status: taskResultStatus(task), task_id: task.task_id, - text: taskResultText(task).slice(0, 20_000), + text: boundedAgentBusReplyText(taskResultText(task)), ...(attachmentRefs.length ? { attachments: attachmentRefs } : {}), ...(importantMessage ? { important_message: importantMessage } : {}) }; diff --git a/control-plane/src/db.ts b/control-plane/src/db.ts index d92f372..bed600e 100644 --- a/control-plane/src/db.ts +++ b/control-plane/src/db.ts @@ -5,7 +5,7 @@ import { writeEmergencyDiagnostic } from './diagnostics.js'; const { Pool } = pg; let pool: pg.Pool | null = null; -export const REQUIRED_SCHEMA_VERSION = '025_arrangement_ledger'; +export const REQUIRED_SCHEMA_VERSION = '026_team_progress'; export interface DatabaseReadiness { ready: boolean; diff --git a/control-plane/src/reply-contract.ts b/control-plane/src/reply-contract.ts index d9c62a9..74d4349 100644 --- a/control-plane/src/reply-contract.ts +++ b/control-plane/src/reply-contract.ts @@ -1,5 +1,6 @@ import { normalizeTaskArtifactFileName } from './artifact-name.js'; import { arrangementSummaryText, isArrangementAction } from './arrangement-ledger.js'; +import { teamProgressText } from './team-progress.js'; export const AGENTBUS_REPLY_CONTRACT_VERSION = 'agentbus-user-reply-v1'; @@ -563,7 +564,8 @@ export function buildCustomerSuccessReceipt( const reply = buildCustomerSuccessReply(baseReceipt, resultValue, operationValue); if (!reply) return baseReceipt; const action = operationAction(operationValue, object(resultValue), resultContext(object(resultValue))); - const summaryText = isArrangementAction(action) ? arrangementSummaryText(baseReceipt.arrangement_summary) : ''; + const summaryText = teamProgressText(baseReceipt.team_progress) + || (isArrangementAction(action) ? arrangementSummaryText(baseReceipt.arrangement_summary) : ''); return { ...baseReceipt, reply: { diff --git a/control-plane/src/task-service.ts b/control-plane/src/task-service.ts index 590a78e..a2be5cb 100644 --- a/control-plane/src/task-service.ts +++ b/control-plane/src/task-service.ts @@ -32,6 +32,7 @@ import { import { convertDocumentToPdf, convertDocumentToXlsx } from './document-converter.js'; import { sanitizeOperationTiming } from './operation-timing.js'; import { ArrangementLedger, arrangementEventFromSuccess, arrangementSummary } from './arrangement-ledger.js'; +import { TeamProgressLedger, progressEventsFromSuccess, partialBatchProgressEvents } from './team-progress.js'; import type { TaskInputAttachmentInput } from './input-attachment.js'; import { diagnosticMetadataKeys, @@ -6618,7 +6619,10 @@ export class TaskService { const operation = decryptedJson(this.config, row.operation_ciphertext) || row.operation || null; const baseSuccessReceipt = status === 'completed' ? successReceiptFromResult(resultObject) : null; // This snapshot is exclusively derived by the server, never by a plugin. - if (baseSuccessReceipt) delete baseSuccessReceipt.arrangement_summary; + if (baseSuccessReceipt) { + delete baseSuccessReceipt.arrangement_summary; + delete baseSuccessReceipt.team_progress; + } let successReceipt = baseSuccessReceipt ? buildCustomerSuccessReceipt(baseSuccessReceipt, resultObject, operation) : null; @@ -6670,6 +6674,20 @@ export class TaskService { arrangement_summary: arrangementSummary(history, arrangementEvent.group_number, arrangementEvent.recorded_at) }, resultObject, operation); } + const progressInput = { taskId, executionId, operation, result: resultObject, + receipt: baseSuccessReceipt, effectiveAt: attemptRow.started_at, recordedAt: new Date().toISOString() }; + const progressEvents = [...progressEventsFromSuccess(progressInput), ...partialBatchProgressEvents(progressInput)]; + if (successReceipt || progressEvents.length) { + const ledger = new ArrangementLedger(this.config); + await ledger.prepare(client, context, successReceiptFromResult); + const progress = new TeamProgressLedger(this.config); + await progress.prepare(client, context, successReceiptFromResult); + for (const event of progressEvents) await progress.insert(client, context, event); + if (successReceipt) successReceipt = buildCustomerSuccessReceipt({ ...successReceipt, + team_progress: await progress.snapshot(client, context, progressEvents, progressInput.recordedAt, + groupNumber => ledger.history(client, context, groupNumber)) + }, resultObject, operation); + } const updated = await client.query( `UPDATE tasks SET execution_result_ciphertext = $1, execution_result = $2, @@ -7300,6 +7318,9 @@ export class TaskService { const ledger = new ArrangementLedger(this.config); await ledger.prepare(client, context, successReceiptFromResult); await ledger.markTasksDeleted(client, context, normalizedTaskIds); + const progress = new TeamProgressLedger(this.config); + await progress.prepare(client, context, successReceiptFromResult); + await progress.markTasksDeleted(client, context, normalizedTaskIds); // Hard delete intentionally has no status or handoff-state gate. The row // locks settle concurrent transitions; existing task foreign keys then diff --git a/control-plane/src/team-progress.ts b/control-plane/src/team-progress.ts new file mode 100644 index 0000000..ab13728 --- /dev/null +++ b/control-plane/src/team-progress.ts @@ -0,0 +1,291 @@ +import type pg from 'pg'; +import type { AppConfig } from './config.js'; +import { decryptText, encryptText, sha256Text } from './crypto.js'; +import { ARRANGEMENT_LABELS, arrangementGroupNumber, currentArrangements, + type ArrangementEvent, type ArrangementScope } from './arrangement-ledger.js'; + +type Json = Record; +type Client = Pick; +const obj = (v: unknown): Json => v && typeof v === 'object' && !Array.isArray(v) ? v as Json : {}; +const txt = (v: unknown): string => typeof v === 'string' || typeof v === 'number' ? String(v).trim() : ''; +const first = (...v: unknown[]): string => v.map(txt).find(Boolean) || ''; +const list = (v: unknown): Json[] => Array.isArray(v) ? v.map(obj) : []; +const strings = (...values: unknown[]): string[] => [...new Set(values.flatMap(v => Array.isArray(v) ? v.map(txt) : [txt(v)]).filter(Boolean))]; +const date = (v: unknown): string => { const d = v instanceof Date ? v : new Date(txt(v)); return Number.isFinite(d.getTime()) ? d.toISOString() : ''; }; +const group = arrangementGroupNumber; +// Lifecycle receipts may fall back to internal tid/ddid values. Those are not +// group numbers and must never make unrelated business records look connected. +const businessGroup = (v: unknown): string => { + const value = group(v); + return /^\d+$/.test(value) || /^D\d+$/.test(value) ? '' : value; +}; +const ACTIONS = new Set(['team_order_create', 'team_order_batch_create', 'shared_plan_create', 'shared_child_order_create', + 'shared_child_order_batch_create', 'passenger_list_import', 'order_update_independent', 'order_update_shared_plan', + 'order_update_shared_child', 'order_cancel', 'order_restore', 'confirmation_export', ...Object.keys(ARRANGEMENT_LABELS)]); +const FILES: Record = { 'xingyou-confirm': '星游客户确认单', 'liantai-confirm': '联泰客户确认单', + 'job-order': '团队申请书', 'visitor-list': '游客名单', 'guide-confirm': '导游通知书', 'hotel-preorder': '酒店预订单', + 'transport-preorder': '团体票申请书', 'filing-current': '备案确认书', 'filing-history': '历史备案确认书', 'pickup-sign': '接机牌' }; +const UPDATE_LABELS: Record = { 'rooms.SGL': '单间', 'rooms.TWN': '标间', 'rooms.DBL': '大床房', 'rooms.TRP': '三人间', + twin_room_count: '标间', planned_capacity: '计划收客数', lodging_note: '订房说明', 'pax.adult': '成人', + 'pax.child_bed': '占床儿童', 'pax.child_no_bed': '不占床儿童', 'pax.leader': '领队人数', departure_date: '出发日期', remark: '备注' }; + +export interface ProgressEvent { + version: 1; group_number: string; subject_kind: 'team' | 'child'; subject_number: string; + action: string; source_task_id: string; execution_id: string; effective_at: string; recorded_at: string; + source_task_deleted_at?: string | null; details: Json; +} +export interface ProgressSnapshot { + version: 1; scope: 'current_account'; recorded_at: string; + groups: Array<{ group_number: string; events: ProgressEvent[]; arrangements: Array<{ action: string; details: Json }> }>; + unlinked_subjects: string[]; +} +export interface ProgressInput { + taskId: string; executionId: string; effectiveAt: unknown; recordedAt: unknown; + operation: unknown; result: unknown; receipt: unknown; +} + +// This is a projection of accepted success evidence, never a copy of operations, +// passenger TSV, attachment contents, URLs or arbitrary plugin metadata. +export function progressEventsFromSuccess(input: ProgressInput): ProgressEvent[] { + const op = obj(input.operation), data = obj(op.data), result = obj(input.result), report = obj(result.report), receipt = obj(input.receipt); + const action = txt(op.action), refs = obj(report.resolved_refs || result.resolved_refs); + if (!ACTIONS.has(action) || result.status !== 'completed' || !Object.keys(receipt).length || receipt.success === false + || result.uncertain === true || result.manual_review_required === true || report.manual_review_required === true + || [result.blockers, report.blockers].some(v => Array.isArray(v) && v.length)) return []; + const recordedAt = date(input.recordedAt), effectiveAt = date(input.effectiveAt) || recordedAt; + if (!recordedAt || !input.taskId || !input.executionId) return []; + const context = obj(result.business_reply_context || report.business_reply_context || report.business_receipt); + const details: Json = {}; + // Search keywords are displayed as input conditions, never as verified names. + for (const [key, value] of Object.entries({ customer: context.customer_name, product: context.product_name, + customer_condition: obj(data.customer).name || obj(data.customer).keyword, + product_condition: obj(data.product).name || obj(data.product).keyword })) if (txt(value)) details[key] = txt(value); + if (Array.isArray(data.departure_dates) && data.departure_dates.length === 1) details.departure_date = txt(data.departure_dates[0]); + if (action === 'passenger_list_import') { + const n = context.imported_count ?? context.row_count ?? result.imported_count ?? obj(data.passenger_list).row_count; + if (n !== undefined && n !== null && txt(n) !== '' && Number.isInteger(Number(n)) && Number(n) >= 0) details.roster_count = Number(n); + } + if (action.startsWith('order_update_')) details.updates = list(obj(data.updates).actions) + .filter(u => Object.hasOwn(UPDATE_LABELS, txt(u.target))) + .map(u => ({ label: UPDATE_LABELS[txt(u.target)], value: txt(u.value), mode: txt(u.operation) })); + if (action === 'order_cancel' || action === 'order_restore') details.order_status = action === 'order_cancel' + ? '已取消' : first(context.to_status, obj(data.transition).to_status, report.to_status) || '已恢复'; + if (action === 'confirmation_export') details.files = [...new Set(list(report.artifacts) + .filter(a => a.agentbus_visible !== false).map(a => txt(a.type)).filter(t => Object.hasOwn(FILES, t)))]; + if (refs.kind === 'shared_plan' || action === 'shared_plan_create' || action === 'order_update_shared_plan') details.team_kind = 'shared_plan'; + const targets: Array<{ group: string; child: string; departure?: string; createdChild?: boolean }> = []; + const batches = list(result.batch_results).length ? list(result.batch_results) : list(report.batch_results); + if (action === 'shared_child_order_batch_create') { + for (const b of batches) if (b.status === 'completed' && !list(b.blockers).length && txt(b.child_order_no)) + targets.push({ group: businessGroup(b.parent_group_no), child: group(b.child_order_no), departure: txt(b.departure_date) }); + } else if (action === 'shared_child_order_create' || refs.kind === 'shared_child_order') { + const children = [...new Set(strings(receipt.order_number, refs.child_order_no, + refs.kind === 'shared_child_order' ? refs.identifier : '', result.order_number).map(group).filter(Boolean))]; + const parents = [...new Set(strings(refs.parent_group_no, receipt.parent_group_no, + obj(report.erp_receipt).parent_group_no).map(businessGroup).filter(Boolean))]; + if (children.length > 1 || parents.length > 1) return []; + if (children.length === 1) targets.push({ group: parents[0] || '', child: children[0] }); + } else { + const groups = strings(receipt.group_numbers, receipt.group_number, refs.identifier, refs.group_no).map(businessGroup); + const unique = [...new Set(groups.filter(Boolean))]; + // A single operation resolving to conflicting group identities is not merged. + const batch = ['team_order_batch_create', 'shared_plan_create'].includes(action); + if (!batch && unique.length > 1) return []; + for (const g of unique) { + const matching = batches.filter(b => strings(b.group_number, obj(b.erp_receipt).group_number).map(group).includes(g)); + targets.push({ group: g, child: '', ...(matching.length === 1 ? { departure: first(matching[0].date, matching[0].departure_date) } : {}) }); + } + // Split-plan creation can prove concrete child identities per parent/date. + for (const row of list(obj(obj(report.erp_receipt || result.erp_receipt).split_order_probe).per_date)) { + if (row.facts_determined !== true || !unique.includes(group(row.parent_group_no))) continue; + for (const child of list(row.child_refs)) if (txt(child.child_order_no)) targets.push({ group: group(row.parent_group_no), + child: group(child.child_order_no), departure: txt(row.departure_date), createdChild: true }); + } + const requery = obj(report.requery || result.requery); + if (action === 'order_cancel' && refs.kind === 'shared_plan' && unique.length === 1 && requery.child_status_propagation_matched === true) { + for (const child of list(requery.child_refs)) if (/^\d+$/.test(txt(child.ddid))) targets.push({ group: unique[0], child: `D${txt(child.ddid)}` }); + } + } + const seen = new Set(); + return targets.filter(t => { const key = `${t.group}/${t.child}`; if (seen.has(key)) return false; seen.add(key); return true; }) + .map(t => ({ version: 1, group_number: t.group, subject_kind: t.child ? 'child' : 'team', subject_number: t.child || t.group, + action: t.createdChild ? 'shared_child_order_create' : action, source_task_id: input.taskId, execution_id: input.executionId, + effective_at: effectiveAt, recorded_at: recordedAt, details: { ...details, ...(t.departure ? { departure_date: t.departure } : {}) } })); +} + +// An interrupted batch is not a successful task. Preserve only independently +// confirmed rows from the existing adapter's per-item completion contract. +export function partialBatchProgressEvents(input: ProgressInput): ProgressEvent[] { + const op = obj(input.operation), result = obj(input.result), report = obj(result.report); + if (result.status === 'completed' || !['blocked', 'failed', 'reconciliation_pending'].includes(txt(result.status))) return []; + const batches = list(result.batch_results).length ? list(result.batch_results) : list(report.batch_results); + if (op.action === 'shared_child_order_batch_create' && report.status === 'split_child_batch_stopped_after_write') { + const verified = batches.filter(b => b.status === 'completed' && b.adapter_status === 'split_child_completed' + && txt(b.parent_group_no) && txt(b.child_order_no) && (!Array.isArray(b.blockers) || !b.blockers.length)); + return progressEventsFromSuccess({ ...input, receipt: { success: true }, result: { status: 'completed', batch_results: verified } }); + } + if (op.action === 'team_order_batch_create' && report.status === 'batch_fallback_incomplete') { + return batches.flatMap(b => { + const receipt = obj(b.erp_receipt); + if (b.status !== 'completed' || receipt.success === false || !txt(receipt.group_number)) return []; + return progressEventsFromSuccess({ ...input, operation: { ...op, + data: { ...obj(op.data), departure_dates: txt(b.date) ? [txt(b.date)] : [] } }, + result: { status: 'completed' }, receipt }); + }); + } + return []; +} + +function ordered(events: ProgressEvent[]): ProgressEvent[] { + return [...events].sort((a,b) => a.effective_at.localeCompare(b.effective_at) || a.recorded_at.localeCompare(b.recorded_at) + || a.execution_id.localeCompare(b.execution_id)); +} +function latest(events: ProgressEvent[], match: (e: ProgressEvent) => boolean): ProgressEvent | undefined { return events.filter(match).at(-1); } +function displayTime(value: string): string { + return new Intl.DateTimeFormat('sv-SE', { timeZone: 'Asia/Shanghai', year: 'numeric', month: '2-digit', day: '2-digit', + hour: '2-digit', minute: '2-digit' }).format(new Date(value)); +} +function arrangementLine(action: string, details: Json): string { + const fields = [txt(details.resource_name), txt(details.item), + details.start_date ? `${txt(details.start_date)} 至 ${txt(details.end_date) || '离店/结束日未记录'}` : txt(details.date), + details.quantity === undefined ? '' : `${txt(details.quantity)}${action === 'arrangement_hotel' ? '间' : '(数量)'}`, + details.phone ? `${action === 'arrangement_vehicle' ? '司机电话' : '电话'}:${txt(details.phone)}` : '', + details.vehicle_number ? `车号:${txt(details.vehicle_number)}` : '', details.grade ? `等级:${txt(details.grade)}` : '', + details.filing_number ? `备案号:${txt(details.filing_number)}` : '', details.filing_date ? `备案日期:${txt(details.filing_date)}` : '', + details.filing_entry_port ? `入境口岸:${txt(details.filing_entry_port)}` : '', details.filing_exit_port ? `出境口岸:${txt(details.filing_exit_port)}` : '', + details.remark ? `备注:${txt(details.remark)}` : ''].filter(Boolean); + if (!details.resource_name) fields.unshift('资源名称未记录'); + return fields.length > 1 ? fields.join('|') : '已安排,明细未记录'; +} + +export function teamProgressText(value: unknown): string { + const snap = value as ProgressSnapshot; + if (!snap || snap.version !== 1 || snap.scope !== 'current_account' || !Array.isArray(snap.groups)) return ''; + const lines = ['本系统团队办理进度(当前账号)']; + for (const entry of snap.groups) { + const events = ordered(entry.events), team = events.filter(e => e.subject_kind === 'team'); + lines.push(`\n团号:${entry.group_number}`); + const facts = Object.assign({}, ...team.map(e => e.details)) as Json; + for (const [key, label] of [['customer', '客户'], ['product', '产品'], ['departure_date', '出发日期']] as const) + if (txt(facts[key])) lines.push(`${label}:${txt(facts[key])}`); + if (!facts.customer && facts.customer_condition) lines.push(`客户条件:${txt(facts.customer_condition)}`); + if (!facts.product && facts.product_condition) lines.push(`产品条件:${txt(facts.product_condition)}`); + const status = latest(team, e => e.action === 'order_cancel' || e.action === 'order_restore' + || ['team_order_create','team_order_batch_create','shared_plan_create'].includes(e.action)); + lines.push(`订单状态:${status ? txt(status.details.order_status) || '已创建' : '暂无本系统状态记录'}`); + const subjects = [...new Set(events.filter(e => e.subject_kind === 'child').map(e => e.subject_number))].sort(); + const shared = facts.team_kind === 'shared_plan' || subjects.length > 0; + const creation = latest(team, e => ['team_order_create','team_order_batch_create','shared_plan_create'].includes(e.action)); + lines.push(`${shared ? '建计划' : '下单'}:${creation ? `已完成(${displayTime(creation.effective_at)})` : '暂无本系统成功记录'}`); + const subjectProgress = (subset: ProgressEvent[], label: string, parent = false): void => { + const roster = latest(subset, e => e.action === 'passenger_list_import'); + if (!parent) lines.push(`${label}名单:${roster ? `已导入,最近一次${roster.details.roster_count === undefined ? '人数未记录' : `${roster.details.roster_count}人`}(${displayTime(roster.effective_at)})` : '暂无本系统导入记录'}`); + const update = latest(subset, e => e.action.startsWith('order_update_')); + lines.push(`${label}最近修改:${update ? `${list(update.details.updates).map(u => `${txt(u.label)}${u.mode === 'append' ? '追加' : '改为'}${txt(u.value)}`).join(';') || '修改成功,明细未记录'}(${displayTime(update.effective_at)})` : '暂无本系统修改记录'}`); + const exports = new Map(); + for (const e of subset) if (e.action === 'confirmation_export') for (const f of strings(e.details.files)) exports.set(f, e); + const unknownExport = latest(subset, e => e.action === 'confirmation_export' && !strings(e.details.files).length); + if (!exports.size && !unknownExport) lines.push(`${label}文件:暂无本系统导出记录`); + if (unknownExport) lines.push(`${label}文件:已导出,文件类型未记录(${displayTime(unknownExport.effective_at)})`); + for (const [file, e] of exports) { + const changed = events.some(change => change.action !== 'confirmation_export' && change.effective_at > e.effective_at + && (change.subject_kind === 'team' || subset.includes(change) || (file === 'visitor-list' && shared))); + lines.push(`${label}${parent && file === 'visitor-list' ? '整团游客信息' : FILES[file] || '团队文件'}:已导出(${displayTime(e.effective_at)})${changed ? ';之后有业务变更,可能需要重新导出' : ''}`); + } + }; + if (!shared) subjectProgress(team, ''); + else { + lines.push(`本系统记录子单:${subjects.length}个(不代表母团全部子单)`); + // Parent modifications and whole-group exports belong to the parent. + subjectProgress(team, '母团', true); + for (const child of subjects) { + const subset = events.filter(e => e.subject_kind === 'child' && e.subject_number === child); + const created = subset.some(e => ['shared_child_order_create','shared_child_order_batch_create'].includes(e.action)); + const state = latest(subset, e => ['order_cancel','order_restore','shared_child_order_create','shared_child_order_batch_create'].includes(e.action)); + lines.push(`子单 ${child}:${created ? '已新增' : '暂无本系统新增记录'};状态:${state ? txt(state.details.order_status) || '已创建' : '未记录'}`); + subjectProgress(subset, ` ${child} `); + } + } + lines.push('【安排情况与明细】'); + for (const [action, label] of Object.entries(ARRANGEMENT_LABELS)) { + const items = entry.arrangements.filter(a => a.action === action); + lines.push(`${label}:${items.length ? `已安排,共${items.length}条` : '未安排'}`); + items.forEach((a,i) => lines.push(` ${i+1}. ${arrangementLine(action, a.details)}`)); + } + } + if (!snap.groups.length || snap.unlinked_subjects.length) lines.push(`\n${snap.unlinked_subjects.length ? `对象:${snap.unlinked_subjects.join('、')}。` : ''}暂无法关联团号,未推断安排状态。`); + lines.push('\n仅统计当前账号通过本系统成功完成的操作;未记录不代表 ERP 中没有。已安排表示至少一条,不代表全行程已安排齐全。'); + return lines.join('\n'); +} + +export class TeamProgressLedger { + constructor(private readonly config: AppConfig) {} + async prepare(client: Client, scope: ArrangementScope, receiptFromResult: (v: unknown) => Json | null): Promise { + if (!scope.organizationId || !scope.userId) throw new Error('team_progress_scope_missing'); + await client.query('INSERT INTO team_progress_owners (organization_id, owner_user_id) VALUES ($1,$2) ON CONFLICT DO NOTHING', [scope.organizationId, scope.userId]); + const locked = await client.query('SELECT backfilled_at FROM team_progress_owners WHERE organization_id = $1 AND owner_user_id = $2 FOR UPDATE', [scope.organizationId, scope.userId]); + if (locked.rows[0]?.backfilled_at) return; + let cursor = ''; + for (;;) { + const rows = await client.query(`SELECT t.task_id, t.operation_ciphertext, t.operation, t.execution_result_ciphertext, t.execution_result, + t.success_receipt_at, t.updated_at, a.id AS execution_id, a.started_at AS effective_at + FROM tasks t JOIN task_attempts a ON a.task_id = t.id AND a.phase = 'erp' + AND a.id::text = COALESCE(t.execution_result->>'execution_id', '') + WHERE t.organization_id = $1 AND t.assigned_user_id = $2 + AND t.status IN ('completed','blocked','failed','reconciliation_pending') AND t.task_id > $3 + ORDER BY t.task_id LIMIT 200`, [scope.organizationId, scope.userId, cursor]); + for (const row of rows.rows) { + const decode = (cipher: unknown, plain: unknown): unknown => cipher ? JSON.parse(decryptText(this.config, txt(cipher))) : plain; + const result = decode(row.execution_result_ciphertext, row.execution_result); + const input = { taskId: row.task_id, executionId: row.execution_id, effectiveAt: row.effective_at, + recordedAt: row.success_receipt_at || row.updated_at, operation: decode(row.operation_ciphertext, row.operation), result, receipt: receiptFromResult(result) }; + const events = [...progressEventsFromSuccess(input), ...partialBatchProgressEvents(input)]; + for (const e of events) await this.insert(client, scope, e); + } + if (rows.rows.length < 200) break; + cursor = txt(rows.rows.at(-1)?.task_id); + } + await client.query('UPDATE team_progress_owners SET backfilled_at = now() WHERE organization_id = $1 AND owner_user_id = $2', [scope.organizationId, scope.userId]); + } + async insert(client: Client, scope: ArrangementScope, e: ProgressEvent): Promise { + await client.query(`INSERT INTO team_progress_events + (organization_id, owner_user_id, group_key, subject_key, subject_kind, event_key, source_task_id, execution_id, effective_at, recorded_at, detail_ciphertext) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) ON CONFLICT (organization_id, owner_user_id, execution_id, event_key) DO NOTHING`, + [scope.organizationId, scope.userId, e.group_number ? sha256Text(e.group_number) : '', sha256Text(e.subject_number), e.subject_kind, + sha256Text(`${e.action}/${e.group_number}/${e.subject_kind}/${e.subject_number}`), e.source_task_id, e.execution_id, + e.effective_at, e.recorded_at, encryptText(this.config, JSON.stringify(e))]); + } + async knownGroups(client: Client, scope: ArrangementScope, child: string): Promise { + const rows = await client.query(`SELECT detail_ciphertext FROM team_progress_events WHERE organization_id = $1 AND owner_user_id = $2 + AND subject_key = $3 AND subject_kind = 'child' AND group_key <> ''`, [scope.organizationId, scope.userId, sha256Text(child)]); + return strings(...rows.rows.map(r => obj(JSON.parse(decryptText(this.config, r.detail_ciphertext))).group_number)); + } + async history(client: Client, scope: ArrangementScope, groupNumber: string): Promise { + const rows = await client.query(`SELECT e.detail_ciphertext, e.source_task_deleted_at FROM team_progress_events e + WHERE e.organization_id = $1 AND e.owner_user_id = $2 AND (e.group_key = $3 OR (e.group_key = '' AND e.subject_kind = 'child' + AND e.subject_key IN (SELECT subject_key FROM team_progress_events WHERE organization_id = $1 AND owner_user_id = $2 AND group_key = $3) + AND NOT EXISTS (SELECT 1 FROM team_progress_events conflict WHERE conflict.organization_id = $1 AND conflict.owner_user_id = $2 + AND conflict.subject_key = e.subject_key AND conflict.group_key <> '' AND conflict.group_key <> $3))) + ORDER BY e.effective_at, e.recorded_at, e.id`, [scope.organizationId, scope.userId, sha256Text(group(groupNumber))]); + return rows.rows.map(r => ({ ...JSON.parse(decryptText(this.config, r.detail_ciphertext)), source_task_deleted_at: date(r.source_task_deleted_at) || null })); + } + async snapshot(client: Client, scope: ArrangementScope, events: ProgressEvent[], recordedAt: string, + arrangements: (g: string) => Promise): Promise { + const groups = new Set(), unlinked = new Set(); + for (const e of events) { + if (e.group_number) groups.add(e.group_number); + else { + const parents = await this.knownGroups(client, scope, e.subject_number); + if (parents.length === 1) groups.add(parents[0]); else unlinked.add(e.subject_number); + } + } + const snapshots: ProgressSnapshot['groups'] = []; + for (const g of groups) snapshots.push({ group_number: g, events: await this.history(client, scope, g), + arrangements: currentArrangements(await arrangements(g)).map(a => ({ action: a.action, details: a.details })) }); + return { version: 1, scope: 'current_account', recorded_at: recordedAt, groups: snapshots, unlinked_subjects: [...unlinked] }; + } + async markTasksDeleted(client: Client, scope: ArrangementScope, taskIds: string[]): Promise { + await client.query(`UPDATE team_progress_events SET source_task_deleted_at = now() WHERE organization_id = $1 AND owner_user_id = $2 + AND source_task_id = ANY($3::text[]) AND source_task_deleted_at IS NULL`, [scope.organizationId, scope.userId, taskIds]); + } +} diff --git a/control-plane/test/arrangement-ledger.test.ts b/control-plane/test/arrangement-ledger.test.ts index 0d67882..4f7c01a 100644 --- a/control-plane/test/arrangement-ledger.test.ts +++ b/control-plane/test/arrangement-ledger.test.ts @@ -8,6 +8,7 @@ import { loadConfig } from '../src/config.js'; import { decryptText, encryptText } from '../src/crypto.js'; import { closePool, getPool } from '../src/db.js'; import { TaskService, successReceiptFromResult } from '../src/task-service.js'; +import { progressMemory } from './support/team-progress-memory.js'; import { buildCustomerSuccessReceipt, replyTextFromReceipt } from '../src/reply-contract.js'; const config = loadConfig({ NODE_ENV: 'test', FIELD_ENCRYPTION_KEY: Buffer.alloc(32, 29).toString('base64'), @@ -146,12 +147,14 @@ test('receipt summary is a frozen snapshot and remains single when public receip }); function memoryClient() { + const progress = progressMemory(config); const owners = new Map(); const records: Json[] = [], tasks: Json[] = []; const queries: string[] = []; const answer = (rows: Json[] = []) => ({ rowCount: rows.length, rows }); const client = { release() {}, async query(sql: string, p: any[] = []) { queries.push(sql); + if (sql.includes('team_progress_') || (sql.includes('FROM tasks t JOIN task_attempts') && sql.includes('COALESCE'))) return progress.query(sql, p); const ownerKey = `${p[0]}/${p[1]}`; if (/INSERT INTO arrangement_ledger_owners/.test(sql)) { if (!owners.has(ownerKey)) owners.set(ownerKey, {}); return answer(); } if (/SELECT backfilled_at/.test(sql)) { assert.match(sql, /FOR UPDATE/); return answer([owners.get(ownerKey)!]); } diff --git a/control-plane/test/control-plane.test.ts b/control-plane/test/control-plane.test.ts index 09098dc..7cb1b2b 100644 --- a/control-plane/test/control-plane.test.ts +++ b/control-plane/test/control-plane.test.ts @@ -447,7 +447,7 @@ test('control plane requires the arrangement ledger migration before readiness', const { readFile } = await import('node:fs/promises'); const db = await readFile(new URL('../src/db.ts', import.meta.url), 'utf8'); const server = await readFile(new URL('../src/server.ts', import.meta.url), 'utf8'); - assert.equal(REQUIRED_SCHEMA_VERSION, '025_arrangement_ledger'); + assert.equal(REQUIRED_SCHEMA_VERSION, '026_team_progress'); assert.match(db, /schema_migrations/); assert.match(db, /databaseReadiness/); assert.match(db, /assertDatabaseSchema/); diff --git a/control-plane/test/support/team-progress-memory.ts b/control-plane/test/support/team-progress-memory.ts new file mode 100644 index 0000000..69255cc --- /dev/null +++ b/control-plane/test/support/team-progress-memory.ts @@ -0,0 +1,30 @@ +import { decryptText } from '../../src/crypto.js'; +import type { AppConfig } from '../../src/config.js'; +import type pg from 'pg'; + +// Stateful test database; scope, unique keys, orphan joins and backfill are all exercised. +export function progressMemory(config: AppConfig) { + const owners = new Map(); + const records: Record[] = [], tasks: Record[] = []; + const answer = (rows: Record[] = []) => ({ rowCount: rows.length, rows }); + const query = async (sql: string, p: any[] = []) => { + const owner = `${p[0]}/${p[1]}`, scoped = records.filter(r => r.organization_id === p[0] && r.owner_user_id === p[1]); + if (sql.includes('INSERT INTO team_progress_owners')) { if (!owners.has(owner)) owners.set(owner,false); return answer(); } + if (sql.includes('SELECT backfilled_at FROM team_progress_owners')) return answer([{backfilled_at:owners.get(owner)?'done':null}]); + if (sql.includes('UPDATE team_progress_owners')) { owners.set(owner,true); return answer(); } + if (sql.includes('FROM tasks t JOIN task_attempts') && sql.includes('COALESCE')) return answer(tasks.filter(t => t.organization_id === p[0] && t.assigned_user_id === p[1] + && ['completed','blocked','failed','reconciliation_pending'].includes(t.status) && t.task_id > p[2] && t.execution_id === t.execution_result.execution_id).sort((a,b) => a.task_id.localeCompare(b.task_id)).slice(0,200)); + if (sql.includes('INSERT INTO team_progress_events')) { + if (!scoped.some(r => r.execution_id === p[7] && r.event_key === p[5])) records.push({organization_id:p[0],owner_user_id:p[1],group_key:p[2],subject_key:p[3], + subject_kind:p[4],event_key:p[5],source_task_id:p[6],execution_id:p[7],effective_at:p[8],recorded_at:p[9],detail_ciphertext:p[10],source_task_deleted_at:null}); + return answer(); + } + if (sql.includes('SELECT detail_ciphertext FROM team_progress_events')) return answer(scoped.filter(r => r.subject_key === p[2] && r.subject_kind === 'child' && r.group_key)); + if (sql.includes('SELECT e.detail_ciphertext')) return answer(scoped.filter(r => r.group_key === p[2] || (!r.group_key && r.subject_kind === 'child' + && scoped.some(link => link.subject_key === r.subject_key && link.group_key === p[2]) + && !scoped.some(link => link.subject_key === r.subject_key && link.group_key && link.group_key !== p[2])))); + if (sql.includes('UPDATE team_progress_events')) { scoped.filter(r => p[2].includes(r.source_task_id)).forEach(r => r.source_task_deleted_at = new Date()); return answer(); } + throw new Error(`Unexpected progress SQL: ${sql}`); + }; + return {query: query as unknown as pg.PoolClient['query'],records,tasks,owners,decode:() => records.map(r => JSON.parse(decryptText(config,r.detail_ciphertext)))}; +} diff --git a/control-plane/test/team-progress.test.ts b/control-plane/test/team-progress.test.ts new file mode 100644 index 0000000..767e40d --- /dev/null +++ b/control-plane/test/team-progress.test.ts @@ -0,0 +1,192 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import { readFile } from 'node:fs/promises'; +import { loadConfig } from '../src/config.js'; +import { encryptText } from '../src/crypto.js'; +import { TeamProgressLedger, progressEventsFromSuccess, teamProgressText, type ProgressInput, type ProgressEvent, type ProgressSnapshot } from '../src/team-progress.js'; +import { successReceiptFromResult } from '../src/task-service.js'; +import { buildCustomerSuccessReceipt, replyTextFromReceipt } from '../src/reply-contract.js'; +import { progressMemory } from './support/team-progress-memory.js'; + +const config = loadConfig({ NODE_ENV:'test', FIELD_ENCRYPTION_KEY:Buffer.alloc(32,24).toString('base64') }); +const scope = {organizationId:'org-a',userId:'user-a'}; +const now = '2026-09-20T01:00:00.123Z'; +function input(action = 'team_order_create', child = '', parent = 'LW-TEST-A'): ProgressInput { + return {taskId:'TASK-test',executionId:'e-1',effectiveAt:now,recordedAt:now, + operation:{action,data:{customer:{name:'客户检索词'},product:{name:'产品检索词'},departure_dates:['2026-10-01']}}, + result:{status:'completed',report:{resolved_refs:{kind:child?'shared_child_order':'independent_order',identifier:child||parent,...(child?{parent_group_no:parent}: {})}}}, + receipt:child?{success:true,order_number:child}:{success:true,group_number:parent}}; +} +function ev(action: string, options: {child?:string,parent?:string,time?:string,details?:Record} = {}): ProgressEvent { + const i=input(action,options.child,options.parent); + i.effectiveAt=options.time||now;i.recordedAt=i.effectiveAt;i.executionId=`e-${action}-${i.effectiveAt}-${options.child||''}`; + const e=progressEventsFromSuccess(i)[0];assert.ok(e);return {...e,details:{...e.details,...options.details}}; +} +function snapshot(events:ProgressEvent[]):ProgressSnapshot { + return {version:1,scope:'current_account',recorded_at:now,groups:[{group_number:'LW-TEST-A',events,arrangements:[]}],unlinked_subjects:[]}; +} + +test('all business actions project accepted success without roster contents or arbitrary result fields',()=>{ + for(const action of ['team_order_create','team_order_batch_create','shared_plan_create','shared_child_order_create','passenger_list_import','order_update_independent', + 'order_update_shared_plan','order_update_shared_child','order_cancel','order_restore','confirmation_export','arrangement_guide','arrangement_vehicle','arrangement_hotel','arrangement_transport','arrangement_other']) { + const i=input(action,action.includes('shared_child')?'D100':''); + (i.operation as any).data.passenger_list={row_count:20,tsv:'SECRET PASSPORT'}; + (i.result as any).raw_html='SECRET HTML'; + const events=progressEventsFromSuccess(i);assert.equal(events.length,1,action); + assert.doesNotMatch(JSON.stringify(events),/SECRET/); + } +}); + +test('failed, uncertain, blocker and unproved results do not create completed business records',()=>{ + for(const status of ['failed','blocked','running','reconciliation_pending','dry_run']) { + const i=input();(i.result as any).status=status;assert.deepEqual(progressEventsFromSuccess(i),[]); + } + for(const patch of [{uncertain:true},{manual_review_required:true},{blockers:['x']}]){ + const i=input();Object.assign(i.result as any,patch);assert.deepEqual(progressEventsFromSuccess(i),[]); + } + const i=input();i.receipt=null;assert.deepEqual(progressEventsFromSuccess(i),[]); +}); + +test('individual group conflicts are rejected and full suffixes remain distinct',()=>{ + const i=input();(i.receipt as any).group_number='LW-TEST-C';assert.deepEqual(progressEventsFromSuccess(i),[]); + assert.equal(ev('team_order_create',{parent:' lw-test-c '}).group_number,'LW-TEST-C'); +}); + +test('internal identifiers and conflicting child-parent evidence cannot become group associations',()=>{ + for (const parent of ['12345','D12345']) assert.deepEqual(progressEventsFromSuccess(input('confirmation_export','',parent)),[]); + const i=input('shared_child_order_create','D1');(i.receipt as any).parent_group_no='LW-OTHER'; + assert.deepEqual(progressEventsFromSuccess(i),[]); + const j=input('shared_child_order_create');j.receipt={success:true}; + assert.deepEqual(progressEventsFromSuccess(j),[],'a parent identifier alone is not a child identity'); +}); + +test('batch children are explicitly mapped to their respective parents, never zipped by array order',()=>{ + const i=input('shared_child_order_batch_create'); + (i.result as any).report.batch_results=[{status:'completed',parent_group_no:'LW-A',child_order_no:'D2',departure_date:'2026-10-01'}, + {status:'completed',parent_group_no:'LW-C',child_order_no:'D1',departure_date:'2026-10-03'}]; + assert.deepEqual(progressEventsFromSuccess(i).map(e=>[e.group_number,e.subject_number,e.details.departure_date]),[['LW-A','D2','2026-10-01'],['LW-C','D1','2026-10-03']]); + (i.result as any).report.batch_results[0].status='not_started';assert.equal(progressEventsFromSuccess(i).length,1); +}); + +test('native child-create receipts carry the verified parent; missing parents remain unlinked',()=>{ + const i=input('shared_child_order_create','D123');(i.result as any).report={erp_receipt:{order_number:'D123',parent_group_no:'LW-P'}}; + assert.equal(progressEventsFromSuccess(i)[0].group_number,'LW-P'); + (i.result as any).report={};assert.equal(progressEventsFromSuccess(i)[0].group_number,''); +}); + +test('roster counts use the latest import per subject and never add repeated uploads',()=>{ + const events=[ev('passenger_list_import',{details:{roster_count:20}}),ev('passenger_list_import',{time:'2026-09-20T02:00:00Z',details:{roster_count:18}})]; + const text=teamProgressText(snapshot(events));assert.match(text,/最近一次18人/);assert.doesNotMatch(text,/38人|20人/); + const late={...events[0],recorded_at:'2026-09-21T00:00:00Z'}; + assert.match(teamProgressText(snapshot([events[1],late])),/最近一次18人/); +}); + +test('shared-plan summaries separate each child and never claim every roster is complete',()=>{ + const events=[ev('shared_plan_create'),ev('passenger_list_import',{child:'D1',details:{roster_count:20}}),ev('shared_child_order_create',{child:'D2'})]; + const text=teamProgressText(snapshot(events));assert.match(text,/不代表母团全部子单/);assert.match(text,/D1 名单:已导入/);assert.match(text,/D2 名单:暂无/); +}); + +test('cancel/restore updates current state without erasing arrangements or uncancelling children',()=>{ + const s=snapshot([ev('shared_plan_create'),ev('order_cancel'),ev('order_cancel',{child:'D1'}),ev('order_restore',{time:'2026-09-20T02:00:00Z'})]); + s.groups[0].arrangements=[{action:'arrangement_guide',details:{resource_name:'导游甲',phone:'000',grade:'A'}}]; + const text=teamProgressText(s);assert.match(text,/订单状态:已恢复/);assert.match(text,/子单 D1.*已取消/);assert.match(text,/导游:已安排/); +}); + +test('verified parent cancellation propagates only proven child references',()=>{ + const i=input('order_cancel');(i.result as any).report.resolved_refs.kind='shared_plan'; + (i.result as any).report.requery={child_status_propagation_matched:true,child_refs:[{ddid:'123'}]}; + assert.deepEqual(progressEventsFromSuccess(i).map(e=>e.subject_number),['LW-TEST-A','D123']); +}); + +test('all current arrangement rows render supported details; unknown legacy detail is explicit',()=>{ + const s=snapshot([]);s.groups[0].arrangements=[ + {action:'arrangement_guide',details:{resource_name:'导游甲',phone:'000',grade:'A'}}, + {action:'arrangement_vehicle',details:{resource_name:'车队甲',item:'25座',start_date:'2026-10-01',end_date:'2026-10-03',quantity:1,vehicle_number:'车牌甲',phone:'001'}}, + {action:'arrangement_hotel',details:{resource_name:'酒店甲',item:'TWN',start_date:'2026-10-01',end_date:'2026-10-02',quantity:12}}, + {action:'arrangement_hotel',details:{}}, + {action:'arrangement_transport',details:{resource_name:'票务甲',item:'C91',date:'2026-10-01',quantity:20}}, + {action:'arrangement_other',details:{resource_name:'备案项目',filing_number:'ABC',filing_entry_port:'口岸甲',filing_exit_port:'口岸乙',remark:'备注'}}]; + const text=teamProgressText(s);for(const term of ['酒店:已安排,共2条','酒店甲','明细未记录','车牌甲','司机电话','备案号:ABC','C91'])assert.ok(text.includes(term),term); +}); + +test('exports are latest per file type and labelled stale after later changes without deleting records',()=>{ + const i=input('confirmation_export');(i.result as any).report.artifacts=[{type:'visitor-list',agentbus_visible:true,name:'SECRET NAME'},{type:'visitor-list',agentbus_visible:false},{type:'hotel-preorder'}]; + const exported=progressEventsFromSuccess(i)[0];assert.deepEqual(exported.details.files,['visitor-list','hotel-preorder']); + const s=snapshot([exported,ev('order_update_independent',{time:'2026-09-20T02:00:00Z',details:{updates:[{label:'标间',value:'12'}]}})]); + const text=teamProgressText(s);assert.match(text,/可能需要重新导出/);assert.doesNotMatch(text,/SECRET NAME/);assert.match(text,/标间改为12/); + const unknown=teamProgressText(snapshot([ev('confirmation_export')])); + assert.match(unknown,/已导出,文件类型未记录/);assert.doesNotMatch(unknown,/暂无本系统导出记录/); +}); + +test('receipts stay frozen and rebuilding never duplicates the progress block',()=>{ + const i=input('order_update_independent'),base={success:true,group_number:'LW-TEST-A',team_progress:snapshot([ev('team_order_create')])}; + const first=buildCustomerSuccessReceipt(base,i.result,i.operation); + assert.equal(replyTextFromReceipt(buildCustomerSuccessReceipt(first,i.result,i.operation)),replyTextFromReceipt(first)); + assert.equal(replyTextFromReceipt(first).match(/本系统团队办理进度/g)?.length,1); +}); + +test('encrypted ledger isolates accounts and organizations, deduplicates, and survives task deletion',async()=>{ + const db=progressMemory(config),ledger=new TeamProgressLedger(config),e=ev('team_order_create'); + await ledger.prepare(db,scope,successReceiptFromResult);await ledger.insert(db,scope,e);await ledger.insert(db,scope,e); + assert.equal(db.records.length,1);assert.doesNotMatch(db.records[0].detail_ciphertext,/LW-TEST/); + assert.equal((await ledger.history(db,scope,e.group_number)).length,1); + assert.equal((await ledger.history(db,{...scope,userId:'other'},e.group_number)).length,0); + assert.equal((await ledger.history(db,{...scope,organizationId:'other'},e.group_number)).length,0); + await ledger.markTasksDeleted(db,scope,[e.source_task_id]);assert.ok((await ledger.history(db,scope,e.group_number))[0].source_task_deleted_at); +}); + +test('unlinked child history joins only a uniquely known parent in the same account',async()=>{ + const db=progressMemory(config),ledger=new TeamProgressLedger(config); + const orphan=ev('passenger_list_import',{child:'D1',parent:'',details:{roster_count:20}});await ledger.insert(db,scope,orphan); + let s=await ledger.snapshot(db,scope,[orphan],now,async()=>[]);assert.equal(s.groups.length,0);assert.match(teamProgressText(s),/暂无法关联团号/); + const link=ev('shared_child_order_create',{child:'D1',parent:'LW-A'});await ledger.insert(db,scope,link); + s=await ledger.snapshot(db,scope,[orphan],now,async()=>[]);assert.equal(s.groups[0].group_number,'LW-A');assert.equal(s.groups[0].events.length,2); + await ledger.insert(db,scope,ev('shared_child_order_create',{child:'D1',parent:'LW-C'})); + s=await ledger.snapshot(db,scope,[orphan],now,async()=>[]);assert.equal(s.groups.length,0); + assert.equal((await ledger.history(db,scope,'LW-A')).length,1,'ambiguous orphan never leaks into either team'); +}); + +test('history backfill preserves only the matching successful execution including archived tasks',async()=>{ + const db=progressMemory(config),ledger=new TeamProgressLedger(config),i=input(); + const result={status:'completed',execution_id:'e-1',erp_receipt:{success:true,group_number:'LW-TEST-A'}}; + const row={task_id:i.taskId,execution_id:'e-1',effective_at:now,success_receipt_at:now,organization_id:scope.organizationId,assigned_user_id:scope.userId, + status:'completed',archived_at:now,operation_ciphertext:encryptText(config,JSON.stringify(i.operation)),execution_result:result,execution_result_ciphertext:encryptText(config,JSON.stringify(result))}; + db.tasks.push(row,{...row,execution_id:'other-attempt'},{...row,task_id:'TASK-other',assigned_user_id:'other'}); + await ledger.prepare(db,scope,successReceiptFromResult);assert.equal(db.records.length,1); + await ledger.prepare(db,scope,successReceiptFromResult);assert.equal(db.records.length,1); +}); + +test('migration preserves business records without task foreign keys and prevents administrator ownership',async()=>{ + const sql=await readFile(new URL('../migrations/026_team_progress.sql',import.meta.url),'utf8'); + assert.doesNotMatch(sql,/REFERENCES\s+(tasks|task_attempts)\b/i);assert.match(sql,/reject_admin_task_principal/);assert.match(sql,/detail_ciphertext text NOT NULL/); +}); + +test('partial batches retain only independently verified completed rows',async()=>{ + const {partialBatchProgressEvents}=await import('../src/team-progress.js'); + const i=input('shared_child_order_batch_create'); + i.result={status:'reconciliation_pending',uncertain:true,blockers:['stop'],report:{status:'split_child_batch_stopped_after_write',batch_results:[ + {status:'completed',adapter_status:'split_child_completed',parent_group_no:'LW-A',child_order_no:'D1',blockers:[]}, + {status:'uncertain',adapter_status:'split_child_live_uncertain',parent_group_no:'LW-C',child_order_no:'D2'}, + {status:'not_started',parent_group_no:'LW-D',child_order_no:''}]}}; + assert.deepEqual(partialBatchProgressEvents(i).map(e=>[e.group_number,e.subject_number]),[['LW-A','D1']]); + assert.deepEqual(progressEventsFromSuccess(i),[],'does not relabel the parent task as successful'); + (i.result as any).report.batch_results[0].adapter_status='unverified';assert.deepEqual(partialBatchProgressEvents(i),[]); + const j=input('team_order_batch_create');j.result={status:'blocked',report:{status:'batch_fallback_incomplete'},batch_results:[ + {status:'completed',erp_receipt:{group_number:'LW-A',success:true},date:'2026-10-01'}, + {status:'saved_unverified',erp_receipt:{group_number:'LW-B'}}]}; + assert.equal(partialBatchProgressEvents(j).length,1); + assert.equal(partialBatchProgressEvents(j)[0].action,'team_order_batch_create'); +}); + +test('batch team dates follow their own receipt rows and never copy a date onto all groups',()=>{ + const i=input('team_order_batch_create');(i.operation as any).data.departure_dates=['2026-10-01','2026-10-02']; + i.receipt={success:true,group_numbers:['LW-A','LW-C']};i.result={status:'completed',batch_results:[ + {date:'2026-10-02',erp_receipt:{group_number:'LW-C'}},{date:'2026-10-01',erp_receipt:{group_number:'LW-A'}}]}; + assert.deepEqual(progressEventsFromSuccess(i).map(e=>[e.group_number,e.details.departure_date]),[['LW-A','2026-10-01'],['LW-C','2026-10-02']]); +}); + +test('channel length limits clearly disclose truncation while leaving short replies unchanged',async()=>{ + const {boundedAgentBusReplyText}=await import('../src/agentbus-delivery.js'); + assert.equal(boundedAgentBusReplyText('已完成'),'已完成'); + const long=boundedAgentBusReplyText('一条明细\n'.repeat(6000));assert.ok(long.length<=20000);assert.match(long,/完整内容请在平台任务回执中查看/); +}); diff --git a/tools/arrangement-ledger-postgres-smoke.mjs b/tools/arrangement-ledger-postgres-smoke.mjs index 3ac369f..40523ab 100644 --- a/tools/arrangement-ledger-postgres-smoke.mjs +++ b/tools/arrangement-ledger-postgres-smoke.mjs @@ -9,6 +9,7 @@ import { getPool, closePool, quoteIdentifier } from '../control-plane/src/db.ts' import { encryptText } from '../control-plane/src/crypto.ts'; import { ArrangementLedger, currentArrangements } from '../control-plane/src/arrangement-ledger.ts'; import { TaskService } from '../control-plane/src/task-service.ts'; +import { TeamProgressLedger } from '../control-plane/src/team-progress.ts'; import { taskResultText } from '../control-plane/src/agentbus.ts'; const rawUrl = process.env.LTJT_ARRANGEMENT_TEST_DATABASE_URL; @@ -20,6 +21,7 @@ const config = loadConfig({ NODE_ENV: 'test', DATABASE_URL: rawUrl, DATABASE_SCH FIELD_ENCRYPTION_KEY: Buffer.alloc(32, 31).toString('base64') }); const pool = getPool(config); const ledger = new ArrangementLedger(config); +const progress = new TeamProgressLedger(config); let created = false; try { await pool.query(`CREATE SCHEMA ${quoteIdentifier(schema)}`); @@ -75,6 +77,7 @@ try { assert.equal(retained.find(e => e.source_task_id === first.taskId).details.resource_name, '实际资源'); assert.equal((await pool.query('SELECT 1 FROM task_attempts WHERE id = $1', [first.executionId])).rowCount, 0); assert.equal(currentArrangements(retained).length, 2, 'task cascade leaves ledger intact'); + assert.ok((await progress.history(pool, ownerA, 'LW-TEST-A')).find(e => e.source_task_id === first.taskId).source_task_deleted_at); const secondHotel = await task(ownerA, 'arrangement_hotel', 'create', '200'); await complete(ownerA, secondHotel); @@ -89,15 +92,17 @@ try { service.emitEvent = emit; assert.equal((await pool.query('SELECT status FROM tasks WHERE task_id = $1', [rollback.taskId])).rows[0].status, 'running'); assert.ok(!(await ledger.history(pool, ownerA, 'LW-TEST-A')).some(e => e.execution_id === rollback.executionId)); + assert.ok(!(await progress.history(pool, ownerA, 'LW-TEST-A')).some(e => e.execution_id === rollback.executionId)); // Deleting an archived, pre-feature success as the first interaction backfills it. const historic = await task(ownerB, 'arrangement_transport'); await pool.query(`UPDATE tasks SET status = 'completed', archived_at = now(), success_receipt_at = now(), - execution_result_ciphertext = $2 WHERE task_id = $1`, [historic.taskId, encryptText(config, JSON.stringify(historic.result))]); + execution_result_ciphertext = $2, execution_result = $3 WHERE task_id = $1`, [historic.taskId, encryptText(config, JSON.stringify(historic.result)), {execution_id:historic.executionId}]); await service.hardDeleteTask(ownerB, historic.taskId); const backfill = await ledger.history(pool, ownerB, 'LW-TEST-A'); assert.equal(backfill.length, 1); assert.equal(backfill[0].action, 'arrangement_transport'); + assert.equal((await progress.history(pool, ownerB, 'LW-TEST-A')).length, 1); assert.ok(backfill[0].source_task_deleted_at); assert.ok(!(await ledger.history(pool, ownerA, 'LW-TEST-A')).some(e => e.execution_id === historic.executionId)); @@ -105,10 +110,58 @@ try { [org, context('ledger-admin').userId]), e => e.code === '23514'); await assert.rejects(pool.query(`INSERT INTO arrangement_ledger_owners (organization_id, owner_user_id) VALUES ($1,$2)`, [otherOrg, ownerA.userId]), e => e.code === '23503'); + // Exercise the same transactional API for non-arrangement operations. + async function business(owner, action, data, refs, extra = {}) { + const item = await task(owner, action); + item.operation = {action, data}; + item.result = {...item.result, summary:{action}, report:{...item.result.report, resolved_refs:refs}, ...extra}; + await pool.query('UPDATE tasks SET operation = $2, operation_ciphertext = $3 WHERE task_id = $1', + [item.taskId,{action},encryptText(config,JSON.stringify(item.operation))]); + return item; + } + const parentRefs = {kind:'shared_plan',identifier:'LW-PROGRESS'}; + const childRefs = {kind:'shared_child_order',identifier:'D300',parent_group_no:'LW-PROGRESS'}; + const plan = await business(ownerA,'shared_plan_create',{},parentRefs); + await complete(ownerA,plan); + const child = await business(ownerA,'shared_child_order_create',{},childRefs,{erp_receipt:{order_number:'D300',parent_group_no:'LW-PROGRESS'}}); + await complete(ownerA,child); + for (const count of [20,18]) { + const item = await business(ownerA,'passenger_list_import',{passenger_list:{row_count:count,tsv:'PRIVATE ROSTER'}},childRefs); + const outcome = await complete(ownerA,item); + assert.match(outcome.important_message.text,new RegExp('D300 名单:已导入,最近一次'+count+'人')); + assert.doesNotMatch(outcome.important_message.text,/PRIVATE ROSTER/); + if(count===18) await service.hardDeleteTask(ownerA,item.taskId); + } + const cancel = await business(ownerA,'order_cancel',{transition:{to_status:'已取消'}},parentRefs); + cancel.result.report.requery={matched:true,child_status_propagation_matched:true,child_refs:[{ddid:'300'}]}; + await complete(ownerA,cancel); + const restore = await business(ownerA,'order_restore',{transition:{to_status:'已恢复'}},parentRefs); + const restored = await complete(ownerA,restore); + assert.match(restored.important_message.text,/订单状态:已恢复/); + assert.match(restored.important_message.text,/子单 D300.*已取消/); + assert.match(restored.important_message.text,/最近一次18人/,'latest roster survives deletion'); + assert.match((await service.getTask(org,plan.taskId,ownerA)).important_message.text,/本系统记录子单:0个/,'old snapshot stays frozen'); + assert.equal((await progress.history(pool,ownerB,'LW-PROGRESS')).length,0); + assert.equal((await progress.history(pool,{...ownerA,organizationId:otherOrg},'LW-PROGRESS')).length,0); + const partial = await business(ownerA,'shared_child_order_batch_create',{},parentRefs,{ + status:'reconciliation_pending',uncertain:true,report:{status:'split_child_batch_stopped_after_write',batch_results:[ + {status:'completed',adapter_status:'split_child_completed',parent_group_no:'LW-PARTIAL',child_order_no:'D401',blockers:[]}, + {status:'uncertain',parent_group_no:'LW-PARTIAL',child_order_no:'D402'}]} + }); + const partialOutcome = await complete(ownerA,partial); + assert.equal(partialOutcome.status,'reconciliation_pending'); + assert.deepEqual((await progress.history(pool,ownerA,'LW-PARTIAL')).map(e=>e.subject_number),['D401']); + assert.ok(!partialOutcome.success_receipt,'partial progress does not produce a success receipt'); + await assert.rejects(pool.query('INSERT INTO team_progress_owners (organization_id, owner_user_id) VALUES ($1,$2)', + [org,context('ledger-admin').userId]),e=>e.code==='23514'); + await assert.rejects(pool.query('INSERT INTO team_progress_owners (organization_id, owner_user_id) VALUES ($1,$2)', + [otherOrg,ownerA.userId]),e=>e.code==='23503'); console.log(JSON.stringify({ ok: true, migrations: migrations.length, verified: [ 'real migrations', 'result ingestion and public/AgentBus receipts', 'deduplication', 'account/org isolation', 'frozen receipts', 'force-delete retention', 'row-specific clear', 'transaction rollback', - 'archived historical backfill before delete', 'admin and cross-org ownership constraints' + 'archived historical backfill before delete', 'admin and cross-org ownership constraints', + 'non-arrangement receipts', 'per-child latest roster without private data', 'retained roster after delete', 'parent restore preserves child cancellation', + 'partial batch proven rows without false task success' ] })); } finally { if (created) await pool.query(`DROP SCHEMA ${quoteIdentifier(schema)} CASCADE`);