feat: auto committed

This commit is contained in:
SpecialX
2026-07-10 15:05:52 +08:00
parent 9ba368477d
commit 21530dc7f6
3 changed files with 612 additions and 85 deletions

View File

@@ -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 是消息服务,提供 NotificationServicePreferenceServiceTemplateService基于 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 多渠道
- P5gRPC 50056 启用 + PushGatewayClient gRPC + 12 类事件 consumer + ES mapping
- P6测试覆盖 ≥ 80% + /readyz 硬化 + 黄金模板对齐 + README 修正
**批次归属**:批次 4P5依赖批次 3 contentai09就绪后启动预估 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-T2schema 迁移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` 表 schemacode + type + title_template + content_template + default_channels + variables + locale + statusUNIQUE INDEX `(code, locale)`
- **验收标准**schema 定义 + 迁移脚本,符合 02-architecture-design.md §3.1.3
#### T4新建 msg_outbox_events 表 + Outbox Publisher worker
- **负责人**ai10
- **依赖**T1-T2
- **交付物**
- `msg_outbox_events` 表 schemaevent_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-003Outbox 强制)
#### 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 + IdempotencyGuardSETNX
- **负责人**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 主体实现
#### T7ChannelDispatcher 多渠道抽象
- **负责人**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-010markAsRead 权限改 READ
#### T12NotificationPreference CRUD
- **负责人**ai10
- **依赖**T1-T2
- **交付物**`preferences/` 目录controller + service + repository + schema + dto
- **验收标准**GET / PUT preferences 端点可用
#### T13NotificationTemplate CRUD + render
- **负责人**ai10
- **依赖**T3
- **交付物**`templates/` 目录controller + service + repository + schema + dto`{{variable}}` 占位符替换渲染
- **验收标准**CreateTemplate / GetTemplate / ListTemplates / RenderTemplate 可用
#### T15createBatch 改批量 INSERT
- **负责人**ai10
- **依赖**T10
- **交付物**createBatch 改为 `db.insert(notifications).values([...])` 批量 INSERT
- **验收标准**1 万条广播通知 < 5s性能验收
### P5 gRPC + 事件 + ES
#### T8PushGatewayClient gRPC
- **负责人**ai10
- **依赖**D5push-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调用方向澄清
#### T912 类 Kafka 事件 consumer
- **负责人**ai10
- **依赖**T5、T6、D1-D2events.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-013events.proto 缺 4 类 message阻塞
#### T14ES 索引 mapping + ensureIndex + 同步
- **负责人**ai10
- **依赖**:无
- **交付物**
- `config/elasticsearch.ts``notifications` 索引 mappingik_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-D4msg.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 返回 SERVING13 RPC 可调用
- **关联 ISSUE**ISSUE-009RPC 数量待仲裁 13 vs 17
### P6 硬化与对齐
#### T16NotificationsModule 补 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 → downRedis/ES/Kafka down → degraded
#### T19统一关闭到 LifecycleService
- **负责人**ai10
- **依赖**:无
- **交付物**:移除 main.ts 重复 closeDb/closeEs统一到 LifecycleService8 步关闭序列
- **验收标准**:符合 02-architecture-design.md §6.7 优雅关闭顺序
#### T20-T21DB 改 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 messagecoord 维护)— 🔴 阻塞 T9
- [ ] **D2**events.proto GradeEvent 补 class_id全部事件补 student_ids[]coord 维护)— 🔴 阻塞 T9 fan-out
- [ ] **D3**msg.proto 补 BatchSendNotification / GetUnreadCount / BatchMarkAsRead / MarkAllAsRead / RecallNotification RPCcoord 维护)— 🔴 阻塞 gRPC依赖 ISSUE-009 仲裁)
- [ ] **D4**msg.proto 补 NotificationPreferenceService + NotificationTemplateServicecoord 维护)— 🔴 阻塞 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 就绪后切换 |