0b55fbdfe2
Звонки 1:1 и в групповых беседах до 10 участников поверх той же голосовой
инфраструктуры LiveKit: те же токены, что у серверных комнат, но без
серверных ролей — позвонить может участник беседы.
Ручки: POST /channels/{id}/call (создание, в ответе токен и адрес
сигналинга), accept, decline, end, PATCH /call/@me для флагов микрофона,
камеры и шаринга, GET /call (текущий звонок) и GET /calls (история: кто
звонил, когда, сколько длился, кто пропустил). Постороннему — 404,
существование чужой беседы не подтверждаем; в беседе больше десяти
участников звонить нельзя (dm.call_limit); выключенный голос отвечает
voice.disabled.
События: DM_CALL_RING приходит приглашённым, DM_CALL_UPDATE — участникам
при смене состояния, DM_CALL_END — с итогом и причиной. Звонок добавлен в
READY и в карточку беседы (active_call), поэтому «идёт звонок» переживает
перезагрузку. Вебхук LiveKit разбирает комнаты dm_call_* и синхронизирует
состояние, сторож голосовых состояний закрывает гудки без ответа и
опустевшие комнаты.
Тесты: жизненный цикл 1:1, отклонение и пропущенный, отказ постороннему
на всех ручках, предел десяти участников, поведение при выключенном
голосе, события Gateway и active_call в READY.
333 lines
11 KiB
Go
333 lines
11 KiB
Go
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))
|
|
}
|
|
}
|