You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Golang Socket的购物应用一对一聊天功能改造需求

购物应用一对一聊天改造需求

我开发了一款购物应用,用户可以发布供应信息,其他用户能查找并参与这些供应活动。目前已实现基于WebSocket的聊天服务,用于客户和购物者沟通细节,但当前聊天逻辑存在问题:同一供应信息(对应reference字段)下,所有用户会收到广播消息,而实际需要的是一对一独立对话——客户A与购物者的对话要和客户B与该购物者的对话完全隔离,购物者可以查看并回复每个独立对话。

以下是现有代码,需要修改实现消息仅发送给指定接收者:


client.go

type Client struct {
    Conn           *websocket.Conn
    ChatRepository chat.ChatRepository
    Message        chan *Message
    ID             string `json:"id"`
    Reference      string `json:"reference"`
    Username       string `json:"username"`
    Sender         string `json:"sender"`
    Recipient      string `json:"recipient"`
}

type Message struct {
    Content   string `json:"content"`
    Reference string `json:"reference"`
    Sender    string `json:"sender"`
}

func (c *Client) writeMessage() {
    defer func() {
        c.Conn.Close()
    }()

    for {
        message, ok := <-c.Message
        if !ok {
            return
        }
        uuid, err := uuid.NewV4()
        if err != nil {
            log.Fatalf("failed to generate UUID: %v", err)
        }
        chatMessage := chat.ChatMessage{
            ID:     uuid.String(),
            Sender: message.Sender,
            Timestamp: time.Now(),
            Content:   message.Content,
        }
        if c.Sender == message.Sender {
            _, errx := c.ChatRepository.AddMessage(message.Reference, chatMessage)
            if err != nil {
                log.Fatalf("failed to generate UUID: %v", errx)
            }
        }
        c.Conn.WriteJSON(chatMessage)
    }
}

func (c *Client) readMessage(hub *Hub) {
    defer func() {
        hub.Unregister <- c
        c.Conn.Close()
    }()

    for {
        _, m, err := c.Conn.ReadMessage()
        if err != nil {
            if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
                log.Printf("error: %v", err)
            }
            break
        }

        msg := &Message{
            Content:   string(m),
            Reference: c.Reference,
            Sender:    c.Sender,
        }

        hub.Broadcast <- msg
    }
}

hub.go

type Room struct {
    ID      string             `json:"id"`
    Name    string             `json:"name"`
    Clients map[string]*Client `json:"clients"`
}

type Hub struct {
    Rooms      map[string]*Room
    Register   chan *Client
    Unregister chan *Client
    Broadcast  chan *Message
    emmiter    events.Emitter
}

func NewHub(emmiter events.Emitter) *Hub {
    return &Hub{
        Rooms:      make(map[string]*Room),
        Register:   make(chan *Client),
        Unregister: make(chan *Client),
        Broadcast:  make(chan *Message, 5),
        emmiter:    emmiter,
    }
}

func (h *Hub) Run() {
    for {
        select {
        case cl := <-h.Register:
            if _, ok := h.Rooms[cl.Reference]; ok {
                r := h.Rooms[cl.Reference]

                if _, ok := r.Clients[cl.ID]; !ok {
                    r.Clients[cl.ID] = cl
                }

            }
        case cl := <-h.Unregister:
            if _, ok := h.Rooms[cl.Reference]; ok {
                if _, ok := h.Rooms[cl.Reference].Clients[cl.ID]; ok {
                    delete(h.Rooms[cl.Reference].Clients, cl.ID)
                    close(cl.Message)
                }
            }

        case m := <-h.Broadcast:
            if _, ok := h.Rooms[m.Reference]; ok {
                for _, cl := range h.Rooms[m.Reference].Clients {
                    cl.Message <- m
                    if m.Sender != cl.Recipient {
                        notifications.SendPush(h.emmiter, cl.Recipient, fmt.Sprintf("New message from %v", cl.Username), m.Content)
                    }
                }
            }
        }
    }
}

handler.go

type Handler struct {
    hub            *Hub
    chatRepository chat.ChatRepository
}

func NewHandler(h *Hub, chatRepository chat.ChatRepository) *Handler {
    return &Handler{
        hub:            h,
        chatRepository: chatRepository,
    }
}

var upgrader = websocket.Upgrader{
    ReadBufferSize:  1024,
    WriteBufferSize: 1024,
    CheckOrigin: func(r *http.Request) bool {
        return true
    },
}

func (h *Handler) JoinRoom(c *gin.Context) {
    conn, err := upgrader.Upgrade(c.Writer, c.Request, nil)
    if err != nil {
        utils.HandleError(c, nil, "error creating chat connection", http.StatusBadGateway)
        return
    }

    reference := c.Param("reference")
    sender := c.Query("sender")
    username := c.Query("username")
    recipient := c.Query("recipient")

    if reference == "" || sender == "" || username == "" || recipient == "" {
        utils.HandleError(c, nil, "required parameters missing", http.StatusBadGateway)
        return
    }
    if _, ok := h.hub.Rooms[reference]; !ok {
        _, err1 := h.chatRepository.GetChatHistory(reference)
        if err1 != nil {
            log.Printf("Failed to retrieve chat history: %s", err1)
            errx := h.chatRepository.CreateChat(reference)
            if errx != nil {
                utils.HandleError(c, nil, "error storing connection", http.StatusBadGateway)
                return
            }
        }
        h.hub.Rooms[reference] = &Room{
            ID:      reference,
            Name:    sender,
            Clients: make(map[string]*Client),
        }
    }

    cl := &Client{
        Conn:           conn,
        ChatRepository: h.chatRepository,
        Message:        make(chan *Message, 10),
        ID:             sender,
        Reference:      reference,
        Sender:         sender,
        Username:       username,
        Recipient:      recipient,
    }

    h.hub.Register <- cl
    go cl.writeMessage()
    cl.readMessage(h.hub)
}

routes.go

hub := ws.NewHub(events.NewEventEmitter(conn))
wsHandler := ws.NewHandler(hub, pr.NewChatRepository(db, client))
go hub.Run()
v1.GET("/chat/ws/:reference", g.Guard([]string{"user", "admin", "dispatcher"}, nil), wsHandler.JoinRoom)

chat.model.go

type Chat struct {
    ID        string        `json:"id,omitempty" bson:"_id,omitempty"`
    Reference string        `json:"reference" bson:"reference"`
    Messages  []ChatMessage `json:"messages" bson:"messages"`
}
type ChatMessage struct {
    ID        string    `json:"id,omitempty" bson:"_id,omitempty"`
    Sender    string    `json:"sender" bson:"sender,omitempty"`
    Timestamp time.Time `json:"timestamp" bson:"timestamp,omitempty"`
    Content   string    `json:"content" bson:"content,omitempty"`
}

改造方案

1. 扩展消息结构体,携带接收者信息

修改client.go中的Message,新增Recipient字段,让消息明确知道要发给谁:

type Message struct {
    Content   string `json:"content"`
    Reference string `json:"reference"`
    Sender    string `json:"sender"`
    Recipient string `json:"recipient"` // 新增:消息接收者ID
}

同时在readMessage方法中,构造消息时带上当前客户端的Recipient:

msg := &Message{
    Content:   string(m),
    Reference: c.Reference,
    Sender:    c.Sender,
    Recipient: c.Recipient, // 补充接收者ID
}

2. 修改Hub的消息分发逻辑,停止广播改为定向发送

修改hub.go中Run方法的Broadcast分支,不再遍历房间所有客户端,而是根据消息的Recipient找到对应客户端发送,同时也给发送者自己发送(用于消息回执):

case m := <-h.Broadcast:
    if room, ok := h.Rooms[m.Reference]; ok {
        // 发送给目标接收者
        if recipientClient, ok := room.Clients[m.Recipient]; ok {
            recipientClient.Message <- m
            // 推送通知只发给接收者
            notifications.SendPush(h.emmiter, m.Recipient, fmt.Sprintf("New message from %v", m.Sender), m.Content)
        }
        // 发送给消息发送者(可选,用于确认消息已发出)
        if senderClient, ok := room.Clients[m.Sender]; ok {
            senderClient.Message <- m
        }
    }

3. 扩展消息模型,存储接收者信息

修改chat.model.go中的ChatMessage,新增Recipient字段,确保聊天历史能记录对话双方:

type ChatMessage struct {
    ID        string    `json:"id,omitempty" bson:"_id,omitempty"`
    Sender    string    `json:"sender" bson:"sender,omitempty"`
    Recipient string    `json:"recipient" bson:"recipient,omitempty"` // 新增
    Timestamp time.Time `json:"timestamp" bson:"timestamp,omitempty"`
    Content   string    `json:"content" bson:"content,omitempty"`
}

4. 调整消息存储逻辑,关联接收者

修改client.go的writeMessage方法,构造ChatMessage时带上Recipient:

chatMessage := chat.ChatMessage{
    ID:        uuid.String(),
    Sender:    message.Sender,
    Recipient: message.Recipient, // 补充接收者
    Timestamp: time.Now(),
    Content:   message.Content,
}

5. 优化房间初始化逻辑(可选)

当前handler.go中房间名称设置为sender不合理,建议改为供应信息相关名称,比如:

h.hub.Rooms[reference] = &Room{
    ID:      reference,
    Name:    fmt.Sprintf("Supply_%s", reference), // 改为供应标识名称
    Clients: make(map[string]*Client),
}

内容的提问来源于stack exchange,提问作者King

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 11:08:10