Files
rogee 0dabb8cfa5 docs: 整理文档目录结构 — 清理过时文档、归集功能子目录、统一命名规范
清理:
- 删除 34 份过时文档(gap reports/QA临时报告/验收报告/阶段性文档)
- 删除 docs/.hermes/skills 第三方 skills 副本(16 文件)
- 删除 skills-lock.json

目录归集:
- 根目录仅保留 README.md 索引
- product/ — 产品与架构设计(PRD + ARCHITECTURE + P2设计文档 + AI/企业路线图)
- tracking/ — Chatwoot parity 开发跟踪
- requirements/ — M01-M12 模块需求
- plans/ — 历史实现计划
- parity/ — 路由 parity 与前端契约
- qa/ — QA 报告与测试计划
- ops/ — 运维部署

命名规范:
- 全小写 kebab-case,禁止全大写文件名
- product/tracking/ops 用 NN- 序号前缀
- requirements 用 MNN- 两位零填充模块号
- plans/qa 用 YYYY-MM-DD- 日期前缀
- requirements M1-M9 零填充为 M01-M09(修复字典序)

同步更新:
- backend/cmd/route_parity/main.go 路径默认值
- backend/scripts/parity_frontend_smoke.sh 报告路径
- 所有 docs 内部交叉引用
- .gitignore 排除编译产物 (backend/gochat, backend/route_parity)
- 新增迁移 000052/000053
- 前端 WS 相关修改
2026-07-09 14:53:27 +08:00

385 lines
11 KiB
Go

package ws
import (
"context"
"encoding/json"
"fmt"
"strings"
"github.com/ThreeDotsLabs/watermill"
"github.com/ThreeDotsLabs/watermill-redisstream/pkg/redisstream"
"github.com/ThreeDotsLabs/watermill/message"
"github.com/redis/go-redis/v9"
wspkg "github.com/gochat/gochat/internal/ws"
"github.com/gochat/gochat/pkg/logger"
)
// Subscriber bridges Watermill PubSub events to WebSocket push notifications.
// It subscribes to relevant Watermill topics (message, conversation, contact events)
// and forwards the events to the appropriate WebSocket rooms via the Hub.
//
// Architecture mapping:
//
// Chatwoot ActionCable broadcasts → Watermill subscriber → Hub.SendToAccount/SendToConversationJSON
// Each GoChat instance runs its own subscriber, so events are delivered to
// locally-connected WebSocket clients. Multi-instance delivery relies on Redis
// PubSub (each instance receives the event and pushes to its own clients).
//
// Reference: P2E §3 — Real-time communication via Redis Pub/Sub
type Subscriber struct {
hub *Hub
subscriber *redisstream.Subscriber
router *message.Router
redisClient redis.UniversalClient
}
// NewSubscriber creates a PubSub-to-WebSocket bridge subscriber.
func NewSubscriber(hub *Hub, redisClient redis.UniversalClient) (*Subscriber, error) {
loggerAdapter := &watermillZapAdapter{}
wmSubscriber, err := redisstream.NewSubscriber(
redisstream.SubscriberConfig{
Client: redisClient,
ConsumerGroup: "ws-gateway-" + watermill.NewUUID(),
},
loggerAdapter,
)
if err != nil {
return nil, fmt.Errorf("failed to create watermill subscriber: %w", err)
}
router, err := message.NewRouter(message.RouterConfig{}, loggerAdapter)
if err != nil {
return nil, fmt.Errorf("failed to create watermill router: %w", err)
}
s := &Subscriber{
hub: hub,
subscriber: wmSubscriber,
router: router,
redisClient: redisClient,
}
// Register handler functions for each event topic
s.registerHandlers()
return s, nil
}
// registerHandlers sets up Watermill router handlers that forward PubSub events
// to WebSocket clients via the Hub.
func (s *Subscriber) registerHandlers() {
// --- Message events → Account room + Conversation room ---
// Using AddNoPublisherHandler because these handlers push directly to WS, not to an output topic.
s.router.AddNoPublisherHandler(
"ws-message-created-handler",
"gochat.message.created",
s.subscriber,
s.forwardToAccountAndConversation(EventMessageCreated),
)
s.router.AddNoPublisherHandler(
"ws-message-updated-handler",
"gochat.message.updated",
s.subscriber,
s.forwardToAccountAndConversation(EventMessageUpdated),
)
s.router.AddNoPublisherHandler(
"ws-message-deleted-handler",
"gochat.message.deleted",
s.subscriber,
s.forwardToAccountAndConversation(EventMessageDeleted),
)
// --- Conversation events → Account room ---
s.router.AddNoPublisherHandler(
"ws-conversation-created-handler",
"gochat.conversation.created",
s.subscriber,
s.forwardToAccount(EventConversationCreated),
)
s.router.AddNoPublisherHandler(
"ws-conversation-updated-handler",
"gochat.conversation.updated",
s.subscriber,
s.forwardToAccountAndConversation(EventConversationUpdated),
)
s.router.AddNoPublisherHandler(
"ws-conversation-resolved-handler",
"gochat.conversation.resolved",
s.subscriber,
s.forwardToAccountAndConversation(EventConversationResolved),
)
s.router.AddNoPublisherHandler(
"ws-conversation-assigned-handler",
"gochat.conversation.assigned",
s.subscriber,
s.forwardToAccountAndConversation(EventConversationAssigned),
)
// --- Contact events → Account room ---
s.router.AddNoPublisherHandler(
"ws-contact-created-handler",
"gochat.contact.created",
s.subscriber,
s.forwardToAccount(EventContactCreated),
)
s.router.AddNoPublisherHandler(
"ws-contact-updated-handler",
"gochat.contact.updated",
s.subscriber,
s.forwardToAccount(EventContactUpdated),
)
// --- Agent typing events → Conversation room ---
s.router.AddNoPublisherHandler(
"ws-agent-typing-on-handler",
"gochat.agent.typing_on",
s.subscriber,
s.forwardToConversation(EventAgentTypingOn),
)
s.router.AddNoPublisherHandler(
"ws-agent-typing-off-handler",
"gochat.agent.typing_off",
s.subscriber,
s.forwardToConversation(EventAgentTypingOff),
)
// --- Inbox events → Account room ---
s.router.AddNoPublisherHandler(
"ws-inbox-created-handler",
"gochat.inbox.created",
s.subscriber,
s.forwardToAccount(EventInboxCreated),
)
s.router.AddNoPublisherHandler(
"ws-inbox-updated-handler",
"gochat.inbox.updated",
s.subscriber,
s.forwardToAccount(EventInboxUpdated),
)
// --- Notification events → Account room (P4 M8 Notification+Webhook) ---
s.router.AddNoPublisherHandler(
"ws-notification-created-handler",
"gochat.notification.created",
s.subscriber,
s.forwardToAccount(EventNotificationCreated),
)
s.router.AddNoPublisherHandler(
"ws-notification-updated-handler",
"gochat.notification.updated",
s.subscriber,
s.forwardToAccount(EventNotificationUpdated),
)
s.router.AddNoPublisherHandler(
"ws-notification-deleted-handler",
"gochat.notification.deleted",
s.subscriber,
s.forwardToAccount(EventNotificationDeleted),
)
// --- Account cache events → Account room ---
s.router.AddNoPublisherHandler(
"ws-account-cache-invalidated-handler",
"gochat.account.cache_invalidated",
s.subscriber,
s.forwardToAccount(EventAccountCacheInvalidated),
)
}
// forwardToAccountAndConversation creates a handler that pushes an event to both
// the account room and the conversation room (if the payload contains conversation_id).
// Used for message events and conversation update events that are relevant at both levels.
func (s *Subscriber) forwardToAccountAndConversation(eventType string) func(msg *message.Message) error {
return func(msg *message.Message) error {
payload, err := extractPayload(msg)
if err != nil {
logger.L().Errorf("ws: failed to extract payload for %s: %v", eventType, err)
return nil // don't retry — bad payload
}
accountID := payload.AccountID
if accountID == 0 {
logger.L().Warnf("ws: event %s missing account_id", eventType)
return nil
}
// Build the WSMessage — the frontend's onReceived handler expects
// { event: "...", data: {...} } which matches WSMessage serialization.
// Hub.SendToAccount wraps this in ActionCable format:
// { identifier: "...", message: { event: "...", data: {...} } }
wsMsg := wspkg.WSMessage{
Event: eventType,
Data: payload.Data,
}
msgData, err := json.Marshal(wsMsg)
if err != nil {
logger.L().Errorf("ws: failed to marshal event for %s: %v", eventType, err)
return nil
}
// Push to account room
s.hub.SendToAccount(accountID, msgData)
// Also push to conversation room if conversation_id is present
if payload.ConversationID > 0 {
s.hub.SendToAccountConversation(accountID, payload.ConversationID, msgData)
}
return nil
}
}
// forwardToAccount creates a handler that pushes an event only to the account room.
// Used for conversation/contact/inbox creation events that don't target a specific conversation.
func (s *Subscriber) forwardToAccount(eventType string) func(msg *message.Message) error {
return func(msg *message.Message) error {
payload, err := extractPayload(msg)
if err != nil {
logger.L().Errorf("ws: failed to extract payload for %s: %v", eventType, err)
return nil
}
accountID := payload.AccountID
if accountID == 0 {
logger.L().Warnf("ws: event %s missing account_id", eventType)
return nil
}
wsMsg := wspkg.WSMessage{
Event: eventType,
Data: payload.Data,
}
msgData, err := json.Marshal(wsMsg)
if err != nil {
logger.L().Errorf("ws: failed to marshal event for %s: %v", eventType, err)
return nil
}
s.hub.SendToAccount(accountID, msgData)
return nil
}
}
// forwardToConversation creates a handler that pushes an event only to the conversation room.
// Used for typing events that only need to reach clients viewing that conversation.
func (s *Subscriber) forwardToConversation(eventType string) func(msg *message.Message) error {
return func(msg *message.Message) error {
payload, err := extractPayload(msg)
if err != nil {
logger.L().Errorf("ws: failed to extract payload for %s: %v", eventType, err)
return nil
}
accountID := payload.AccountID
conversationID := payload.ConversationID
if accountID == 0 || conversationID == 0 {
logger.L().Warnf("ws: event %s missing account_id or conversation_id", eventType)
return nil
}
wsMsg := wspkg.WSMessage{
Event: eventType,
Data: payload.Data,
}
msgData, err := json.Marshal(wsMsg)
if err != nil {
logger.L().Errorf("ws: failed to marshal event for %s: %v", eventType, err)
return nil
}
s.hub.SendToAccountConversation(accountID, conversationID, msgData)
return nil
}
}
// --- Event Payload Extraction ---
// EventPayload represents the common structure of PubSub event messages.
// Watermill message payloads are JSON-encoded with account_id, conversation_id,
// and an arbitrary data field.
type EventPayload struct {
AccountID uint `json:"account_id"`
ConversationID uint `json:"conversation_id,omitempty"`
Data interface{} `json:"data"`
}
// extractPayload parses a Watermill message payload into an EventPayload.
func extractPayload(msg *message.Message) (*EventPayload, error) {
var payload EventPayload
if err := json.Unmarshal(msg.Payload, &payload); err != nil {
return nil, fmt.Errorf("json unmarshal failed: %w", err)
}
return &payload, nil
}
// --- Watermill Zap Logger Adapter ---
// Reuses the same pattern as pubsub/event_bus.go's zapLoggerAdapter.
type watermillZapAdapter struct{}
func (l *watermillZapAdapter) Error(msg string, err error, fields watermill.LogFields) {
if err != nil {
logger.L().Errorf("watermill-ws: %s err=%v fields=%v", msg, err, fields)
} else {
logger.L().Errorf("watermill-ws: %s fields=%v", msg, fields)
}
}
func (l *watermillZapAdapter) Info(msg string, fields watermill.LogFields) {
logger.L().Infof("watermill-ws: %s fields=%v", msg, fields)
}
func (l *watermillZapAdapter) Debug(msg string, fields watermill.LogFields) {
// Skip debug logs to reduce noise
}
func (l *watermillZapAdapter) Trace(msg string, fields watermill.LogFields) {
// Skip trace logs
}
func (l *watermillZapAdapter) With(fields watermill.LogFields) watermill.LoggerAdapter {
return l
}
// --- Lifecycle Management ---
// Run starts the Watermill router (blocking). Should be called in a goroutine.
func (s *Subscriber) Run(ctx context.Context) error {
logger.L().Info("ws: starting PubSub subscriber router")
return s.router.Run(ctx)
}
// Close gracefully shuts down the Watermill router and subscriber.
func (s *Subscriber) Close() error {
logger.L().Info("ws: shutting down PubSub subscriber")
if err := s.router.Close(); err != nil {
logger.L().Errorf("ws: router close error: %v", err)
return err
}
if err := s.subscriber.Close(); err != nil {
if !strings.Contains(err.Error(), "already closed") {
logger.L().Errorf("ws: subscriber close error: %v", err)
return err
}
}
return nil
}
// Running returns whether the router is currently running.
func (s *Subscriber) Running() bool {
return s.router.IsRunning()
}