feat(voice): вебхуки LiveKit и сторож голосовых состояний
AGENT.md 7.14: - POST /api/v1/livekit/webhook: подпись проверяется JWT-токеном, подписанным API-секретом, с проверкой sha256 тела (иначе 401), при выключенном голосе — 503; обработка идемпотентна по event.id (повторная доставка → duplicate); - события: participant_joined восстанавливает состояние после рестарта приложения, participant_left и participant_connection_aborted убирают участника, room_finished очищает комнату целиком, track_published и track_unpublished обновляют флаги камеры и демонстрации экрана; каждое изменение рассылает VOICE_STATE_UPDATE участникам сервера; - сторож `StartVoiceWatchdog` раз в 15 минут удаляет состояния без обновлений дольше шести часов (сбои без вебхука) и уведомляет клиентов; - тесты: подпись (чужой секрет, другое тело, отсутствие заголовка), полный цикл событий SFU и идемпотентность повторной доставки.
This commit is contained in:
@@ -1053,3 +1053,105 @@ func TestVoiceStateFlow(t *testing.T) {
|
||||
t.Fatalf("после выхода = %+v", empty.States)
|
||||
}
|
||||
}
|
||||
|
||||
// TestLiveKitWebhook проверяет подпись, идемпотентность и синхронизацию состояний.
|
||||
func TestLiveKitWebhook(t *testing.T) {
|
||||
f := newMessagingFixture(t)
|
||||
f.srv.cfg.LiveKitAPIKey = "testkey"
|
||||
f.srv.cfg.LiveKitAPISecret = "testsecret"
|
||||
f.srv.cfg.LiveKitURL = "wss://gl.test/rtc"
|
||||
f.srv.voice = voice.NewIssuer("testkey", "testsecret", time.Hour)
|
||||
|
||||
voiceRec := doJSON(t, f.srv, http.MethodPost, "/api/v1/guilds/"+f.guildID+"/channels",
|
||||
`{"name":"Голос","type":"voice"}`, f.ownerCookie)
|
||||
voiceChannel := decodeResponse[struct {
|
||||
Channel struct {
|
||||
ID string `json:"id"`
|
||||
} `json:"channel"`
|
||||
}](t, voiceRec)
|
||||
channelID := guildIDOf(t, voiceChannel.Channel.ID)
|
||||
guildID := guildIDOf(t, f.guildID)
|
||||
room := voice.RoomName(guildID, channelID)
|
||||
|
||||
// Событие без подписи отклоняется.
|
||||
noAuth := doJSON(t, f.srv, http.MethodPost, "/api/v1/livekit/webhook",
|
||||
`{"event":"participant_joined"}`, f.ownerCookie)
|
||||
if noAuth.Code != http.StatusUnauthorized {
|
||||
t.Fatalf("webhook без подписи = %d, want 401", noAuth.Code)
|
||||
}
|
||||
|
||||
post := func(body string) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
signature, err := f.srv.voice.IssueWebhookForTest([]byte(body))
|
||||
if err != nil {
|
||||
t.Fatalf("подпись вебхука: %v", err)
|
||||
}
|
||||
req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/api/v1/livekit/webhook",
|
||||
strings.NewReader(body))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Authorization", signature)
|
||||
rec := httptest.NewRecorder()
|
||||
f.srv.Handler().ServeHTTP(rec, req)
|
||||
return rec
|
||||
}
|
||||
|
||||
// Участник «появился» на стороне SFU — состояние восстанавливается.
|
||||
joined := post(`{"event":"participant_joined","id":"evt-1","room":{"name":"` + room + `"},` +
|
||||
`"participant":{"identity":"` + f.memberID + `"}}`)
|
||||
if joined.Code != http.StatusOK {
|
||||
t.Fatalf("participant_joined = %d, body = %s", joined.Code, joined.Body.String())
|
||||
}
|
||||
states := doJSON(t, f.srv, http.MethodGet, "/api/v1/guilds/"+f.guildID+"/voice-states", "", f.ownerCookie)
|
||||
list := decodeResponse[struct {
|
||||
States []struct {
|
||||
UserID string `json:"user_id"`
|
||||
ChannelID string `json:"channel_id"`
|
||||
} `json:"voice_states"`
|
||||
}](t, states)
|
||||
if len(list.States) != 1 || list.States[0].UserID != f.memberID {
|
||||
t.Fatalf("состояния после вебхука = %+v", list.States)
|
||||
}
|
||||
|
||||
// Повторная доставка того же события ничего не меняет.
|
||||
duplicate := post(`{"event":"participant_joined","id":"evt-1","room":{"name":"` + room + `"},` +
|
||||
`"participant":{"identity":"` + f.memberID + `"}}`)
|
||||
var duplicateBody struct {
|
||||
Duplicate bool `json:"duplicate"`
|
||||
}
|
||||
if err := json.Unmarshal(duplicate.Body.Bytes(), &duplicateBody); err != nil {
|
||||
t.Fatalf("decode duplicate: %v", err)
|
||||
}
|
||||
if !duplicateBody.Duplicate {
|
||||
t.Fatalf("повторное событие должно помечаться duplicate: %s", duplicate.Body.String())
|
||||
}
|
||||
|
||||
// Публикация камеры отмечает флаг.
|
||||
camera := post(`{"event":"track_published","id":"evt-2","room":{"name":"` + room + `"},` +
|
||||
`"participant":{"identity":"` + f.memberID + `"},"track":{"source":"CAMERA"}}`)
|
||||
if camera.Code != http.StatusOK {
|
||||
t.Fatalf("track_published = %d, body = %s", camera.Code, camera.Body.String())
|
||||
}
|
||||
cameraState, err := f.srv.store.GetVoiceState(t.Context(), guildID, guildIDOf(t, f.memberID))
|
||||
if err != nil {
|
||||
t.Fatalf("GetVoiceState: %v", err)
|
||||
}
|
||||
if !cameraState.Camera {
|
||||
t.Fatal("флаг камеры не выставлен по вебхуку")
|
||||
}
|
||||
|
||||
// Уход участника убирает состояние.
|
||||
left := post(`{"event":"participant_left","id":"evt-3","room":{"name":"` + room + `"},` +
|
||||
`"participant":{"identity":"` + f.memberID + `"}}`)
|
||||
if left.Code != http.StatusOK {
|
||||
t.Fatalf("participant_left = %d", left.Code)
|
||||
}
|
||||
after := doJSON(t, f.srv, http.MethodGet, "/api/v1/guilds/"+f.guildID+"/voice-states", "", f.ownerCookie)
|
||||
empty := decodeResponse[struct {
|
||||
States []struct {
|
||||
UserID string `json:"user_id"`
|
||||
} `json:"voice_states"`
|
||||
}](t, after)
|
||||
if len(empty.States) != 0 {
|
||||
t.Fatalf("после participant_left = %+v", empty.States)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user