# msg 工作排期 > 负责人:ai10 > 关联:[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 是消息通知中台(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) ```mermaid gantt title ai10 msg 全阶段排期(13 天) dateFormat YYYY-MM-DD axisFormat %m-%d 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 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 - **依赖**:无(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 依赖与就绪信号 ### 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 就绪后切换 |