23 KiB
服务端数据库规范
当前状态
服务端使用 drizzle-orm@^0.45.2 + postgres@^3.4.9 建立连接边界,并使用 drizzle-kit@^0.31.10 生成 migration;OneTalk 事实 schema 与面向领域的 repository 已建立。
当前已建立 OneTalk Bright 事实存储 schema,定义位于
apps/server/src/database/schema/onetalk.ts:
onetalk_message:页面事实消息。channel_account_id + conversation_id + message_id复合主键负责幂等;收件和确认发件通过direction区分。content是唯一内容事实,只承载 shared contract 的 versionedtext | image | file | business_card | inquiry | orderJSON;其中business_card严格只保存{ version: 1, kind: "business_card" }marker,不保存客户资料、包含认证信息的完整 OneTalk envelope、业务卡 raw 正文或 SDK payload。onetalk_conversation:插件发现的技术会话和共享同步锚点。channel_account_id + conversation_id复合主键,不按 binding 或设备复制;conversation_kind只接受显式direct,未知历史会话保持null;sync_phase、sync_result、latest_message_id和history_complete表达同步进度及锚点状态,并允许零消息会话。onetalk_contact_profile:Bright 当前联系人资料事实。channel_account_id + conversation_id复合主键,不建立到技术会话表的外键;资料字段允许显式null,只有严格较新的observed_at_ms才能覆盖整行。它是名片读取 view 的唯一客户资料来源,不是消息事实的嵌入列。onetalk_message_anomaly:缺字段、协议和同步异常的独立诊断事实。fingerprint仅用于诊断合并;payload必须由写入边界清洗,不能被消息读取、发送或锚点流程消费。
生成的初始迁移为 apps/server/drizzle/0000_rapid_winter_soldier.sql,其中显式维护 PostgreSQL 表/字段 COMMENT ON 备注(Drizzle 当前版本不会从 TypeScript 注释自动生成数据库备注)。未确认发送不进入任何一张表,普通运行路径不提供物理删除;profile 当前行由读取服务另行受限读取后在内存组合,不复制进 conversation 或 message。媒体切换 migration 0005_young_squadron_supreme 是一次性开发数据重置:仅 DELETE 本仓库拥有的 OneTalk message/anomaly/profile/conversation 事实,再删除 text/content_type 并为 content 加 v1 kind CHECK;不触及授权、binding 或其它渠道。后续 0011_mushy_baron_strucker 在建立名片 exact CHECK 前,将已有带客户资料字段的 business_card content 归一化为 marker。
Scenario: Schema 注释与 PostgreSQL 备注
1. Scope / Trigger
- Trigger:新增或修改任何服务端 PostgreSQL 表、字段、enum 或索引时。
- Scope:schema 源码的可读注释、迁移 SQL 的数据库备注,以及备注内容的安全审查。
2. Signatures
schema source -> apps/server/src/database/schema/**/*.ts
db:generate -> apps/server/drizzle/*.sql + apps/server/drizzle/meta/*_snapshot.json
COMMENT ON TABLE/COLUMN -> PostgreSQL catalog comments
每张业务表和每个业务字段都必须同时具备:
- schema 中紧邻声明的 TypeScript/JSDoc 说明;
- 对应迁移中的
COMMENT ON TABLE或COMMENT ON COLUMN。
3. Contracts
- 备注使用中文或明确的业务术语,说明数据来源、用途、空值语义和是否参与主键/幂等。
- OneTalk 原始字段备注必须标明来源路径,例如
cid、messageId、sender.uid、createAt。 - schema 注释与 PostgreSQL comment 必须语义一致;Drizzle 当前版本不会自动同步两者,数据库备注需在迁移 SQL 中显式维护。
- 备注、
COMMENT ON文本和诊断说明不得包含 Cookie、session、token、连接串或完整敏感 payload。 - 已执行的迁移不可修改;备注变化必须通过新的补偿 migration 更新。
4. Validation & Error Matrix
| Condition | Result |
|---|---|
| 新增表/字段但缺少源码说明 | Review 阻止合并;补齐 JSDoc 后再生成迁移 |
迁移缺少对应 COMMENT ON |
db:check 之外的 schema review 失败;补齐数据库备注 |
| 源码注释与数据库 comment 语义不一致 | 以当前 schema 为准修正并生成补偿 migration,不静默保留漂移 |
| 备注包含认证秘密或完整敏感 payload | 立即删除敏感内容并重新生成/修订未执行的迁移;不得写入数据库 |
| 已执行 migration 需要改备注 | 禁止编辑旧文件,新增只改 comment 的 migration |
5. Good / Base / Bad Cases
- Good:字段旁写清
messageId的 OneTalk 来源和幂等约束,migration 同步写COMMENT ON COLUMN ...。 - Base:本地执行
db:generate、人工审查备注、db:check,再在临时 PostgreSQL 执行db:migrate。 - Bad:只写 TypeScript 注释、只在数据库手工改备注,或把
chatToken等 envelope 内容写进 comment。
6. Tests Required
- Static:检查每个新增表/字段存在源码说明,并在迁移中存在对应
COMMENT ON。 - Migration integration:执行真实迁移后,用
obj_description和col_description断言表/字段备注已落库。 - Migration repeatability:重复执行迁移,备注和
drizzle.__drizzle_migrations不产生重复或漂移。 - Security:备注文本、异常 payload 和日志断言不包含 token、Cookie、连接串或无必要正文。
7. Wrong vs Correct
Wrong
messageId: text("message_id").notNull();
CREATE TABLE "onetalk_message" (...);
Correct
/** OneTalk 消息原始 ID;禁止写入客户端合成 ID。 */
messageId: text("message_id").notNull();
COMMENT ON COLUMN "onetalk_message"."message_id"
IS 'OneTalk 消息原始 ID;禁止写入客户端合成 ID。';
当前边界
- 普通应用数据库连接只能通过
apps/server/src/database/index.ts的createDatabase(databaseUrl)创建。 createDatabase返回{ db, close };close幂等,应用通过 FastifyonClose调用。- 普通应用连接仍由
apps/server/src/database/index.ts创建;迁移使用独立的一次性入口和单连接客户端。 - 业务 schema 位于
apps/server/src/database/schema/,迁移目录和日志标识的单一配置源位于apps/server/src/database/migration-config.ts,生成文件位于apps/server/drizzle/;后续事实存储任务必须在此边界内增加表和事务。
Migration Workflow
采用 code-first 流程:TypeScript schema 是结构源,开发者运行 pnpm --filter @trade-message-center/server db:generate 生成 SQL 和 meta/_journal.json,提交前运行 db:check 审查迁移一致性。生成的 SQL 与元数据必须纳入 Git;已经执行的 migration 不得修改。
db:migrate 是独立的一次性部署 job,使用 drizzle-orm/postgres-js/migrator 和 postgres.js { max: 1 } 客户端读取根目录 .env* 中的 DATABASE_URL。成功后 job 正常退出,连接或 SQL 失败则返回非零状态并释放连接。应用 entry.ts 不自动执行迁移,部署系统必须保证同一环境只有一个 migration job,再启动或滚动更新 server。
Drizzle 的迁移记录固定在 drizzle.__drizzle_migrations。空数据库首次接入时,先添加真实业务 schema,再生成并审查初始 migration;若数据库已有手工表,先确认现状并建立 baseline,不得盲目执行破坏性初始 migration。生产破坏性变更采用 expand-contract 和补偿 migration,不提供自动 down migration。
Migration Contract
1. Scope / Trigger
- Trigger:新增或修改服务端 PostgreSQL 表结构,需要将 schema 变更版本化并由部署 job 应用。
- Scope:
apps/server/src/database/schema/、migration-config.ts、migrate.ts、drizzle.config.ts和数据库脚本。
2. Signatures
db:generate -> SQL migration + drizzle/meta metadata
db:check -> migration consistency check
db:migrate -> Promise<void> / non-zero process exit on failure
readDatabaseUrl(environment) -> string
runMigrations(databaseUrl, dependencies?) -> Promise<void>
3. Contracts
- Required environment key:
DATABASE_URL; blank values fail withMissing DATABASE_URL。 - Schema source is
apps/server/src/database/schema/**/*.ts; generated output isapps/server/drizzle/。 - Migration journal is
drizzle.__drizzle_migrations; migration client usespostgres.js{ max: 1 }。 db:migratecloses its client in both success and failure paths; it never changes application startup behavior。
4. Validation & Error Matrix
| Condition | Result |
|---|---|
Missing/blank DATABASE_URL |
Non-zero exit; Missing DATABASE_URL; no secret value |
| Connection or SQL failure | Non-zero exit; safe error message; client closed |
| Previously applied migration | No duplicate SQL; journal remains stable |
entry.ts starts server |
No migration is executed |
5. Good / Base / Bad Cases
- Good: deployment runs one
db:migratejob, then starts server replicas; generated SQL is reviewed and committed。 - Base: local developer runs
db:generate,db:check, thendb:migrateagainst a disposable PostgreSQL database。 - Bad: route handler or
entry.tscalls migrator,drizzle-kit pushchanges a shared environment, or an applied SQL file is edited。
6. Tests Required
- Unit: missing/blank URL and URL trimming assertions。
- Unit:
runMigrationscloses the injected client after both successful and failed apply。 - Integration: temporary PostgreSQL applies a real migration twice; assert target table exists and journal count is one。
- Existing server format, typecheck, build and dual-source/compiled test commands remain green。
7. Wrong vs Correct
Wrong
// 在应用监听前隐式修改共享数据库结构
await migrate(drizzle(client), { migrationsFolder: "drizzle" });
await app.listen(options);
Correct
deployment migration job: db:migrate (once)
server process: entry.ts -> createApp -> listen
首次引入数据库时必须记录
- 驱动或 ORM 版本、连接配置来源和本地开发启动方式。
- schema/migration 的目录、命名、执行顺序和回滚办法。
- 查询与事务的边界,特别是连接释放、超时和批量操作行为。
- 表、列、索引和外键命名规则,以及敏感数据的存储策略。
- 可重复执行的迁移和测试数据库验证命令。
代码边界
数据库访问应在明确的数据访问模块内完成,路由处理器不应直接拼接查询语句。当前实现位于 apps/server/src/database/index.ts。
查询组合约束
服务端所有关系读取还必须遵守项目级 数据库查询组合:默认禁止 SQL/Drizzle JOIN,先以同一授权 scope 和完整复合键读取各表,再在 read service/repository 的内存投影边界组合。read-repository.ts 的会话/profile 读取采用两次受限查询和复合键 Map 合并;历史读取先独立确认 direct 会话,再读取同 scope 的消息事实。除非本次需求明确指定 JOIN 或有该文档要求的单查询可复核证据,不得增加 JOIN。
Scenario: OneTalk 事实 repository 与提交后发布
1. Scope / Trigger
- Trigger:服务端接收 OneTalk 页面观察、技术会话发现、同步完成或当前联系人资料,并需要写入对应的四类事实表。
- Scope:
apps/server/src/onetalk/repository.ts的事务、幂等和错误边界;WebSocket handler 只消费 service,不直接访问 Drizzle。
2. Signatures
createOneTalkRepository(database) -> OneTalkRepository
repository.insertMessage(context, observationSource, message) -> accepted | duplicate | rejected
repository.discoverConversation(context, conversationId, conversationKind) -> conversation state
repository.recordAnomaly(input) -> Promise<void>
repository.updateSyncState(context, update, conversationId, conversationKind) -> conversation state | null
消息事实的数据库唯一键固定为
channel_account_id + conversation_id + message_id;技术会话和共享锚点的唯一键固定为
channel_account_id + conversation_id。workspace_id、mind_user_id、binding 和
device_id 都是来源/授权上下文,不参与消息唯一性;同一个
channel_account_id 切换到不同 workspace 后,仍落在同一条消息事实的唯一键范围内。
3. Contracts
insertMessage必须在一个数据库事务内先确认会话存在,再ON CONFLICT DO NOTHING写消息;只有实际插入才递增message_count。- 重复消息返回
duplicate,不得覆盖首次事实或再次触发外部事件;允许只更新last_observed_at。 - 跨 workspace 收到相同
channel_account_id + conversation_id + message_id时,必须沿用同一条已存在事实:返回duplicate,保留首次写入的workspace_id、mind_user_id、binding和device_id,不得因后续 workspace 改写来源上下文。 - 消息读取按
channel_account_id + conversation_id读取共享事实;workspace 隔离由 Mind 授权 scope 负责,不能把workspace_id加入消息事实主键或作为第二套消息副本维度。 content jsonb必须是对象,且version=1、kind in (text,image,file,business_card,inquiry,order);应用边界再用 shared exact decoder 验证完整字段。kind=business_card时 JSON 必须精确等于{ "version": 1, "kind": "business_card" },不得把contactName等 view 字段写回数据库。不得保留顶层text、content_type、raw content、params、sign、完整contact或平行投影列。- 名片 view 只能由读取服务先按
channel_account_id + conversation_id读取onetalk_contact_profile,再以内存方式投影四个批准字段;profile 缺失返回 marker,单字段缺失返回null,不得 JOIN 其他账号/会话,也不得回退到登录人资料。 content或文本相同本身不构成重复;只要message_id或conversation_id不同,就按新的 OneTalk 事实入库。discoverConversation只按账号/会话幂等 upsert,不清空已有消息计数、同步结果或锚点。- anomaly 以
fingerprint唯一合并并递增occurrence_count;payload 必须是领域层清洗后的 JSON。 OneTalkDatabaseError只在 repository 边界包装数据库异常;WebSocket 必须返回database_unavailable,不得伪造成功 ACK。- 数据库事务提交后才发送 plugin
message.ack,再向当前授权的精确 Mind scope 发布message.created;发布失败不回滚事实。
4. Validation & Error Matrix
| Condition | Result |
|---|---|
| 未发现技术会话 | 不写消息,返回 rejected(conversation_not_discovered) |
| 消息复合键已存在 | 不新增行、不递增计数、不发布,返回 duplicate |
| 同一消息从不同 workspace 重复上报 | 与普通重复相同:复用原行,保留首次来源上下文,不新增 workspace 副本 |
正文相同但 message_id 或 conversation_id 不同 |
视为不同事实,允许新增行 |
| 观测字段缺失/类型错误 | 只 upsert anomaly,不写消息,返回 anomaly |
| 消息/异常/同步事务失败 | OneTalkDatabaseError;不发成功 ACK 或事件 |
| sync latest ID 为 null 且会话零消息 | 允许成功保持空锚点 |
| sync latest ID 缺失或不属于同账号/会话 | 记录同步 anomaly,状态为 incomplete,不推进锚点 |
| business_card content 含 marker 以外字段 | 数据库 CHECK 拒绝;旧数据须由 0011 迁移先归一化再建立 CHECK |
| business_card 读取时没有同账号同会话 profile | 返回 marker;不伪造空资料或使用消息 contact |
| profile 只有部分批准字段 | 读取 view 对缺失字段返回 null,消息事实不变 |
5. Good / Base / Bad Cases
- Good:先
await service.observeMessage(...),收到accepted后发送 ACK,再调用 publisher;重复结果不调用 publisher。 - Good:workspace A 已写入消息后,workspace B 以相同三元组再次观察;结果为
duplicate,消息表仍只有一行,首条来源上下文不变。 - Base:无
TEST_DATABASE_URL时只运行 fake repository/domain 测试,并将真实 PostgreSQL 用例明确标记 skipped。 - Bad:把
workspace_id拼进消息主键、按 workspace 复制消息,或在 route 中直接insert(onetalkMessage)、用 hash/时间戳补messageId,或在 ACK 前向 Mind 广播。
6. Tests Required
- Domain:断言未知会话 reject、缺字段/未知 content anomaly、fingerprint 合并、metadata-only anomaly、四类同步结果和 latest ID 所属校验。
- WebSocket:断言逐条 ACK、accepted/duplicate/anomaly/rejected 分流、plugin-only 写边界、数据库错误和精确 Mind 发布。
- PostgreSQL:使用显式
TEST_DATABASE_URL执行真实 migration;断言0005只清空 OneTalk owned facts、旧列被删除,0010只扩大 v1 kind CHECK,0011将已有 business_card 归一化为 marker 后建立 exact CHECK;六类 content 都能 round-trip 且 raw card/profile 字段不能入库。另测附带资料的名片插入被拒绝、读取服务按账号/会话 profile 组合 view、无 profile 返回 marker、部分 profile 返回null;同一复合键只有一行、消息计数为 1、首次 observation type 和来源 workspace 上下文保留,并在跨 workspace 重复上报后仍只有一行;测试账号结束后清理。 - Static:
db:check、无 legacy outbox/dispatch 引用、database commit -> ACK -> publish数据流检查。
7. Wrong vs Correct
Wrong
await publisher.publish(message);
await repository.insertMessage(context, source, message);
sendAck({ status: "accepted" });
Correct
const result = await service.observeMessage(context, source, normalizedMessage);
if (result.status === "accepted") {
sendAck(result);
await registry.publishMessageCreated({
message: toOneTalkCenterMessage(result.message),
requestId,
scope,
});
}
Cross-workspace uniqueness boundary
// Wrong: workspaceId creates a second message fact for the same OneTalk message.
const wrongKey = [context.workspaceId, message.conversationId, message.messageId];
// Correct: workspaceId is source context; the fact key is account-scoped.
const factKey = [context.channelAccountId, message.conversationId, message.messageId];
Scenario: OneTalk 联系人资料当前事实与严格时间前进
1. Scope / Trigger
- Trigger:Bright 接收联系人资料观察,需要把 profile 作为可供 conversation read 使用的独立当前事实。
- Scope:
onetalk_contact_profileschema、0003_onetalk_contact_profile_facts.sql、profile repository/service,以及读取时按复合键的实时内存组合。
2. Signatures
createOneTalkProfileRepository(database) -> OneTalkProfileRepository
repository.storeProfiles({ channelAccountId, profiles, receivedAt, commitGuard }) -> { writtenProfileCount, staleProfileCount }
profileService.ingestProfiles({ channelAccountId, profiles, commitGuard }) -> accepted | profile_observed_at_future
表的主键是 (channel_account_id, conversation_id);ali_id 是普通资料字段,不参与唯一性,也不对 onetalk_conversation 建外键。
3. Contracts
- repository 必须在一个事务中写入完整白名单快照,使用
ON CONFLICT的严格条件existing.observed_at_ms < excluded.observed_at_ms。相等或更旧输入不能覆盖任何列。 - 严格较新的输入可以把任意资料列更新为显式
null;current_time_zone用double precision,不得把合法小数截断为整数。 - schema 源码、SQL comment 和 snapshot 必须同步;profile row 允许先于 conversation row 写入。读取服务分别受限读取当前 row,并按
(channelAccountId, conversationId)在内存组合,不复制 profile 历史。 - 事务前后都检查
OneTalkCommitGuard;guard 失败不得返回伪造成功或 ACK。
4. Validation & Error Matrix
| Condition | Result |
|---|---|
duplicate (account, conversation) with newer observation |
full row update |
| duplicate with equal/older observation | no update;increment staleProfileCount |
| explicit null in a strictly newer snapshot | null is persisted and later read as null |
fractional timezone such as 5.5 |
exact PostgreSQL round-trip |
| transaction/guard/database failure | rollback or stable database failure;no success ACK |
missing TEST_DATABASE_URL |
PG integration skipped and reported external_unverified;static db:check still required |
5. Good / Base / Bad Cases
- Good:schema and
COMMENT ONagree, migration is additive, repository uses a composite key and strict timestamp predicate, and read projection combines the current row in memory. - Base:fake repository/domain tests run locally; real migration/upsert/two scoped reads plus in-memory composition use an explicitly configured PostgreSQL URL.
- Bad:key by
ali_id, add a profile-to-conversation foreign key that blocks ordering, update only non-null fields, or store profile JSON in message/anomaly payload.
6. Tests Required
db:checkand static schema/migration/comment review.- PostgreSQL migration round-trip for nulls, fractional timezone, newer/older/equal snapshots, no-FK ordering and two scoped reads plus read-model in-memory composition; run only with
TEST_DATABASE_URL. - Repository/service transaction and commit-guard tests plus WebSocket no-write/no-ACK future-skew tests.
- Assert no profile deletion/TTL, no Mind profile delivery table/path and no raw profile fields in public messages.
7. Wrong vs Correct
Wrong
ON CONFLICT (channel_account_id, conversation_id)
DO UPDATE SET name = COALESCE(EXCLUDED.name, onetalk_contact_profile.name);
Correct
ON CONFLICT (channel_account_id, conversation_id)
DO UPDATE SET name = EXCLUDED.name, current_time_zone = EXCLUDED.current_time_zone
WHERE onetalk_contact_profile.observed_at_ms < EXCLUDED.observed_at_ms;
Design Decision: Drizzle + postgres.js
选择 Drizzle ORM 和 postgres.js 是因为它们满足 PostgreSQL 类型查询与连接池需求;drizzle-kit 仅作为开发期 schema/migration 生成工具,运行时 migrator 复用既有驱动。连接 URL 只作为创建参数传入,不写入日志或响应。
const connection = createDatabase(process.env.DATABASE_URL ?? "");
await connection.close();
迁移命令复用同一驱动,但必须保持独立连接、固定 schema/migration 目录、明确执行顺序和回滚策略。
常见风险
- 在没有迁移记录的情况下手工修改共享数据库。
- 把连接、查询和业务规则写在同一个 HTTP 回调里。
- 将完整连接串、令牌或用户敏感字段写入日志或错误响应。