package server import ( "context" "encoding/json" "io" "log/slog" "net/http" "strings" "time" "github.com/go-chi/chi/v5" "glchat/internal/store" "glchat/internal/voice" ) // livekitWebhookEvent — событие вебхука LiveKit (AGENT.md 7.14). type livekitWebhookEvent struct { Event string `json:"event"` ID string `json:"id"` Room struct { Name string `json:"name"` SID string `json:"sid"` } `json:"room"` Participant struct { Identity string `json:"identity"` SID string `json:"sid"` } `json:"participant"` Track struct { Type string `json:"type"` Source string `json:"source"` } `json:"track"` } // webhookTTL — сколько помним идентификаторы обработанных событий. const webhookTTL = time.Hour // registerVoiceWebhook вешает обработчик вебхуков LiveKit (AGENT.md 7.14). func (s *Server) registerVoiceWebhook(router chi.Router) { router.Post("/livekit/webhook", s.handleLiveKitWebhook) } // handleLiveKitWebhook принимает события SFU: подпись проверяется, обработка // идемпотентна по идентификатору события, состояния синхронизируются с БД. func (s *Server) handleLiveKitWebhook(w http.ResponseWriter, r *http.Request) { if !s.voice.Enabled() { httpxWriteJSONError(w, http.StatusServiceUnavailable, "voice.disabled", "голос не настроен") return } body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) if err != nil { httpxWriteJSONError(w, http.StatusBadRequest, "request.bad", "cannot read body") return } if err := s.voice.VerifyWebhook(r.Header.Get("Authorization"), body); err != nil { s.logger.WarnContext(r.Context(), "livekit webhook rejected", slog.Any("error", err)) httpxWriteJSONError(w, http.StatusUnauthorized, "voice.webhook_unauthorized", "invalid webhook signature") return } var event livekitWebhookEvent if err := json.Unmarshal(body, &event); err != nil { httpxWriteJSONError(w, http.StatusBadRequest, "request.bad", "malformed event") return } if event.ID != "" && !s.markWebhookProcessed(event.ID) { // Повторная доставка: отвечаем успехом, ничего не меняя. httpxWriteJSON(w, http.StatusOK, map[string]any{"ok": true, "duplicate": true}) return } s.applyLiveKitEvent(r.Context(), event) // Пишем в лог: по нему видно, что SFU действительно доставляет события. s.logger.InfoContext(r.Context(), "livekit webhook handled", slog.String("event", event.Event), slog.String("room", event.Room.Name)) httpxWriteJSON(w, http.StatusOK, map[string]any{"ok": true}) } // markWebhookProcessed возвращает false, если событие уже обрабатывалось. func (s *Server) markWebhookProcessed(id string) bool { now := time.Now() s.webhookMu.Lock() defer s.webhookMu.Unlock() if len(s.webhookSeen) > 4096 { cutoff := now.Add(-webhookTTL) for key, at := range s.webhookSeen { if at.Before(cutoff) { delete(s.webhookSeen, key) } } } if _, ok := s.webhookSeen[id]; ok { return false } s.webhookSeen[id] = now return true } // applyLiveKitEvent обновляет голосовые состояния по событию SFU. func (s *Server) applyLiveKitEvent(ctx context.Context, event livekitWebhookEvent) { guildID, channelID, ok := voice.ParseRoomName(event.Room.Name) if !ok { return } switch event.Event { case "participant_joined": userID := parseVoiceIdentity(event.Participant.Identity) if userID == 0 { return } // Восстанавливаем состояние, если приложение перезапускалось: комната // уже существует, значит участник действительно в голосовой комнате. current, err := s.store.GetVoiceState(ctx, guildID, userID) if err == nil && current.ChannelID == channelID { return } params := store.VoiceStateParams{UserID: userID, GuildID: guildID, ChannelID: channelID} if err == nil { params.SelfMute = current.SelfMute params.SelfDeaf = current.SelfDeaf params.ServerMute = current.ServerMute params.ServerDeaf = current.ServerDeaf } state, err := s.store.UpsertVoiceState(ctx, params) if err != nil { s.logger.WarnContext(ctx, "livekit webhook: upsert state failed", slog.Any("error", err)) return } s.dispatchVoiceState(ctx, state) case "participant_left", "participant_connection_aborted": userID := parseVoiceIdentity(event.Participant.Identity) if userID == 0 { return } if err := s.store.DeleteVoiceState(ctx, guildID, userID); err != nil { return } s.dispatchVoiceLeave(ctx, guildID, channelID, userID) case "room_finished": states, err := s.store.ListVoiceStates(ctx, guildID) if err != nil { return } for _, state := range states { if state.ChannelID != channelID { continue } if err := s.store.DeleteVoiceState(ctx, guildID, state.UserID); err != nil { continue } s.dispatchVoiceLeave(ctx, guildID, channelID, state.UserID) } case "track_published", "track_unpublished": userID := parseVoiceIdentity(event.Participant.Identity) if userID == 0 { return } state, err := s.store.GetVoiceState(ctx, guildID, userID) if err != nil || state.ChannelID != channelID { return } published := event.Event == "track_published" params := store.VoiceStateParams{ UserID: userID, GuildID: guildID, ChannelID: channelID, SessionID: state.SessionID, SelfMute: state.SelfMute, SelfDeaf: state.SelfDeaf, ServerMute: state.ServerMute, ServerDeaf: state.ServerDeaf, Camera: state.Camera, Screen: state.Screen, } switch strings.ToUpper(event.Track.Source) { case "CAMERA": params.Camera = published case "SCREEN_SHARE", "SCREEN_SHARE_AUDIO": params.Screen = published default: return } updated, err := s.store.UpsertVoiceState(ctx, params) if err != nil { return } s.dispatchVoiceState(ctx, updated) } } // parseVoiceIdentity превращает identity участника LiveKit в идентификатор. func parseVoiceIdentity(identity string) uint64 { id, err := parseID("identity", identity) if err != nil { return 0 } return id } // startVoiceWatchdog периодически убирает «призрачные» состояния: записи без // обновлений дольше шести часов (сервер мог упасть и не получить вебхук). func (s *Server) startVoiceWatchdog(ctx context.Context, interval time.Duration) { if interval <= 0 { interval = 15 * time.Minute } go func() { ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: s.cleanupStaleVoiceStates(ctx) } } }() } // cleanupStaleVoiceStates удаляет давно не обновлявшиеся голосовые состояния. func (s *Server) cleanupStaleVoiceStates(ctx context.Context) { cutoff := time.Now().UTC().Add(-6 * time.Hour) guilds, err := s.store.ListAllGuilds(ctx) if err != nil { return } removed := 0 for _, guild := range guilds { states, err := s.store.ListVoiceStates(ctx, guild.ID) if err != nil { continue } for _, state := range states { if state.UpdatedAt.After(cutoff) { continue } if err := s.store.DeleteVoiceState(ctx, guild.ID, state.UserID); err != nil { continue } removed++ s.dispatchVoiceLeave(ctx, guild.ID, state.ChannelID, state.UserID) } } if removed > 0 { s.logger.InfoContext(ctx, "stale voice states cleaned", slog.Int("removed", removed)) } }