Files
glchat/internal/server/api_voice_webhook.go
T

242 lines
7.4 KiB
Go
Raw Normal View History

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)
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))
}
}