|
| 1 | +package server |
| 2 | + |
| 3 | +import ( |
| 4 | + "log/slog" |
| 5 | + "net" |
| 6 | + "strings" |
| 7 | + "time" |
| 8 | + "unicode/utf8" |
| 9 | + |
| 10 | + "github.com/NicolasHaas/gospeak/pkg/datastore" |
| 11 | + "github.com/NicolasHaas/gospeak/pkg/model" |
| 12 | + "github.com/NicolasHaas/gospeak/pkg/protocol/pb" |
| 13 | + "github.com/NicolasHaas/gospeak/pkg/rbac" |
| 14 | +) |
| 15 | + |
| 16 | +const chatPageLimit = 40 // JSON can escape each of 2000 runes as six bytes; 40 messages fit 512 KiB. |
| 17 | + |
| 18 | +func (s *Server) chatChannel(sessionID uint32, channelID int64, st datastore.DataProviderFactory) bool { |
| 19 | + session, ok := s.sessions.GetSnapshot(sessionID) |
| 20 | + if !ok || channelID <= 0 || session.ChannelScope != 0 && session.ChannelScope != channelID { |
| 21 | + return false |
| 22 | + } |
| 23 | + channel, err := st.NonTx().GetChannel(channelID) |
| 24 | + return err == nil && channel != nil |
| 25 | +} |
| 26 | + |
| 27 | +func (s *Server) handleChatMessage(handler *ControlHandler, sessionID uint32, chat *pb.ChatMessage, st datastore.DataProviderFactory, conn net.Conn) { |
| 28 | + session, ok := s.sessions.GetSnapshot(sessionID) |
| 29 | + if !ok { |
| 30 | + return |
| 31 | + } |
| 32 | + channelID := chat.ChannelID |
| 33 | + if channelID == 0 { |
| 34 | + channelID = s.channels.ChannelOf(sessionID) |
| 35 | + } // old clients use their voice channel |
| 36 | + if !s.chatChannel(sessionID, channelID, st) { |
| 37 | + sendError(conn, 3, "channel not found or inaccessible") |
| 38 | + return |
| 39 | + } |
| 40 | + text := sanitizeText(strings.TrimSpace(chat.Text)) |
| 41 | + if text == "" || utf8.RuneCountInString(text) > model.MessageMaxBodyLength { |
| 42 | + return |
| 43 | + } |
| 44 | + var id, timestamp int64 |
| 45 | + if s.cfg.ChatHistoryLimit > 0 { |
| 46 | + message := &model.Message{ChannelID: channelID, SenderID: session.UserID, SenderName: session.Username, Body: text} |
| 47 | + if err := st.NonTx().CreateMessageWithRetention(message, s.cfg.ChatHistoryLimit, s.cfg.ChatMaxAge); err != nil { |
| 48 | + slog.Error("store chat message", "err", err) |
| 49 | + sendError(conn, 3, "could not save message") |
| 50 | + return |
| 51 | + } |
| 52 | + id, timestamp = message.ID, message.CreatedAt.Unix() |
| 53 | + } else { |
| 54 | + timestamp = time.Now().Unix() |
| 55 | + } |
| 56 | + handler.broadcastChat(channelID, &pb.ControlMessage{ChatEvent: &pb.ChatMessage{ |
| 57 | + ID: id, ChannelID: channelID, SenderID: session.UserID, SenderName: session.Username, |
| 58 | + Text: text, Timestamp: timestamp, |
| 59 | + }}, true) |
| 60 | + s.metrics.ChatMessagesSent.Add(1) |
| 61 | +} |
| 62 | + |
| 63 | +func (s *Server) handleChatHistory(handler *ControlHandler, sessionID uint32, req *pb.ChatHistoryRequest, st datastore.DataProviderFactory, conn net.Conn) { |
| 64 | + if req.BeforeID < 0 || req.Limit < 0 || req.Limit > chatPageLimit || !s.chatChannel(sessionID, req.ChannelID, st) { |
| 65 | + sendError(conn, 3, "channel not found or inaccessible") |
| 66 | + return |
| 67 | + } |
| 68 | + // Select before reading history so a concurrent write reaches either the page or live fanout. |
| 69 | + // A message can appear in both; clients should deduplicate by its stored ID. |
| 70 | + handler.mu.Lock() |
| 71 | + if _, registered := handler.connMap[sessionID]; registered { |
| 72 | + handler.chatSelection[sessionID] = req.ChannelID |
| 73 | + } |
| 74 | + handler.mu.Unlock() |
| 75 | + limit := req.Limit |
| 76 | + if limit == 0 { |
| 77 | + limit = chatPageLimit |
| 78 | + } |
| 79 | + if s.cfg.ChatHistoryLimit == 0 { |
| 80 | + if err := writeControlMessage(conn, &pb.ControlMessage{ChatHistoryResp: &pb.ChatHistoryResponse{ChannelID: req.ChannelID, Messages: []pb.ChatMessage{}}}); err != nil { |
| 81 | + slog.Warn("send empty chat history", "session", sessionID, "err", err) |
| 82 | + } |
| 83 | + return |
| 84 | + } |
| 85 | + fetch := limit + 1 |
| 86 | + filters := model.MessageFilters{LimitToChannelID: &req.ChannelID, BeforeID: req.BeforeID, PageSize: &fetch} |
| 87 | + if s.cfg.ChatMaxAge > 0 { |
| 88 | + filters.Since = time.Now().Add(-s.cfg.ChatMaxAge) |
| 89 | + } |
| 90 | + rows, err := st.NonTx().ListMessages(filters) |
| 91 | + if err != nil { |
| 92 | + slog.Error("load chat history", "err", err) |
| 93 | + sendError(conn, 3, "could not load history") |
| 94 | + return |
| 95 | + } |
| 96 | + resp := &pb.ChatHistoryResponse{ChannelID: req.ChannelID, Messages: []pb.ChatMessage{}, HasMore: int64(len(rows)) > limit} |
| 97 | + if resp.HasMore { |
| 98 | + rows = rows[:limit] |
| 99 | + } |
| 100 | + // ponytail: cap legacy rows at current wire limits; stored originals remain untouched. |
| 101 | + for _, m := range rows { |
| 102 | + resp.Messages = append(resp.Messages, pb.ChatMessage{ID: m.ID, ChannelID: m.ChannelID, |
| 103 | + SenderID: m.SenderID, SenderName: truncateRunes(m.SenderName, model.MaxUsernameLength), |
| 104 | + Text: truncateRunes(m.Body, model.MessageMaxBodyLength), Timestamp: m.CreatedAt.Unix()}) |
| 105 | + } |
| 106 | + if err := writeControlMessage(conn, &pb.ControlMessage{ChatHistoryResp: resp}); err != nil { |
| 107 | + slog.Warn("send chat history", "session", sessionID, "err", err) |
| 108 | + } |
| 109 | +} |
| 110 | + |
| 111 | +func (s *Server) handleChatDelete(handler *ControlHandler, sessionID uint32, req *pb.ChatDeleteRequest, st datastore.DataProviderFactory, conn net.Conn) { |
| 112 | + if s.cfg.ChatHistoryLimit == 0 { |
| 113 | + sendError(conn, 3, "chat history disabled") |
| 114 | + return |
| 115 | + } |
| 116 | + if req.MessageID <= 0 || !s.chatChannel(sessionID, req.ChannelID, st) { |
| 117 | + sendError(conn, 3, "channel not found or inaccessible") |
| 118 | + return |
| 119 | + } |
| 120 | + // Role updates use the same lock: no demotion can interleave with this delete. |
| 121 | + s.remoteModerationMu.Lock() |
| 122 | + session, ok := s.sessions.GetSnapshot(sessionID) |
| 123 | + if !ok || !rbac.HasPermission(session.Role, rbac.PermDeleteChatMessage) { |
| 124 | + s.remoteModerationMu.Unlock() |
| 125 | + sendError(conn, 30, "permission denied") |
| 126 | + return |
| 127 | + } |
| 128 | + actor, err := st.NonTx().GetUserByID(session.UserID) |
| 129 | + if err != nil || actor == nil || !rbac.HasPermission(actor.Role, rbac.PermDeleteChatMessage) { |
| 130 | + s.remoteModerationMu.Unlock() |
| 131 | + sendError(conn, 30, "permission denied") |
| 132 | + return |
| 133 | + } |
| 134 | + deleted, err := st.NonTx().DeleteMessageInChannel(req.MessageID, req.ChannelID) |
| 135 | + s.remoteModerationMu.Unlock() |
| 136 | + if err != nil { |
| 137 | + slog.Error("delete chat message", "err", err) |
| 138 | + sendError(conn, 3, "could not delete message") |
| 139 | + return |
| 140 | + } |
| 141 | + if !deleted { |
| 142 | + sendError(conn, 3, "message not found") |
| 143 | + return |
| 144 | + } |
| 145 | + event := &pb.ControlMessage{ChatDeleteEvent: &pb.ChatDeleteEvent{ChannelID: req.ChannelID, MessageID: req.MessageID}} |
| 146 | + if err := writeControlMessage(conn, event); err != nil { |
| 147 | + slog.Warn("send chat deletion", "session", sessionID, "err", err) |
| 148 | + } |
| 149 | + // Old clients reject unknown envelope fields; only history-capable subscribers receive deletion events. |
| 150 | + handler.broadcastChat(req.ChannelID, event, false) |
| 151 | +} |
| 152 | + |
| 153 | +func (handler *ControlHandler) broadcastChat(channelID int64, event *pb.ControlMessage, includeLegacy bool) { |
| 154 | + members := handler.server.channels.Members(channelID) |
| 155 | + legacy := make(map[uint32]bool, len(members)) |
| 156 | + for _, sessionID := range members { |
| 157 | + legacy[sessionID] = true |
| 158 | + } |
| 159 | + handler.mu.RLock() |
| 160 | + clients := make(map[uint32]*controlClient) |
| 161 | + for sessionID, client := range handler.connMap { |
| 162 | + selected, explicit := handler.chatSelection[sessionID] |
| 163 | + if explicit && selected == channelID || includeLegacy && !explicit && legacy[sessionID] { |
| 164 | + // Membership and explicit selection both remain subject to the account's scope. |
| 165 | + if session, ok := handler.server.sessions.GetSnapshot(sessionID); ok && |
| 166 | + (session.ChannelScope == 0 || session.ChannelScope == channelID) { |
| 167 | + clients[sessionID] = client |
| 168 | + } |
| 169 | + } |
| 170 | + } |
| 171 | + handler.mu.RUnlock() |
| 172 | + for sessionID, client := range clients { |
| 173 | + if err := client.send(event); err != nil { |
| 174 | + slog.Error("chat event write failed", "session", sessionID, "err", err) |
| 175 | + } |
| 176 | + } |
| 177 | +} |
| 178 | + |
| 179 | +// ponytail: one bounded batch per minute; old imported backlogs drain over multiple ticks. |
| 180 | +func (s *Server) runChatJanitor(st datastore.DataProviderFactory) { |
| 181 | + ticker := time.NewTicker(time.Minute) |
| 182 | + defer ticker.Stop() |
| 183 | + for { |
| 184 | + select { |
| 185 | + case <-s.ctx.Done(): |
| 186 | + return |
| 187 | + case <-ticker.C: |
| 188 | + if _, err := st.NonTx().PruneExpiredMessages(s.cfg.ChatMaxAge, 1000); err != nil { |
| 189 | + slog.Error("prune expired chat", "err", err) |
| 190 | + } |
| 191 | + } |
| 192 | + } |
| 193 | +} |
0 commit comments