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) { // Комнаты звонков в беседах живут отдельно от серверных голосовых комнат: // по префиксу имени понятно, куда относить событие (AGENT.md 7.8). if callID, ok := voice.ParseCallRoomName(event.Room.Name); ok { s.applyLiveKitCallEvent(ctx, callID, event) return } 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 } // applyLiveKitCallEvent синхронизирует звонок беседы с состоянием комнаты: // вебхук исправляет расхождения после перезапуска приложения и закрывает // звонок, когда в комнате никого не осталось (AGENT.md 7.14). func (s *Server) applyLiveKitCallEvent(ctx context.Context, callID uint64, event livekitWebhookEvent) { call, err := s.store.GetDMCall(ctx, callID) if err != nil || call.Status == store.DMCallEnded { return } userID := parseVoiceIdentity(event.Participant.Identity) switch event.Event { case "participant_joined": if userID == 0 || dmCallParticipant(call, userID) == nil { return } if _, err := s.store.JoinDMCallParticipant(ctx, callID, userID); err != nil { s.logger.WarnContext(ctx, "livekit webhook: join dm call failed", slog.Any("error", err)) return } if updated, err := s.store.GetDMCall(ctx, callID); err == nil { s.dispatchDMCallUpdate(updated) } case "participant_left", "participant_connection_aborted": if userID == 0 { return } current := dmCallParticipant(call, userID) if current == nil || current.State != store.DMCallJoined { return } reason := store.DMCallReasonCompleted if call.Status == store.DMCallRinging { reason = store.DMCallReasonMissed } if _, err := s.leaveDMCall(ctx, call, userID, store.DMCallLeft, reason); err != nil { s.logger.DebugContext(ctx, "livekit webhook: leave dm call skipped", slog.Any("error", err)) } case "room_finished": reason := store.DMCallReasonCompleted if call.Status == store.DMCallRinging { reason = store.DMCallReasonMissed } if err := s.finishDMCall(ctx, call, reason, nil); err != nil { s.logger.WarnContext(ctx, "livekit webhook: finish dm call failed", slog.Any("error", err)) } case "track_published", "track_unpublished": if userID == 0 { return } current := dmCallParticipant(call, userID) if current == nil || current.State != store.DMCallJoined { return } flags := store.DMCallFlags{ SelfMute: current.SelfMute, SelfDeaf: current.SelfDeaf, Camera: current.Camera, Screen: current.Screen, } published := event.Event == "track_published" switch strings.ToUpper(event.Track.Source) { case "MICROPHONE": flags.SelfMute = !published case "CAMERA": flags.Camera = published case "SCREEN_SHARE", "SCREEN_SHARE_AUDIO": flags.Screen = published default: return } if err := s.store.UpdateDMCallFlags(ctx, callID, userID, flags); err != nil { return } if updated, err := s.store.GetDMCall(ctx, callID); err == nil { s.dispatchDMCallUpdate(updated) } } } // 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) s.expireStaleDMCalls(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)) } }