fix: unify OneTalk buyer facts IndexedDB

This commit is contained in:
YBF
2026-09-15 12:12:46 +08:00
parent aea6085d73
commit cd888acd16
13 changed files with 302 additions and 54 deletions
@@ -84,7 +84,7 @@ OneTalkContactProfileStore.discardPendingProfile(input: {
### Durable lifecycle
- `ONE_TALK_SYNC_DATABASE_VERSION``9`。独立的 `onetalk_contact_profiles` store 以账号/会话对为键;消息、candidate、checkpoint、anomaly 各 store 保持独立。v8→v9 只建立查询 index 并保留 ledger record`hasProfileRecord(account)` 走账号 index countpending flush 以 `[channelAccountId, pending.observedAtMs]` 的合法账号前缀范围读取。无 `pending.observedAtMs` 的 uploaded/rejected record 不得出现在该 index 结果中,也不得回退到 object-store 全表扫描。
- `ONE_TALK_SYNC_DATABASE_VERSION``10`。独立的 `onetalk_contact_profiles` store 以账号/会话对为键;消息、candidate、checkpoint、anomaly 各 store 保持独立。v8→v9 只建立查询 index 并保留 ledger recordv9→v10 只在统一主 schema 新增 `onetalk_buyer_facts`,保留 profile 与所有既有 OneTalk ledger records/indexes`hasProfileRecord(account)` 走账号 index countpending flush 以 `[channelAccountId, pending.observedAtMs]` 的合法账号前缀范围读取。无 `pending.observedAtMs` 的 uploaded/rejected record 不得出现在该 index 结果中,也不得回退到 object-store 全表扫描。
- durable profile 指纹是上传去重的唯一闸门。每收到一条 profile,先在 `putPendingProfile()` 之前读取它的 `[channelAccountId, conversationId]` 记录。若 `pending.fingerprint``lastUploadedFingerprint` 与来入指纹相同,读后即结束:不更新时间戳/高水位、不写 IndexedDB、不发 `profile_observed` 诊断、不 flush、不发送。指纹不同时走既有的 durable-first pending 写入。
- 同键观察对 `getProfile → putPendingProfile` 转换做串行化。并发的相同指纹可以各自执行读取,但最多一个能写入、诊断或 flush。读取失败对调用方保持可见;tracker 清理不得制造第二个未处理 rejection。durable 读进行期间协调器被 dispose 时,其完成不产生写入、诊断或 flush。
- 只有 readwrite transaction 完成且当前 pending 的 fingerprint 与 `observedAtMs` 仍然匹配时,ACK 才标记 uploaded。更旧、未知、重复或不匹配的 ACK 一律 no-op。future-skew 错误只移除精确匹配的 pending snapshot,并记录其被拒观察水位。
@@ -123,7 +123,7 @@ OneTalkContactProfileStore.discardPendingProfile(input: {
- Contract:精确 profile/frame 字段、协议版本、方向/scope、direct discovery 类型、敏感/未知字段拒绝、空/超限批次和 `256 KiB` 字节上限。
- Observerinitial snapshot、可重复的 `syncData` 发布、精确 snapshot/collect payload 校验、targeted direct 查找、群聊/未知目标排除、登录身份、登出/切账号、CRM 客户匹配和头像 URL 校验。
- Ledgerpre-v7`oldVersion < 7`)升级先清空全部旧 OneTalk store 再重建当前 profile 与消息同步状态,v7→v8 升级保留五个既有 storev8→v9 逐条保留六个 OneTalk store record 并创建 scope index;账号/会话键、durable-first 顺序、同 pending/已上传指纹一次读零写、同键并发写、读中 dispose、显式读取失败、ACK/CAS、future-skew 丢弃、重连与重启恢复。测试必须区分 profile account/pending index 与 `objectStore.getAll()`,并确认 uploaded record 不会混入 pending flush。升级必须保留 configuration/deviceId,且不得 rekey 或重试旧 pending ledger。
- Ledgerpre-v7`oldVersion < 7`)升级先清空全部旧 OneTalk store 再重建当前 profile 与消息同步状态,v7→v8 升级保留五个既有 storev8→v9 逐条保留六个 OneTalk store record 并创建 scope indexv9→v10 保留所有既有 OneTalk ledger records/indexes 并新增 buyer facts store;账号/会话键、durable-first 顺序、同 pending/已上传指纹一次读零写、同键并发写、读中 dispose、显式读取失败、ACK/CAS、future-skew 丢弃、重连与重启恢复。测试必须区分 profile account/pending index 与 `objectStore.getAll()`,并确认 uploaded record 不会混入 pending flush。升级必须保留 configuration/deviceId,且不得 rekey 或重试旧 pending ledger。
- Service Worker:既有 Bright binding/read/sync 授权、page-first 与 auth-first 两种顺序的首次认证 setup、ledger 读取后授权丢失、实时 sent/received targeted collect、历史排除、50/51 条 profile 批次、ACK 计数、future 错误映射和过期回调/页面身份 fencedirect discovery 始终携带 `conversationType: "direct"`
- 必须执行定向 typecheck、contract/extension focused 与全量测试、format check 和 `git diff --check`。真实 Chromium、Bright PostgreSQL 和生产 Mind 集成是另行执行的外部检查。
@@ -61,7 +61,7 @@ rejected
### Profile ledger (independent state machine)
联系人资料不使用消息 candidate/checkpoint/anomaly store。`ONE_TALK_SYNC_DATABASE_VERSION=9`:从 `oldVersion < 7` 升级时,transaction 必须先清空五个 OneTalk store(消息、candidate、checkpoint、anomaly、`onetalk_contact_profiles`),不重键或回补旧 ledger;从 v7 升级到 v8 时只新增 `onetalk_conversation_bootstraps` store,必须保留五个既有 store 的记录;从 v7/v8 升级到 v9 时只为既有消息、checkpoint、candidate、anomaly、profile records 建立账号、会话与状态查询索引,六个 store 的记录逐条保留。legacy clear 完成后才建 index,避免为必然删除的 records 建索引。Chrome 配置、deviceId、binding 与其它渠道存储不属于该删除范围;首次 pre-v7 升级后必须重新采集 profile 并从 clean state 执行 full syncv7→v8/v9 不得因为 schema 演进清理既有事实。v7 checkpoint 新增 OneTalk 会话列表来源的 `latestMessageAtMs`,用于同步完成后更新服务端会话活动时间。
联系人资料不使用消息 candidate/checkpoint/anomaly store。`ONE_TALK_SYNC_DATABASE_VERSION=10`:从 `oldVersion < 7` 升级时,transaction 必须先清空五个 OneTalk store(消息、candidate、checkpoint、anomaly、`onetalk_contact_profiles`),不重键或回补旧 ledger;从 v7 升级到 v8 时只新增 `onetalk_conversation_bootstraps` store,必须保留五个既有 store 的记录;从 v7/v8 升级到 v9 时只为既有消息、checkpoint、candidate、anomaly、profile records 建立账号、会话与状态查询索引,六个 store 的记录逐条保留;从 v9 升级到 v10 时只在统一 `trade-message-center` schema 新增 `onetalk_buyer_facts` store,必须保留所有既有 OneTalk stores、records 与 indexes。数据库名、版本与 upgrade handler 只由 `storage.ts` 拥有;buyer ledger 通过该主库路径打开,不得创建独立数据库。legacy clear 完成后才建 index,避免为必然删除的 records 建索引。Chrome 配置、deviceId、binding 与其它渠道存储不属于该删除范围;首次 pre-v7 升级后必须重新采集 profile 并从 clean state 执行 full syncv7→v10 不得因为 schema 演进清理既有事实。v7 checkpoint 新增 OneTalk 会话列表来源的 `latestMessageAtMs`,用于同步完成后更新服务端会话活动时间。
清空后的 `onetalk_contact_profiles` 业务键为 `channelAccountId + conversationId`;记录包含 key、账号、conversationId、资料字段中的 aliId、lastUploadedFingerprint、uploaded/rejected observed high-water mark、updatedAt、lastUploadedAt 和最新 pending 清洗 profile。资料上传到 Bright persistence,不调用 Mind profile HTTP。
@@ -71,7 +71,7 @@ rejected
- 已知 ACK `candidateKey` 必须直接主键读取,且只接纳仍为 `pending_ack` 的 candidate;没有 request key 时才兼容查找 `(channelAccountId, conversationId, messageId)`。确认后只比较 checkpoint 当前 `latestMessageId` candidate 与本次 confirmed candidate 的 `candidateOrder`,不得重建全会话候选列表。
- profile pending 读取必须使用 `[channelAccountId, pending.observedAtMs]` 的合法账号前缀 `IDBKeyRange`;缺少 `pending.observedAtMs` 的 uploaded/rejected record 没有该 index entry,不得由调用方再过滤整表。
- 会话重建必须在一个 readwrite transaction 中,对目标 scope index 执行 `getAllKeys()` 并删除 message/candidate/anomaly,再直接删除该会话 checkpoint;不得在 transaction 外读后删,也不得使用 object-store `getAll()`
- v9 profile 一旦打开,v8 binary 打开同一数据库会因版本较低失败。发布回滚只能通过 v10 前向迁移,或保留 v9 schema 并仅禁用非 schema 行为;不得重发 v8 或清除 durable ledger。
- v10 schema 一旦打开,v9 binary 打开同一数据库会因版本较低失败。发布回滚只能通过后续前向迁移,或保留 v10 schema 并仅禁用非 schema 行为;不得重发 v9 或清除 durable ledger。
Observation 先写最新 pending,再由现有 Bright WebSocket 发送 contact.profile.observed。同 fingerprint 且无 pending 时只推进更高 observed high-water mark;任何不高于 uploaded/rejected high-water mark 的不同 fingerprint 也跳过。断线、Service Worker 重启或新页面连接只从 pending 重建发送。收到 contact.profile.ack 后,必须等待 readwrite transaction oncomplete,且只确认仍匹配的 fingerprint 和 observedAtMs;迟到旧 ACK 不得删除新 pending。收到 `ws.error` `profile_observed_at_future` 时只丢弃该 request 的精确 pending,避免无限重试。
@@ -179,7 +179,7 @@ Service Worker 重启后必须从 IndexedDB 恢复:
- 断线恢复会重新发现会话并恢复未确认事实和 completion 声明。
- \`delivery_unknown\` 不创建发送任务、不自动重发;迟到消息继续进入普通 observation。
- 页面同步路由、Bright ACK 和 IndexedDB 事务顺序在重启/断线下保持一致。
- `oldVersion < 7` 升级到 v9 必须断言五个 OneTalk store 均为空、没有 profile 重键/重试分支,且下一次 bootstrap 重新采集 profile、会话活动时间并执行 full sync;v7→v8 必须断言新增 bootstrap store 且五个既有 store 逐条保留;v8→v9 必须断言六个 OneTalk store 的 records 逐条保留、目标 index 存在且读取不调用 `objectStore.getAll()`;配置、deviceId 与其它渠道数据在所有升级路径保持不变。
- `oldVersion < 7` 升级到 v10 必须断言五个 OneTalk store 均为空、没有 profile 重键/重试分支,且下一次 bootstrap 重新采集 profile、会话活动时间并执行 full sync;v7→v8 必须断言新增 bootstrap store 且五个既有 store 逐条保留;v8→v9 必须断言六个 OneTalk store 的 records 逐条保留、目标 index 存在且读取不调用 `objectStore.getAll()`v9→v10 必须断言所有既有 OneTalk records/indexes 逐条保留且新增 `onetalk_buyer_facts`,并验证 buyer ledger 以主数据库名与 v10 通过同一 upgrade handler 打开;配置、deviceId 与其它渠道数据在所有升级路径保持不变。
## 7. Wrong vs Correct
@@ -17,7 +17,7 @@ OneTalk 插件同时需要以下能力时,遵循本总览和对应子规范:
| --- | --- |
| [OneTalk 页面桥、Port 与命令路由](./page-bridge.md) | MAIN/ISOLATED/SW 页面桥、Port 注册、页面身份和 command 路由 |
| [OneTalk 耐久同步与连接生命周期](./durable-sync.md) | IndexedDB、full/incremental/live、ACK、checkpoint、重启恢复和连接生命周期 |
| [OneTalk 联系人资料 Bright 持久化](./contact-profile-sync.md) | profile 白名单、Bright profile frame、pre-v7 清空后重采集、v9 scope index、ACK/HWM/future-skew 和账号/epoch 隔离 |
| [OneTalk 联系人资料 Bright 持久化](./contact-profile-sync.md) | profile 白名单、Bright profile frame、pre-v7 清空后重采集、v9 scope index、v10 统一 buyer store、ACK/HWM/future-skew 和账号/epoch 隔离 |
| [OneTalk 扩展安装实例设备身份](./device-identity.md) | deviceId 生成、迁移、独立存储、配置清除和生命周期 |
| [OneTalk Service Worker 状态与诊断](./runtime-diagnostics.md) | getSnapshot、错误投影、敏感信息脱敏和 development 构建 |
| [OneTalk PWA 出站发送 SOP](./send-sop.md) | sendUIMessages 输入、SDK-only 发送和 WebSocket 旁路事实确认 |
@@ -36,7 +36,7 @@ OneTalk MAIN world
-> Bright WebSocket upload
-> per-message ACK
联系人资料事实使用独立路径:OneTalk MAIN snapshot/syncData → safe profile envelope + page identity → Service Worker pre-v7 clean-state / v9 scope-indexed profile ledger → existing Bright plugin WebSocket `contact.profile.observed` → Bright guarded profile transaction → `contact.profile.ack`
联系人资料事实使用独立路径:OneTalk MAIN snapshot/syncData → safe profile envelope + page identity → Service Worker pre-v7 clean-state / v9 scope-indexed profile ledger(位于 v10 统一 schema→ existing Bright plugin WebSocket `contact.profile.observed` → Bright guarded profile transaction → `contact.profile.ack`
\`\`\`
服务端命令:
@@ -247,7 +247,7 @@ writer.sendSendConfirmation(result);
- `authorization_unavailable` 表示授权依赖暂时不可用,只关闭当前 Bright socket 并沿既有连接退避自动重连;只有凭证、授权版本、binding、scope 或协议版本等确定性错误才阻断自动重连并进入 unauthorized。
- Service Worker 重启从 IndexedDB 恢复 checkpoint、候选和模式,不信任旧内存 cursor。
- live `messageType: "new"`sent 或 received)只触发所属 conversation 的 profile collecthistory 不触发。profile 相同 fingerprint 在 coordinator 的唯一 durable read 后静默结束,不能通过 MAIN `seen`、时间水位或 timer 再建第二去重状态。
- `oldVersion < 7` 的数据库先执行全量 OneTalk state 清空(v7→v8 只新增 bootstrap store、保留既有事实;v8→v9 只新增 scope index、逐条保留 records),再重新采集 profile、会话活动时间并执行 full sync;之后的重启/重连才从当前 ledger 恢复 pending。v9 profile pending flush 只读取 `[channelAccountId, pending.observedAtMs]` index,不得扫描其它账号或已上传 record。ACK 只在当前 `[channelAccountId, conversationId, fingerprint, observedAtMs]` 的 IndexedDB transaction commit 后生效;future-skew 整批拒绝并只丢弃匹配 pending。
- `oldVersion < 7` 的数据库先执行全量 OneTalk state 清空(v7→v8 只新增 bootstrap store、保留既有事实;v8→v9 只新增 scope index、逐条保留 recordsv9→v10 只在统一主 schema 新增 buyer facts store 并保留既有 stores、records 与 indexes),再重新采集 profile、会话活动时间并执行 full sync;之后的重启/重连才从当前 ledger 恢复 pending。v9 profile pending flush 只读取 `[channelAccountId, pending.observedAtMs]` index,不得扫描其它账号或已上传 record。ACK 只在当前 `[channelAccountId, conversationId, fingerprint, observedAtMs]` 的 IndexedDB transaction commit 后生效;future-skew 整批拒绝并只丢弃匹配 pending。
### Send and protocol boundaries
@@ -0,0 +1,3 @@
{"file":".trellis/spec/project/architecture.md","reason":"检查 schema 打开路径和业务 store 职责是否保持单一所有者。"}
{"file":".trellis/spec/chrome-extension/frontend/onetalk/durable-sync.md","reason":"检查 v10 升级不破坏 durable-first、ACK 与前向 schema 语义。"}
{"file":".trellis/spec/chrome-extension/frontend/quality-guidelines.md","reason":"检查扩展测试、严格类型、构建与格式化门禁。"}
@@ -0,0 +1,46 @@
# 统一 buyer facts IndexedDB 设计
## 决策
`onetalk_buyer_facts` 进入 `trade-message-center` 的 OneTalk 统一 schema。
`storage.ts` 是数据库名称、schema version 与 `onupgradeneeded` 的唯一 owner
buyer facts 的领域读写逻辑继续留在 `buyer-fact-store.ts`
不迁移历史数据,也不把旧库清理写入扩展运行时。用户确认扩展尚未上线,因此旧的
`trade-message-center-onetalk-buyer-facts` 仅在当前开发浏览器 profile 中直接
删除。这样不会为不存在的生产兼容需求引入第二套状态、启动 I/O 或失败分支。
## Schema 与调用边界
1. `storage.ts` 将主库 version 从 9 升至 10,并在统一 store registry 加入
`onetalk_buyer_facts`
2. 主库打开 helper 由 `storage.ts` 导出给 buyer facts store 使用;所有调用者通过
同一数据库名、version 和 upgrade handler 打开数据库。
3. `buyer-fact-store.ts` 删除私有的数据库名、version 与 upgrade handler,只保留
`onetalk_buyer_facts` 为对象存储的 pending、ACK 和 high-water 写入语义。
4. `configured-sync-session.ts` 的 coordinator 装配和 Bright 协议不变;它继续只
依赖 `OneTalkBuyerFactStore` 接口。
这避免了“仅把数据库名改为主库”造成的 version 竞争:如果 buyer facts 自行请求
不同 version 或自行 upgrade,旧 Service Worker 打开主库时可能抛出
`VersionError`,并且无法保证既有 stores/indexes 同步创建。
## 删除与回滚
- 代码切换后,用浏览器 IndexedDB 工具直接删除旧独立库;删除前不读取、不迁移其
records,符合未上线且无数据保留要求。
- 清理后重新加载扩展并触发 buyer facts 写入,验证旧库不会被重建。
- v10 一旦打开,v9 或更旧源码无法安全打开同一主库;本任务不允许通过恢复旧独立库名
回滚。出现问题时,只能发布后续兼容 v10 schema 的前向修复;本任务不提供数据回滚,
因为用户明确放弃了旧开发数据,也不得将浏览器数据库降级到 v9。
## 验证设计
- buyer facts focused test 使用记录 `open(name, version)` 的 fake,断言主库名和
主 schema version,并覆盖 buyer store 随主库 upgrade 创建。
- storage focused test 从 v9 升至 v10,断言既有 stores/records/indexes 保留且新增
buyer store。
- 现有 buyer pending/ACK/contact-details regressions 保持通过,证明持久化位置改变
未改变领域状态机。
- 实现后运行 focused tests、typecheck、format check、build 与相关 Chrome runtime
smoke;最后直接查看当前 profile 的 IndexedDB 名称。
@@ -0,0 +1,3 @@
{"file":".trellis/spec/project/architecture.md","reason":"统一 IndexedDB schema owner、模块职责与状态唯一所有者约束。"}
{"file":".trellis/spec/chrome-extension/frontend/onetalk/durable-sync.md","reason":"OneTalk durable ledger 的 v9 升级行为、前向 schema 与 ACK 不变量。"}
{"file":".trellis/spec/chrome-extension/frontend/quality-guidelines.md","reason":"扩展包类型检查、构建、测试与格式化验证要求。"}
@@ -0,0 +1,32 @@
# 实施计划
## 预检
1. 读取 `trellis-before-dev` 与本任务相关的扩展规范。
2. 对将被修改的 `openSyncDatabase`(或其最终共享 helper)和
`createOneTalkBuyerFactStore` 执行 GitNexus upstream impact analysis;若风险为
HIGH/CRITICAL,先向用户报告再编辑。
3. 确认工作树只含本任务规划文件变更。
## 实现
1. 先扩展 buyer facts store 的 fake/test,令其断言实际主数据库名与 version,且
能运行统一 upgrade handler。
2.`storage.ts` 集中加入 buyer object store、升级 schema 到 v10,并提供唯一的
主库打开路径;新增 v9→v10 保留既有数据的回归测试。
3.`buyer-fact-store.ts` 复用该打开路径和 store 常量,删除私有数据库配置与
upgrade handler;保留所有事实合并、pending 和 ACK 行为。
4. 更新 `durable-sync.md`,将 v10 store/升级不变量记为当前基线。
5. 在已更新的开发扩展中直接删除错误独立库;重新加载后进行 browser runtime
smoke,确认新的 buyer write 只创建/使用主库。
## 验证与评审
1. 运行 buyer facts store、sync storage 与 Service Worker storage focused tests。
2. 运行 `pnpm format:check``pnpm typecheck``pnpm build` 与适用的扩展测试。
3. 运行 GitNexus `detect_changes()`,确认范围只包括 OneTalk storage、buyer store、
对应测试及 durable-sync 规范。
4. 检查 diff:不存在旧数据库名、迁移/兼容分支、无关 schema 删除或 buyer 状态机
行为改变。
5. 通过 Trellis check 后提交一个与本任务对应的 commit;不推送或创建 PR,除非用户
后续明确授权。
@@ -0,0 +1,59 @@
# 统一 OneTalk buyer facts IndexedDB
## Goal
将 buyer facts 账本纳入统一 OneTalk IndexedDB schema,避免再创建错误的独立
数据库;由于该扩展尚未上线,直接删除当前开发环境中的旧库,不保留迁移兼容层。
## Background
- 当前 `apps/chrome-extension/src/onetalk/service-worker/buyer-fact-store.ts:12-14`
`onetalk_buyer_facts` 打开在独立数据库
`trade-message-center-onetalk-buyer-facts`(版本 1);这是截图中多出的
IndexedDB。
- 统一 OneTalk 数据库是
`apps/chrome-extension/src/onetalk/service-worker/storage.ts:18-26`
`trade-message-center`(当前 schema 版本 9)。消息、同步、联系人资料与
rendered-card 账本均在该库中。
- 联系人资料账本虽在业务语义上独立,仍通过
`storage.ts:1189-1201` 的统一数据库打开路径使用主库。这说明 buyer facts
的独立数据库不是“独立 ledger”的必要条件。
- buyer facts 的初始实现(commit `24e0228`)即使用独立库;测试 fake 的
`open(_name, version)` 忽略数据库名称,故此前未覆盖该错误。
## Requirements
- R1`onetalk_buyer_facts` 必须属于统一 `trade-message-center` IndexedDB
schema,由一个版本与升级入口管理;不得再创建
`trade-message-center-onetalk-buyer-facts`
- R2:不实现旧独立库的数据迁移、兼容读取或运行时清理分支。旧库仅作为当前
未上线开发环境的本地状态直接删除。
- R3:合库后 buyer facts 的 key、同账号隔离、pending/ACK 语义,以及已确认
字段不被后续 pending 快照清除的现有行为保持不变。
- R4:在源码不再引用旧库后,直接删除当前浏览器 profile 中的
`trade-message-center-onetalk-buyer-facts`;该删除已获用户授权。
- R5:测试必须断言实际打开的数据库名和升级路径,而非只使用忽略库名的 fake。
## Acceptance Criteria
- [ ] 新安装或当前开发扩展升级后,DevTools 只将 buyer facts 写入
`trade-message-center``onetalk_buyer_facts` object store。
- [ ] 不存在旧库迁移、兼容读取或自动清理的生产代码;当前开发 profile 的旧
独立库被直接删除。
- [ ] 升级主库不会遗漏现有 sync/profile/rendered-card stores 和 indexes,也不
触发版本冲突。
- [ ] 清理旧库后重新加载扩展并触发 buyer facts 写入,不会重新创建旧库。
- [ ] buyer fact store 的 focused tests、扩展 typecheck 与相关存储测试通过。
## Out of Scope
- 不改变 buyer facts 的 DOM 观察、Bright wire contract、服务端持久化或授权。
- 不改变其他 OneTalk ledger 的数据模型、清理策略或索引,除非统一 schema
升级必需。
- 不清除浏览器中与此次错误库无关的 IndexedDB 数据。
## Notes
- 用户已确认“直接删除,我没有发上线”。本任务仍涉及统一 IndexedDB schema
version 与 durable ledger,属于复杂任务;需补充 `design.md`
`implement.md`,再进行实现审批。
@@ -0,0 +1,26 @@
{
"id": "unify-buyer-facts-indexeddb",
"name": "unify-buyer-facts-indexeddb",
"title": "统一 OneTalk buyer facts IndexedDB",
"description": "将 buyer facts 账本纳入统一 OneTalk IndexedDB schema,并删除未上线开发环境的错误独立库。",
"status": "in_progress",
"dev_type": null,
"scope": null,
"package": null,
"priority": "P2",
"creator": "ybf",
"assignee": "ybf",
"createdAt": "2026-09-15",
"completedAt": null,
"branch": null,
"base_branch": "main",
"worktree_path": null,
"commit": null,
"pr_url": null,
"subtasks": [],
"children": [],
"parent": null,
"relatedFiles": [],
"notes": "",
"meta": {}
}
@@ -1,4 +1,4 @@
// 持久化独立的 OneTalk 买家事实待投递账本
// 持久化 OneTalk 买家事实待投递账本
import {
cloneOneTalkBuyerFact,
@@ -9,9 +9,7 @@ import {
type OneTalkBuyerFactSource,
} from "@trade-message-center/onetalk-contract";
export const ONE_TALK_BUYER_FACT_STORE_NAME = "onetalk_buyer_facts";
const DATABASE_NAME = "trade-message-center-onetalk-buyer-facts";
const DATABASE_VERSION = 1;
import { ONE_TALK_BUYER_FACT_STORE_NAME, openSyncDatabase } from "./storage.ts";
export type OneTalkStoredBuyerFact = {
key: string;
@@ -143,20 +141,6 @@ const transactionResult = (transaction: IDBTransaction): Promise<void> =>
reject(transaction.error ?? new Error("IndexedDB write aborted"));
});
const openDatabase = (factory: IDBFactory): Promise<IDBDatabase> =>
new Promise((resolve, reject) => {
const request = factory.open(DATABASE_NAME, DATABASE_VERSION);
request.onupgradeneeded = () => {
if (!request.result.objectStoreNames.contains(ONE_TALK_BUYER_FACT_STORE_NAME)) {
request.result.createObjectStore(ONE_TALK_BUYER_FACT_STORE_NAME, {
keyPath: "key",
});
}
};
request.onsuccess = () => resolve(request.result);
request.onerror = () => reject(request.error ?? new Error("IndexedDB open failed"));
});
/** 创建只属于 buyer facts 的 durable-first 本地账本。 */
export const createOneTalkBuyerFactStore = (
factory: IDBFactory = indexedDB,
@@ -164,7 +148,7 @@ export const createOneTalkBuyerFactStore = (
): OneTalkBuyerFactStore => {
let databasePromise: Promise<IDBDatabase> | null = null;
const database = (): Promise<IDBDatabase> => {
databasePromise ??= openDatabase(factory).catch((error) => {
databasePromise ??= openSyncDatabase(factory).catch((error: unknown) => {
databasePromise = null;
throw error;
});
@@ -16,7 +16,7 @@ import {
import { traceOneTalkObservedMessages } from "../diagnostics/message-trace.ts";
export const ONE_TALK_MESSAGE_DATABASE_NAME = "trade-message-center";
export const ONE_TALK_SYNC_DATABASE_VERSION = 9;
export const ONE_TALK_SYNC_DATABASE_VERSION = 10;
export const ONE_TALK_MESSAGE_STORE_NAME = "onetalk_messages";
export const ONE_TALK_CHECKPOINT_STORE_NAME = "onetalk_sync_checkpoints";
export const ONE_TALK_CANDIDATE_STORE_NAME = "onetalk_sync_candidates";
@@ -24,6 +24,7 @@ export const ONE_TALK_ANOMALY_STORE_NAME = "onetalk_sync_anomalies";
export const ONE_TALK_CONTACT_PROFILE_STORE_NAME = "onetalk_contact_profiles";
export const ONE_TALK_CONVERSATION_BOOTSTRAP_STORE_NAME = "onetalk_conversation_bootstraps";
export const ONE_TALK_RENDERED_CARD_LEDGER_STORE_NAME = "onetalk_rendered_card_ledger";
export const ONE_TALK_BUYER_FACT_STORE_NAME = "onetalk_buyer_facts";
export const ONE_TALK_CONVERSATION_BOOTSTRAP_MIGRATION_ID = "direct-discovery-before-history-v1";
export type OneTalkConversationBootstrapPhase =
@@ -362,6 +363,7 @@ const ensureSyncStores = (database: IDBDatabase): void => {
ONE_TALK_CONTACT_PROFILE_STORE_NAME,
ONE_TALK_CONVERSATION_BOOTSTRAP_STORE_NAME,
ONE_TALK_RENDERED_CARD_LEDGER_STORE_NAME,
ONE_TALK_BUYER_FACT_STORE_NAME,
];
for (const storeName of stores) {
if (!database.objectStoreNames.contains(storeName)) {
@@ -425,7 +427,7 @@ const clearLegacyOneTalkStores = (transaction: IDBTransaction): void => {
}
};
const openSyncDatabase = (factory: IDBFactory): Promise<IDBDatabase> => {
export const openSyncDatabase = (factory: IDBFactory): Promise<IDBDatabase> => {
return new Promise((resolve, reject) => {
const request = factory.open(
ONE_TALK_MESSAGE_DATABASE_NAME,
@@ -5,6 +5,11 @@ import test from "node:test";
import { createOneTalkBuyerFactFingerprint } from "@trade-message-center/onetalk-contract";
import { createOneTalkBuyerFactStore } from "../src/onetalk/service-worker/buyer-fact-store.ts";
import {
ONE_TALK_BUYER_FACT_STORE_NAME,
ONE_TALK_MESSAGE_DATABASE_NAME,
ONE_TALK_SYNC_DATABASE_VERSION,
} from "../src/onetalk/service-worker/storage.ts";
const fact = (observedAtMs, fingerprint, tag) => ({
conversationId: "conversation-1",
@@ -40,6 +45,10 @@ class Request {
class Store {
constructor() {
this.records = new Map();
this.indexNames = {
values: new Set(),
contains: (name) => this.indexNames.values.has(name),
};
}
attach(transaction) {
@@ -57,6 +66,14 @@ class Store {
put(record) {
this.transaction.stage(this.records, record.key, record);
}
clear() {
this.records.clear();
}
createIndex(name) {
this.indexNames.values.add(name);
}
}
class Transaction {
@@ -108,14 +125,20 @@ class Database {
class Factory {
constructor() {
this.database = new Database();
this.opens = [];
}
open(_name, version) {
open(name, version) {
this.opens.push({ name, version });
const request = new Request(undefined, false);
request.transaction = new Transaction(this.database);
queueMicrotask(() => {
if (this.database.version < version) {
request.result = this.database;
request.onupgradeneeded?.();
request.onupgradeneeded?.({
oldVersion: this.database.version,
newVersion: version,
});
this.database.version = version;
}
setImmediate(() => request.onsuccess?.());
@@ -124,6 +147,20 @@ class Factory {
}
}
test("buyer ledger opens the unified database and creates its store through the main upgrade", async () => {
const factory = new Factory();
const store = createOneTalkBuyerFactStore(factory, () => 100);
assert.deepEqual(await store.listPendingFacts("account-1"), []);
assert.deepEqual(factory.opens, [
{
name: ONE_TALK_MESSAGE_DATABASE_NAME,
version: ONE_TALK_SYNC_DATABASE_VERSION,
},
]);
assert.equal(factory.database.stores.has(ONE_TALK_BUYER_FACT_STORE_NAME), true);
});
test("buyer ledger retains the latest pending and uploaded high-water mark across delayed observations", async () => {
const store = createOneTalkBuyerFactStore(new Factory(), () => 100);
const current = fact(20, 'tags=["current"]|features=null', "current");
@@ -6,6 +6,7 @@ import test from "node:test";
import { createOneTalkRenderedCardContentFingerprint } from "@trade-message-center/onetalk-contract";
import {
ONE_TALK_ANOMALY_STORE_NAME,
ONE_TALK_BUYER_FACT_STORE_NAME,
ONE_TALK_CANDIDATE_STORE_NAME,
ONE_TALK_CHECKPOINT_STORE_NAME,
ONE_TALK_CONTACT_PROFILE_STORE_NAME,
@@ -31,6 +32,36 @@ const validMessage = {
unreadCount: 0,
};
const v9IndexesByStore = new Map([
[
ONE_TALK_MESSAGE_STORE_NAME,
[["by_account_conversation", ["channelAccountId", "conversationId"]]],
],
[ONE_TALK_CHECKPOINT_STORE_NAME, [["by_account", "channelAccountId"]]],
[
ONE_TALK_CANDIDATE_STORE_NAME,
[
["by_account_status", ["channelAccountId", "status"]],
["by_account_conversation_status", ["channelAccountId", "conversationId", "status"]],
["by_account_conversation", ["channelAccountId", "conversationId"]],
],
],
[
ONE_TALK_ANOMALY_STORE_NAME,
[
["by_account", "channelAccountId"],
["by_account_conversation", ["channelAccountId", "conversationId"]],
],
],
[
ONE_TALK_CONTACT_PROFILE_STORE_NAME,
[
["by_account", "channelAccountId"],
["by_account_pending_observed_at", ["channelAccountId", "pending.observedAtMs"]],
],
],
]);
class FakeRequest {
constructor(result, schedule = true) {
this.result = result;
@@ -278,6 +309,7 @@ test("creates durable stores, keeps confirmed candidates, and merges anomalies",
"onetalk_contact_profiles",
ONE_TALK_CONVERSATION_BOOTSTRAP_STORE_NAME,
ONE_TALK_RENDERED_CARD_LEDGER_STORE_NAME,
ONE_TALK_BUYER_FACT_STORE_NAME,
].sort(),
);
});
@@ -321,8 +353,33 @@ test("preserves existing stores while adding new stores from v7 and v8", async (
}
});
test("preserves v8 ledger records while creating v9 query indexes", async () => {
test("preserves all v8 ledger records while creating v9 query indexes", async () => {
const factory = new FakeFactory(8);
const v8StoreNames = [
ONE_TALK_MESSAGE_STORE_NAME,
ONE_TALK_CANDIDATE_STORE_NAME,
ONE_TALK_CHECKPOINT_STORE_NAME,
ONE_TALK_ANOMALY_STORE_NAME,
ONE_TALK_CONTACT_PROFILE_STORE_NAME,
ONE_TALK_CONVERSATION_BOOTSTRAP_STORE_NAME,
];
for (const storeName of v8StoreNames) {
factory.database.createObjectStore(storeName).put({ key: storeName, v8: true });
}
const store = createOneTalkSyncStore(factory, () => 500);
await store.getCheckpoint("account-1", "conversation-1");
for (const storeName of v8StoreNames) {
assert.equal(factory.database.stores.get(storeName).records.get(storeName).v8, true);
}
for (const [storeName, indexes] of v9IndexesByStore) {
assert.deepEqual([...factory.database.stores.get(storeName).indexes], indexes);
}
});
test("preserves v9 ledger records while adding the buyer facts store in v10", async () => {
const factory = new FakeFactory(9);
for (const storeName of [
ONE_TALK_MESSAGE_STORE_NAME,
ONE_TALK_CANDIDATE_STORE_NAME,
@@ -330,35 +387,34 @@ test("preserves v8 ledger records while creating v9 query indexes", async () =>
ONE_TALK_ANOMALY_STORE_NAME,
ONE_TALK_CONTACT_PROFILE_STORE_NAME,
ONE_TALK_CONVERSATION_BOOTSTRAP_STORE_NAME,
ONE_TALK_RENDERED_CARD_LEDGER_STORE_NAME,
]) {
factory.database.createObjectStore(storeName).put({ key: storeName, v8: true });
factory.database.createObjectStore(storeName).put({ key: storeName, v9: true });
}
for (const [storeName, indexes] of v9IndexesByStore) {
const store = factory.database.stores.get(storeName);
for (const [name, keyPath] of indexes) store.createIndex(name, keyPath);
}
const store = createOneTalkSyncStore(factory, () => 500);
await store.listCheckpoints("account-1");
await store.getCheckpoint("account-1", "conversation-1");
assert.equal(factory.database.version, ONE_TALK_SYNC_DATABASE_VERSION);
for (const storeName of factory.database.stores.keys()) {
assert.equal(factory.database.stores.get(storeName).records.get(storeName).v8, true);
for (const storeName of [
ONE_TALK_MESSAGE_STORE_NAME,
ONE_TALK_CANDIDATE_STORE_NAME,
ONE_TALK_CHECKPOINT_STORE_NAME,
ONE_TALK_ANOMALY_STORE_NAME,
ONE_TALK_CONTACT_PROFILE_STORE_NAME,
ONE_TALK_CONVERSATION_BOOTSTRAP_STORE_NAME,
ONE_TALK_RENDERED_CARD_LEDGER_STORE_NAME,
]) {
assert.equal(factory.database.stores.get(storeName).records.get(storeName).v9, true);
}
assert.equal(factory.database.stores.has(ONE_TALK_BUYER_FACT_STORE_NAME), true);
for (const [storeName, indexes] of v9IndexesByStore) {
assert.deepEqual([...factory.database.stores.get(storeName).indexes], indexes);
}
assert.equal(
factory.database.stores
.get(ONE_TALK_CANDIDATE_STORE_NAME)
.indexNames.contains("by_account_status"),
true,
);
assert.equal(
factory.database.stores
.get(ONE_TALK_CONTACT_PROFILE_STORE_NAME)
.indexNames.contains("by_account_pending_observed_at"),
true,
);
assert.equal(
factory.database.stores
.get(ONE_TALK_ANOMALY_STORE_NAME)
.indexNames.contains("by_account_conversation"),
true,
);
});
test("clears only one conversation history ledger after its transaction commits", async () => {