Merge worktree branch merge-15-modules-to-main-5ug5xJ
This commit is contained in:
@@ -1,45 +1,316 @@
|
||||
# msg 工作排期
|
||||
|
||||
> 负责人:ai10
|
||||
> 关联:[workline.md](../workline.md)、[coord.md](../coord.md)、[contracts/msg_contract.md](../contracts/msg_contract.md)
|
||||
> 模式:全并行(各 AI 一口气完成 P2-P6 全部代码,最后统一集成测试)
|
||||
> 关联:[workline.md](../workline.md)、[coord.md](../coord.md)、[contracts/msg_contract.md](../contracts/msg_contract.md)、[02-architecture-design.md §10](../../../services/msg/docs/02-architecture-design.md)
|
||||
> 模式:全并行(各 AI 一口气完成 P2-P6 全部代码,开发期间用 mock,最后统一集成测试)
|
||||
|
||||
---
|
||||
|
||||
## §1 总览
|
||||
|
||||
msg 是消息服务,提供 NotificationService、PreferenceService、TemplateService,并基于 Outbox 模式发布消息事件。全阶段目标:P2 服务骨架+Outbox → P3 三大 Service 实现 → P4-P6 持续优化。
|
||||
msg 是消息通知中台(P5),提供 NotificationService + NotificationPreferenceService + NotificationTemplateService 三服务,基于 Outbox 模式发布通知事件,消费 iam/core-edu/data-ana 共 12 类事件触发多渠道通知。
|
||||
|
||||
**全阶段目标**:
|
||||
- P2-P3:服务骨架补全(schema 迁移 + Outbox + Kafka 基础设施 + mock 消费)
|
||||
- P4:三大 Service 主体实现(Notification + Preference + Template)+ ChannelDispatcher 多渠道
|
||||
- P5:gRPC 50056 启用 + PushGatewayClient gRPC + 12 类事件 consumer + ES mapping
|
||||
- P6:测试覆盖 ≥ 80% + /readyz 硬化 + 黄金模板对齐 + README 修正
|
||||
|
||||
**批次归属**:批次 4(P5),依赖批次 3 content(ai09)就绪后启动,预估 13 天。
|
||||
|
||||
---
|
||||
|
||||
## §2 全阶段甘特图(P2-P6,各 AI 自行细化)
|
||||
## §2 全阶段甘特图(P2-P6)
|
||||
|
||||
```mermaid
|
||||
gantt
|
||||
title ai10 msg 全阶段排期
|
||||
title ai10 msg 全阶段排期(13 天)
|
||||
dateFormat YYYY-MM-DD
|
||||
axisFormat %m-%d
|
||||
|
||||
section P2-P6
|
||||
[阶段任务] :a10a, 2026-07-10, Xd
|
||||
```
|
||||
section P2-P3 骨架补全
|
||||
T1-T2 schema迁移(notifications+preferences字段) :crit, a1, 2026-07-10, 1d
|
||||
T3 新建msg_notification_templates表 :a2, after a1, 1d
|
||||
T4 新建msg_outbox_events+Publisher worker :crit, a3, after a1, 2d
|
||||
T5 新建shared/kafka(producer+consumer骨架) :crit, a4, after a3, 1d
|
||||
T6 引入ioredis+IdempotencyGuard(SETNX) :a5, after a3, 1d
|
||||
|
||||
> **注意**:以上为 coord 初始规划,ai10 接管后必须自行细化为完整 P2-P6 排期。
|
||||
section P4 主体实现
|
||||
T7 ChannelDispatcher多渠道抽象 :crit, a6, after a4, 2d
|
||||
T10 重构notifications.service(移除同步fetch) :a7, after a6, 1d
|
||||
T11 batchMarkAsRead/markAllAsRead/recall/getUnreadCount :a8, after a7, 1d
|
||||
T12 NotificationPreference CRUD :a9, after a7, 1d
|
||||
T13 NotificationTemplate CRUD+render :a10, after a7, 1d
|
||||
T15 createBatch改批量INSERT :a11, after a7, 1d
|
||||
|
||||
section P5 gRPC+事件+ES
|
||||
T8 PushGatewayClient gRPC(替代fetch降级) :crit, a12, after a8, 1d
|
||||
T9 12类Kafka事件consumer(iam/core-edu/data-ana) :crit, a13, after a12, 2d
|
||||
T14 ES索引mapping+ensureIndex+同步 :a14, after a12, 1d
|
||||
gRPC 50056启用+3 Service 13 RPC :crit, a15, after a13, 1d
|
||||
|
||||
section P6 硬化与对齐
|
||||
T16 NotificationsModule补exports :a16, after a15, 1d
|
||||
T17 /readyz多依赖(DB/ES/Redis/Kafka/PushGW) :a17, after a15, 1d
|
||||
T19 统一关闭到LifecycleService :a18, after a15, 1d
|
||||
T20-T21 DB改getDb()+ID改cuid2 :a19, after a15, 1d
|
||||
T22 单元测试覆盖≥80% :crit, a20, after a19, 2d
|
||||
T23 修正README与实现对齐 :a21, after a20, 1d
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## §3 详细任务
|
||||
|
||||
### 全阶段任务
|
||||
### P2-P3 骨架补全
|
||||
|
||||
#### T1-T2:schema 迁移(notifications + preferences 字段扩展)
|
||||
|
||||
- **负责人**:ai10
|
||||
- **交付物**:⚠️ 由 ai10 自行补充
|
||||
- **依赖**:见 [contracts/msg_contract.md](../contracts/msg_contract.md)
|
||||
- **验收标准**:⚠️ 由 ai10 自行补充
|
||||
- **依赖**:无(msg 独占 DB)
|
||||
- **交付物**:
|
||||
- `msg_notifications` 表新增 status / metadata / related_entity_type / related_entity_id / group_id / sender_id / template_id / event_id / updated_at 字段 + 6 个索引
|
||||
- `msg_notification_preferences` 表补齐 created_at / updated_at + 新增 frequency_limit / quiet_hours_start / quiet_hours_end / quiet_hours_timezone 字段
|
||||
- **验收标准**:Drizzle schema 定义更新,迁移脚本可执行,索引符合 02-architecture-design.md §3.1.1 / §3.1.2
|
||||
|
||||
#### T3:新建 msg_notification_templates 表
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T1-T2
|
||||
- **交付物**:`msg_notification_templates` 表 schema(code + type + title_template + content_template + default_channels + variables + locale + status),UNIQUE INDEX `(code, locale)`
|
||||
- **验收标准**:schema 定义 + 迁移脚本,符合 02-architecture-design.md §3.1.3
|
||||
|
||||
#### T4:新建 msg_outbox_events 表 + Outbox Publisher worker
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T1-T2
|
||||
- **交付物**:
|
||||
- `msg_outbox_events` 表 schema(event_id PK + aggregate_type + aggregate_id + event_type + topic + payload + status + retry_count + created_at + published_at + next_retry_at)
|
||||
- `shared/outbox/outbox.publisher.ts`(独立 worker,每 1s 轮询 PENDING 事件投递 Kafka)
|
||||
- `shared/outbox/outbox.schema.ts`(Drizzle schema)
|
||||
- **验收标准**:Outbox Publisher 可轮询 + 投递 + 更新 status=SENT,符合 02-architecture-design.md §3.1.4 / §5.4
|
||||
- **关联 ISSUE**:ISSUE-003(Outbox 强制)
|
||||
|
||||
#### T5:新建 shared/kafka/(producer + consumer 骨架)
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T4
|
||||
- **交付物**:
|
||||
- `shared/kafka/kafka.producer.ts`(idempotent=true + transactionalId=msg-producer)
|
||||
- `shared/kafka/kafka.consumer.ts`(consumer group = msg-service)
|
||||
- `shared/kafka/topic-map.ts`(PRODUCER_TOPIC_MAP + CONSUMER_TOPICS)
|
||||
- **验收标准**:producer 可投递消息,consumer 可订阅 topic,符合 02-architecture-design.md §5.3
|
||||
|
||||
#### T6:引入 ioredis + IdempotencyGuard(SETNX)
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:无
|
||||
- **交付物**:
|
||||
- `shared/redis/redis.client.ts`(ioredis 客户端,连接 REDIS_URL)
|
||||
- `shared/redis/idempotency.guard.ts`(SETNX `msg:processed:{event_id}` TTL 7 天)
|
||||
- `shared/redis/read-bitmap.ts`(已读位图 BITCOUNT / GETBIT)
|
||||
- env.ts 补 REDIS_URL 必填校验
|
||||
- **验收标准**:IdempotencyGuard SETNX 原子去重,Redis 不可用时降级到 DB 唯一索引
|
||||
- **关联 ISSUE**:ISSUE-012(三层幂等防线,补 msg_idempotency 表中间层)
|
||||
|
||||
### P4 主体实现
|
||||
|
||||
#### T7:ChannelDispatcher 多渠道抽象
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T5、T6
|
||||
- **交付物**:
|
||||
- `channels/notification-channel.interface.ts`(NotificationChannel 接口)
|
||||
- `channels/in-app.channel.ts`(站内信,写 MySQL)
|
||||
- `channels/email.channel.ts`(邮件,SMTP,异步队列)
|
||||
- `channels/sms.channel.ts`(短信,HTTP API)
|
||||
- `channels/wechat.channel.ts`(微信,HTTP API)
|
||||
- `channels/push.channel.ts`(推送,调 PushGatewayClient)
|
||||
- `channels/channel-dispatcher.ts`(Promise.allSettled 并行投递 + in_app 总是发送)
|
||||
- **验收标准**:新增渠道只需实现接口 + 注册,符合 02-architecture-design.md §12
|
||||
|
||||
#### T10:重构 notifications.service.ts
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T4、T5、T7
|
||||
- **交付物**:重构 notifications.service.ts,移除同步 fetch push-gateway,改为 ChannelDispatcher + Outbox 事务
|
||||
- **验收标准**:send 方法走 BEGIN TX → INSERT notifications + INSERT outbox → COMMIT → ChannelDispatcher.dispatch
|
||||
|
||||
#### T11:新增端点(batchMarkAsRead / markAllAsRead / recall / getUnreadCount)
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T10
|
||||
- **交付物**:controller + service 新增 4 个端点
|
||||
- **验收标准**:符合 02-architecture-design.md §4.1 REST API 表
|
||||
- **关联 ISSUE**:ISSUE-010(markAsRead 权限改 READ)
|
||||
|
||||
#### T12:NotificationPreference CRUD
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T1-T2
|
||||
- **交付物**:`preferences/` 目录(controller + service + repository + schema + dto)
|
||||
- **验收标准**:GET / PUT preferences 端点可用
|
||||
|
||||
#### T13:NotificationTemplate CRUD + render
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T3
|
||||
- **交付物**:`templates/` 目录(controller + service + repository + schema + dto),含 `{{variable}}` 占位符替换渲染
|
||||
- **验收标准**:CreateTemplate / GetTemplate / ListTemplates / RenderTemplate 可用
|
||||
|
||||
#### T15:createBatch 改批量 INSERT
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T10
|
||||
- **交付物**:createBatch 改为 `db.insert(notifications).values([...])` 批量 INSERT
|
||||
- **验收标准**:1 万条广播通知 < 5s(性能验收)
|
||||
|
||||
### P5 gRPC + 事件 + ES
|
||||
|
||||
#### T8:PushGatewayClient gRPC
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:D5(push-gateway 提供 gRPC PushService.Push)
|
||||
- **交付物**:`shared/push/push-gateway.client.ts`(gRPC 调用,替代 fetch POST /internal/push 降级)
|
||||
- **验收标准**:gRPC 调用 push-gateway PushService.Push,降级模式保留(push-gateway 不可用时走 in_app)
|
||||
- **关联 ISSUE**:ISSUE-005(调用方向澄清)
|
||||
|
||||
#### T9:12 类 Kafka 事件 consumer
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T5、T6、D1-D2(events.proto 补齐 message + 字段)
|
||||
- **交付物**:
|
||||
- `shared/kafka/consumers/iam.consumer.ts`(6 类 user/role 事件)
|
||||
- `shared/kafka/consumers/core-edu.consumer.ts`(5 类 exam/homework/grade/attendance 事件)
|
||||
- `shared/kafka/consumers/data-ana.consumer.ts`(1 类 mastery 事件)
|
||||
- 每个 consumer 走 IdempotencyGuard → NotificationService.createNotificationFromEvent
|
||||
- **验收标准**:12 类事件均可消费 + 幂等去重 + fan-out 通知
|
||||
- **关联 ISSUE**:ISSUE-013(events.proto 缺 4 类 message,阻塞)
|
||||
|
||||
#### T14:ES 索引 mapping + ensureIndex + 同步
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:无
|
||||
- **交付物**:
|
||||
- `config/elasticsearch.ts` 补 `notifications` 索引 mapping(ik_max_word 分词)
|
||||
- ensureIndex 幂等创建
|
||||
- 数据同步:Outbox 事件触发 ES 索引更新(替代当前同步 safeIndex)
|
||||
- **验收标准**:ES 检索可用,mapping 符合 02-architecture-design.md §3.2.1
|
||||
- **关联 ISSUE**:ISSUE-011(降级方向待仲裁)
|
||||
|
||||
#### gRPC 50056 启用 + 3 Service 13 RPC
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:D3-D4(msg.proto 补 RPC + Service)
|
||||
- **交付物**:
|
||||
- `notifications.grpc.controller.ts`(NotificationService gRPC)
|
||||
- `preferences.grpc.controller.ts`(NotificationPreferenceService gRPC)
|
||||
- `templates.grpc.controller.ts`(NotificationTemplateService gRPC)
|
||||
- app.module.ts 注册 gRPC server :50056
|
||||
- **验收标准**:HealthService.Check 返回 SERVING,13 RPC 可调用
|
||||
- **关联 ISSUE**:ISSUE-009(RPC 数量待仲裁 13 vs 17)
|
||||
|
||||
### P6 硬化与对齐
|
||||
|
||||
#### T16:NotificationsModule 补 exports
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:无
|
||||
- **交付物**:notifications.module.ts 补 `exports: [NotificationsService]`
|
||||
- **验收标准**:BFF 可注入 NotificationsService
|
||||
|
||||
#### T17:/readyz 多依赖检查
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:T4、T5、T6、T8
|
||||
- **交付物**:health.controller.ts /readyz 检查 DB / ES / Redis / Kafka producer / Kafka consumer / PushGateway 6 项依赖
|
||||
- **验收标准**:符合 02-architecture-design.md §6.6 判定规则(DB down → down;Redis/ES/Kafka down → degraded)
|
||||
|
||||
#### T19:统一关闭到 LifecycleService
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:无
|
||||
- **交付物**:移除 main.ts 重复 closeDb/closeEs,统一到 LifecycleService,8 步关闭序列
|
||||
- **验收标准**:符合 02-architecture-design.md §6.7 优雅关闭顺序
|
||||
|
||||
#### T20-T21:DB 改 getDb() + ID 改 cuid2
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:无
|
||||
- **交付物**:database.ts 改 getDb() 函数式;service 层 randomUUID → cuid2
|
||||
- **验收标准**:与 classes 黄金模板对齐
|
||||
|
||||
#### T22:单元测试覆盖 ≥ 80%
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:全部 P4-P5 任务
|
||||
- **交付物**:`*.spec.ts`(Service / Repository / ChannelDispatcher / IdempotencyGuard / OutboxPublisher)
|
||||
- **验收标准**:覆盖率 ≥ 80%
|
||||
|
||||
#### T23:修正 README
|
||||
|
||||
- **负责人**:ai10
|
||||
- **依赖**:全部任务
|
||||
- **交付物**:README.md API 表与实现对齐(PUT /:id/read、补 batch/user/:userId/page 端点、补 env 变量表)
|
||||
- **验收标准**:无文档脱节
|
||||
|
||||
---
|
||||
|
||||
## §4 依赖与就绪信号
|
||||
|
||||
- **我依赖**:⚠️ 由 ai10 自行补充(见 contract.md)
|
||||
- **我的就绪信号**:⚠️ 由 ai10 自行补充
|
||||
### 4.1 我依赖的上游就绪标志
|
||||
|
||||
- [ ] **D1**:events.proto 补 UserEvent / RoleEvent / NotificationEvent / MasteryEvent message(coord 维护)— 🔴 阻塞 T9
|
||||
- [ ] **D2**:events.proto GradeEvent 补 class_id;全部事件补 student_ids[](coord 维护)— 🔴 阻塞 T9 fan-out
|
||||
- [ ] **D3**:msg.proto 补 BatchSendNotification / GetUnreadCount / BatchMarkAsRead / MarkAllAsRead / RecallNotification RPC(coord 维护)— 🔴 阻塞 gRPC(依赖 ISSUE-009 仲裁)
|
||||
- [ ] **D4**:msg.proto 补 NotificationPreferenceService + NotificationTemplateService(coord 维护)— 🔴 阻塞 gRPC
|
||||
- [ ] **D5**:push-gateway 提供 gRPC PushService.Push 方法(ai02)— 🔴 阻塞 T8
|
||||
- [ ] **D6**:iam 发布 6 类 user/role 事件(ai06)— 🟡 不阻塞开发(用 mock),阻塞集成验证
|
||||
- [ ] **D7**:core-edu 发布 5 类教学事件(ai08)— 🟡 同上
|
||||
- [ ] **D8**:data-ana 发布 mastery 事件(ai11)— 🟡 同上
|
||||
- [ ] **D9**:Redis 集群可用(infra 部署)— 🟡 降级到 DB 唯一索引
|
||||
|
||||
### 4.2 我的就绪标志(供下游消费)
|
||||
|
||||
- [ ] msg gRPC 50056 启用(HealthService.Check 返回 SERVING)
|
||||
- [ ] NotificationService RPC 可调用(数量待 ISSUE-009 仲裁)
|
||||
- [ ] NotificationPreferenceService RPC 可调用
|
||||
- [ ] NotificationTemplateService RPC 可调用(含 RenderTemplate)
|
||||
- [ ] `edu.notification.*` topic 可发布(供 push-gateway / data-ana 消费,命名待 ISSUE-008 仲裁)
|
||||
- [ ] /readyz 返回 6 项依赖状态
|
||||
- [ ] 测试覆盖率 ≥ 80%
|
||||
|
||||
---
|
||||
|
||||
## §5 Mock 策略
|
||||
|
||||
### 5.1 我提供的 mock(供下游)
|
||||
|
||||
在 msg 真实服务就绪前,为下游(teacher-bff / student-bff / parent-bff / push-gateway)提供 mock:
|
||||
|
||||
- **gRPC mock**:grpc-mock 拦截 50056 端口
|
||||
- ListNotifications 返回固定 10 条未读通知
|
||||
- MarkAsRead 返回 success=true
|
||||
- GetPreference 返回默认偏好(in_app + email 开启,sms + push 关闭)
|
||||
- RenderTemplate 返回固定 title + content
|
||||
- **Kafka mock**:msg 就绪前不发布真实通知事件,push-gateway 使用本地 stub 推送
|
||||
|
||||
### 5.2 我消费的 mock(开发期间)
|
||||
|
||||
在真实上游就绪前,msg 使用以下 mock:
|
||||
|
||||
- **业务事件**:core-edu / data-ana 就绪前,msg 内置定时器发布本地 stub 事件(ExamEvent / HomeworkEvent),触发 mock 通知流程验证 consumer 链路
|
||||
- **用户偏好**:iam 就绪前使用默认偏好(所有用户 in_app 开启)
|
||||
- **模板渲染**:内置 5 个常用模板(exam.published / homework.graded / grade.recorded / mastery.warning / system.notice)
|
||||
- **Push Gateway**:ai02 就绪前用 fetch POST /internal/push 降级(当前实现保留)
|
||||
|
||||
---
|
||||
|
||||
## §6 风险与缓解
|
||||
|
||||
| 风险 | 影响 | 缓解 |
|
||||
| ---- | ---- | ---- |
|
||||
| events.proto 补齐延迟(D1-D2) | T9 consumer 无法验证 | 开发期用 stub 事件,proto 补齐后切换 |
|
||||
| RPC 数量仲裁延迟(ISSUE-009) | gRPC controller 实现范围不确定 | 先实现 13 RPC 基线,仲裁后增减 |
|
||||
| topic 命名仲裁延迟(ISSUE-008) | Kafka producer/consumer topic 不确定 | 开发期用 02-architecture-design.md §5.3 的 TOPIC_MAP,仲裁后统一 |
|
||||
| push-gateway gRPC 延迟(D5) | T8 无法验证 | 保留 fetch POST 降级,gRPC 就绪后切换 |
|
||||
|
||||
Reference in New Issue
Block a user