Files
JiuContinent eb8267ed18 ```
feat(feishu): 添加飞书入站事件inbox和混合数据库协调功能

- 实现飞书入站事件持久化inbox机制,支持状态管理、租约锁定和重试退避
- 添加混合数据库基线协调工具,确保平台PostgreSQL结构安全对齐
- 增加运行组件心跳检测和readiness就绪检查机制
- 实现app_ticket事件的安全轮换和验证处理
- 添加生产环境运行编排和fail-closed安全机制
- 支持webhook快速确认和长连接独立进程处理
- 完善个人数据擦除时的待处理事件清理功能
```
2026-07-27 17:14:37 +08:00

16 KiB
Raw Permalink Blame History

设计文档

概述

本设计在现有 FastAPI 模块化单体内增加飞书终端用户身份、个人画像与持久化订阅三个 业务域。内部 API 继续使用服务 API Key飞书终端权限只能来自已经通过验签或 SDK 验证的事件,身份主键固定为 (tenant_key, open_id)

平台只写自身数据库。旧业务 MySQL 保持 SELECT-only。个人定时提示词不允许调用工具、 读取公司报表或写入问答历史;群订阅也不加载创建者个人画像。

架构设计

系统架构图

flowchart LR
    F["飞书 webhook / 长连接"] --> V["事件验证与解密边界"]
    V --> I["feishu_users<br/>身份与权限"]
    I --> C["飞书命令编排"]
    C --> P["personalization<br/>规则/偏好/会话"]
    C --> S["subscriptions<br/>计划/投递"]
    C --> A["ai_agent<br/>受限问答"]
    API["内部 API Key"] --> I
    API --> S
    DB[("平台数据库")] --- I
    DB --- P
    DB --- S
    T["每分钟持久化扫描"] --> S
    S --> A
    S --> FC["FeishuClient"]

数据流图

sequenceDiagram
    participant F as 飞书
    participant E as 事件服务
    participant U as 用户服务
    participant C as 命令服务
    participant P as 个性化服务
    participant A as AI

    F->>E: 已签名消息事件
    E->>E: 验签/解密、事件去重
    E->>U: tenant_key + sender open_id
    U->>U: 查找或自动注册普通用户
    U-->>E: FeishuPrincipal
    E->>C: principal + chat context + mentions
    C->>P: 仅加载 principal 所有的数据
    P->>A: 安全规则→公司规则→个人规则→请求→偏好/记忆→历史
    A-->>C: 回答或明确不可用
    C-->>F: 原会话回复

组件与接口

app/modules/feishu_users

  • models.py 定义 FeishuUser,公开随机 code,内部使用整数主键。
  • principal.py 定义不可伪造的 FeishuPrincipal包含用户编号、租户、open_id、角色、 状态、当前 chat_id/chat_type 和结构化 mentions。
  • services/identity.py 只接收验证边界传入的身份,负责首次注册、初始管理员引导、 最后活跃时间和旧自选认领。
  • services/management.py 负责列表、角色/状态修改、最后管理员保护和审计。
  • routes.py 暴露服务 API Key 保护的用户管理接口;正文中的 actor/open_id 不参与授权。

app/modules/personalization

  • UserPreference 保存白名单类别:languagetonedetailtopicinterest
  • AIConversation 唯一标识“用户 + 私聊/群聊”;AIConversationMessage 保存消息, 每次读写时清理 30 天外记录并裁剪到最近 20 轮。
  • services/preferences.py 提供显式偏好管理和真实 AI 可用时的结构化提取;敏感类别及 密钥样式在入库前拒绝。
  • services/conversations.py 提供历史加载、追加、重置。
  • services/context.py 按固定优先级组装问答上下文。
  • services/erasure.py 实现一次性确认码和事务性“忘记我”,审计只保留随机匿名主体。

app/modules/subscriptions

  • services/schedule_parser.py 实现受控中文语法,不使用不确定的自由文本推断。
  • services/subscriptions.py 管理创建、列表、暂停、恢复、退订、时区和安静时段。
  • services/scanner.py 每分钟领取到期订阅,使用行锁(支持时)与唯一投递键创建投递, 并推进订阅的下一执行时间。
  • services/delivery.py 生成受限 AI 内容并发送;个人订阅仅使用个人上下文,群订阅仅使用 系统/公司规则。HTTP 错误、429、5xx 和飞书业务 code != 0 都是失败。
  • tasks/task_queue/ 仅作为执行适配层;关闭 Celery 时由扫描器直接处理持久化投递。

共享集成

  • FeishuEventService 在验证后、注册事件前解析 tenant_key/open_id;缺失身份不创建数据。
  • FeishuCommandService.handle_text 接收 FeishuPrincipal | None。内部预览接口不构造终端 用户身份,因此只能预览不涉及个人数据的安全路径。
  • 管理员能力由统一权限守卫保护:公司规则、公司/财务/风险/考勤命令、用户管理、群订阅。
  • FeishuClient.send_message 接收租户键和可选稳定 uuid,并把非零业务码转换为可分类失败。
  • 自建应用使用 /auth/v3/tenant_access_token/internal;商店应用先使用最近一次已验证的 app_ticket 获取 app_access_token,再按目标 tenant_key 获取 tenant_access_token。应用令牌按应用缓存,租户令牌按 (app_id, tenant_key) 隔离缓存。
  • 只有通过 webhook 验真或长连接 SDK 验证的 app_ticket 事件可以轮换持久化票据; 环境变量票据仅作为启动兜底,票据和访问令牌不得进入日志、审计或响应。
  • 命令回复、卡片、图片和个人/群订阅都显式携带事件主体的 tenant_key;固定默认群报表 使用 FEISHU_DEFAULT_TENANT_KEY,不得复用另一租户的令牌。
  • AI adapter context 增加每用户 session idnoop 明确返回不可用且不写会话、偏好或记忆。

混合数据库协调与运行编排

  • app/tools/reconcile_platform_schema.py 只面向平台 PostgreSQL。它先生成结构指纹和允许列表 dry-run获取 PostgreSQL advisory transaction lock 后执行一次性基线协调;实际结构、 行数、依赖、权限或版本任一不符即中止。
  • 协调 revision 只补齐当前 metadata 缺失的表、列、索引、外键和唯一语义;除已审核为空且 无依赖的遗留表外不删除远端对象。完成后在同一事务内验证零漂移并登记 Alembic revision。
  • FEISHU_EVENT_TRANSPORT 明确选择 disabled|webhook|long_connection。生产环境开启用户 功能时必须具备相应凭据;长连接使用独立受管进程,不嵌入 API 工作线程。
  • API、scheduler、worker 和长连接进程使用 heartbeat 表报告存活。readiness 同时校验 Alembic head、期望组件 heartbeat、队列和飞书依赖缺失或陈旧时返回 HTTP 503。
  • Compose 提供内部 PostgreSQL默认编排和不覆盖 DATABASE_URL 的外部数据库覆盖方式, 并为持续进程配置重启策略。运行配置只使用本地且不提交的 .env

飞书入站 inbox

  • 已验证的消息事件先写入 FeishuEventReceipt inbox再由持久化扫描任务执行。
  • inbox 保存处理所需的短期规范化载荷、状态、尝试次数、下次尝试时间、租约和最小错误摘要; 成功或最终失败后清除原始载荷,避免长期保存个人消息。
  • 命令事务内通过 reply outbox 捕获唯一的文本、卡片或图片回复意图、目标租户和稳定 UUID 只有命令写入、receipt 成功状态与回复意图一起提交后才允许访问飞书网络。
  • 回复使用独立状态、租约和退避时间。发送失败或发送成功后进程在状态提交前退出时,只以 同一内容和 UUID 重试回复,不重新执行用户命令;终态后清除回复载荷。
  • webhook 在验真、挑战处理和 inbox 入库后立即确认;长连接 SDK 回调复用同一入库路径。
  • 执行器通过唯一事件键、数据库锁、租约 token 和状态条件更新防止并发重复执行。失败进入 有界重试,进程崩溃后由过期租约重新领取;成功事件的后续重复投递只返回已有状态。
  • app_ticket 仍在已验证边界内同步轮换,不把票据写入普通消息 inbox。

数据模型

FeishuUser

  • id, code, tenant_key, open_id, union_id, user_id
  • role (user|admin), status (active|disabled)
  • timezone, quiet_hours_start, quiet_hours_end
  • last_active_at, created_at, updated_at
  • 唯一约束:(tenant_key, open_id)code 全局唯一。

FeishuAdminBootstrapTombstone

  • 仅保存带域隔离的 tenant_key + open_id 不可逆摘要和创建时间。
  • 不保存原始飞书身份、owner、角色或其他画像仅用于阻止已执行“忘记我”的 配置初始管理员在重新联系时被自动重授管理员。

UserPreference

  • id, code, owner_id, category, value, source, created_at, updated_at
  • 唯一约束:(owner_id, category, normalized_value)

AIConversation / AIConversationMessage

  • 会话:id, code, owner_id, chat_type, chat_key, created_at, updated_at
  • 消息:id, conversation_id, role, content, created_at
  • 唯一约束:(owner_id, chat_type, chat_key)

现有模型扩展

  • AIMemoryEntry.owner_id 可空;公司规则 owner_id=NULL,个人规则/记忆必须有 owner。
  • AIMemoryEntry.kind 区分 company_rule|personal_rule|memory,旧显式规则迁为 company_rule/legacy_company,旧无所有者自动记忆归档。
  • 指纹唯一性改为 (owner_id, fingerprint),公司数据使用空 owner 的独立范围。
  • MarketWatchlist.owner_id 可空;旧 actor 暂存为 legacy claim key首次注册时认领 认领前个人查询不可见。

PushSubscription

  • id, code, owner_id, target_type (user|chat), target_id
  • prompt, schedule_type, schedule_config, timezone, next_run_at
  • status, consented_at, last_run_at, created_at, updated_at
  • 群目标仅从当前已验证群事件写入。

PushDelivery

  • id, code, subscription_id, scheduled_for, idempotency_key, message_uuid
  • status, attempt_count, next_attempt_at, provider_message_id
  • last_error, sent_at, created_at, updated_at
  • idempotency_keymessage_uuid 全局唯一。

FeishuAppTicket

  • id, app_id, app_ticket, received_at, updated_at
  • app_id 全局唯一;只保留当前有效票据,不保留票据历史。

FeishuEventReceipt inbox 扩展

  • event_key, source, event_id, message_id, received_at
  • status, payload, attempt_count, next_attempt_at
  • locked_until, lock_token, last_error, processed_at, updated_at
  • reply_payload, reply_status, reply_attempt_count, reply_next_attempt_at
  • reply_locked_until, reply_locked_by, reply_last_error, reply_sent_at
  • 旧 receipt 迁为已成功状态;新消息使用唯一 event_key 保证幂等。

业务流程

身份与权限

flowchart TD
    E["收到事件"] --> V{"来源已验证?"}
    V -- 否 --> R["拒绝且不创建个人数据"]
    V -- 是 --> K{"tenant_key/open_id 完整?"}
    K -- 否 --> R
    K -- 是 --> U["查找或注册普通用户"]
    U --> B["应用 FEISHU_ADMIN_IDENTITIES 引导"]
    B --> S{"用户有效且有权限?"}
    S -- 否 --> D["拒绝并审计"]
    S -- 是 --> C["执行命令"]

初始管理员配置按 tenant_key:open_id 精确匹配,多条配置均可各自完成首次引导。若该 配置身份曾执行“忘记我”,系统在同一删除事务内保留不可逆 bootstrap tombstone再次 联系时只注册为普通用户。飞书管理员命令只读取事件中的结构化 mention open_id内部 API 修改由服务 principal 审计。降级或停用前锁定目标并统计有效管理员,禁止移除最后 一个。

问答与偏好

  1. 验证 principal 并按用户 + chat 建立会话。
  2. 加载启用的公司规则、当前用户个人规则、偏好/兴趣、相关个人记忆和最近历史。
  3. 按固定顺序传给 AI使用由用户与会话派生的 provider session id。
  4. 真实 AI 成功后保存本轮消息,再做一次受限偏好提取;定时任务跳过两步。
  5. 每次读写裁剪到最近 20 轮并删除 30 天前消息。

计划解析

支持:

  • 单次:今天/明天 HH:mmYYYY-MM-DD HH:mm
  • 周期:每天、工作日、每周一至周日、每月 1-31 号
  • 间隔:每隔 N 分钟/小时,折算后不得短于 15 分钟

解析输出 schedule_type + schedule_config + timezone + next_run_at(UTC)。日期不存在、时间 已过、模糊表达、无效 IANA 时区均拒绝。月末没有目标日期时跳到下个有效月份。

扫描、投递与重试

flowchart TD
    S["每分钟扫描 next_run_at<=now"] --> L["领取订阅并创建唯一投递"]
    L --> Q{"用户/订阅/频控/安静时段允许?"}
    Q -- 否 --> P["跳过或延后,并推进计划"]
    Q -- 是 --> G["生成受限内容"]
    G --> F["以 delivery UUID 发送飞书"]
    F --> O{"HTTP 与业务 code 成功?"}
    O -- 是 --> X["标记 sent"]
    O -- 可重试 --> B["1/5/15 分钟后 retry"]
    O -- 最终失败 --> Z["标记 failed"]

投递尝试总计最多四次(初次 + 三次重试)。扫描器重启后继续处理 pending/retry 状态。唯一投递键由 subscription_id + scheduled_for 派生;同一键重复任务 返回已有结果,不再次调用飞书。

入站事件处理

flowchart TD
    F["飞书已验证事件"] --> C{"挑战或 app_ticket"}
    C -- 是 --> S["立即安全处理并确认"]
    C -- 否 --> I["幂等写入 inbox pending"]
    I --> A["立即 ACK"]
    W["持久化扫描器"] --> L["领取租约 processing"]
    L --> E["执行身份、权限和命令"]
    E --> O{"结果"}
    O -- 成功 --> X["原子提交 succeeded 与 reply outbox"]
    O -- 可重试 --> R["retry + 退避时间"]
    O -- 超限 --> Z["failed + 最小错误并清空 payload"]
    X --> Q["独立领取并发送 reply"]
    Q --> Y{"发送结果"}
    Y -- 成功 --> C["reply succeeded 并清空载荷"]
    Y -- 可重试 --> QR["仅重试同一内容与 UUID"]
    QR --> W
    R --> W

忘记我

首次命令仅保存哈希确认码与短期过期时间。确认后在一个事务内删除个人规则/记忆、偏好、 兴趣、会话、订阅及投递中的个人内容,删除身份映射;相关审计 actor/target 替换为随机 匿名标识,并清空目标、请求、响应和请求 ID。若被删除身份仍属于初始管理员配置仅保留 不可逆 bootstrap tombstone 防止重新授予;用户再次联系时按新普通用户注册。删除成功 确认消息不再额外写入包含目标摘要或提供方响应的发送审计。

错误处理

  • 未验证事件、缺身份、停用用户、权限不足:拒绝且只审计最小元数据。
  • AI 未配置或 noop返回“AI 当前不可用”,不创建虚假回答、历史、偏好或记忆。
  • 计划解析失败、额度超限、手填群 id使用可操作示例响应且不部分写入。
  • 飞书 429/5xx/超时及业务非零码:分类为可重试;其他 4xx 为最终失败。
  • 缺少飞书凭据且存在启用订阅readiness 返回 degraded不影响基础 health。
  • 商店应用缺少可用 app_ticket 或默认租户时 readiness 返回 degraded自建应用若启用订阅 涉及多个租户时 readiness 返回 degraded避免把单租户令牌错误用于其他租户。
  • 数据库竞争:依赖唯一约束兜底;冲突后回滚到保存点并读取已存在投递。
  • 混合数据库基线不匹配:协调工具中止并输出不含凭据/数据的结构差异,不自动猜测或 stamp。
  • 期望运行组件无 heartbeat、Alembic 非 head 或事件 transport 不可用readiness 返回 503。
  • 入站命令失败或进程丢失租约:保留短期 inbox payload 并按退避重试;成功或最终失败后清除。
  • 入站回复失败或进程在发送后退出:保留短期 reply payload 并以稳定 UUID 重试,绝不重跑命令。

测试策略

  • 单元测试:身份解析、权限矩阵、敏感偏好过滤、上下文顺序、会话裁剪、计划解析、 安静时段、下一执行时间和重试分类。
  • 服务测试:两租户/两用户隔离、最后管理员保护、群订阅绑定、忘记我级联与匿名审计。
  • 并发/幂等测试:多扫描器只产生一条投递、相同 UUID 不重复发送、业务非零码不成功。
  • API/事件测试:验签失败不注册、内部 API Key 保护、正文 actor 不参与权限。
  • 多租户认证测试A/B 租户令牌缓存隔离、app_ticket 只能由已验证事件更新、文本/卡片/ 图片/订阅均使用目标租户、固定群任务使用默认租户。
  • 迁移测试Alembic head 与元数据一致,旧规则/记忆/自选按既定策略迁移。
  • 回归测试:现有固定群报表调度与现有内部接口继续工作。
  • 基线测试:用远端结构的脱敏快照验证 dry-run 指纹、允许列表、事务回滚和零漂移。
  • 运行测试Compose 契约、transport 配置、Alembic head、组件 heartbeat 与 fail-closed。
  • 入站测试:快速 ACK、失败重试、租约回收、并发重复只执行一次、成功后不重复。