diff --git a/internal/gateway/gateway.go b/internal/gateway/gateway.go index 0ce0383..a3824b3 100644 --- a/internal/gateway/gateway.go +++ b/internal/gateway/gateway.go @@ -150,6 +150,34 @@ type ReadyChannel struct { MemberCount int `json:"member_count,omitempty"` MemberIDs []string `json:"member_ids,omitempty"` OwnerID string `json:"owner_id,omitempty"` + // ActiveCall — звонок, идущий в личной беседе (AGENT.md 7.8, Фаза 7): + // список бесед и шапка показывают «идёт звонок» сразу после подключения, + // не дожидаясь события DM_CALL_UPDATE. + ActiveCall *ReadyDMCall `json:"active_call,omitempty"` +} + +// ReadyDMCall — звонок в личной беседе в снапшоте. Форма совпадает с +// dmCallPayload REST-ручек и событиями DM_CALL_*, чтобы клиент разбирал +// состояние звонка одним кодом (AGENT.md 8.3). +type ReadyDMCall struct { + ID string `json:"id"` + ChannelID string `json:"channel_id"` + InitiatorID string `json:"initiator_id"` + Status string `json:"status"` + Media string `json:"media"` + CreatedAt string `json:"created_at"` + StartedAt string `json:"started_at,omitempty"` + Participants []ReadyDMCallParticipant `json:"participants"` +} + +// ReadyDMCallParticipant — участник звонка в снапшоте. +type ReadyDMCallParticipant struct { + UserID string `json:"user_id"` + State string `json:"state"` + SelfMute bool `json:"self_mute"` + SelfDeaf bool `json:"self_deaf"` + Camera bool `json:"camera"` + Screen bool `json:"screen"` } type ReadyRole struct { diff --git a/internal/gateway/ready.go b/internal/gateway/ready.go index 4ebb57e..4c2a701 100644 --- a/internal/gateway/ready.go +++ b/internal/gateway/ready.go @@ -112,6 +112,11 @@ func (s *Snapshot) Build(ctx context.Context, user *store.User) (*Ready, error) // У 1:1 иконка беседы — аватар собеседника (AGENT.md 7.8). channel.IconFileID = formatID(*summary.AvatarFileID) } + // Звонок в беседе: клиент рисует маркер «идёт звонок» сразу после + // подключения, а не только по событию (AGENT.md 7.8). + if call, err := s.store.GetLiveDMCall(ctx, summary.Channel.ID); err == nil { + channel.ActiveCall = readyDMCall(call) + } if summary.LastMessageID != 0 { channel.LastMessageID = formatID(summary.LastMessageID) } @@ -427,6 +432,34 @@ func optionalFormatID(id *uint64) string { return formatID(*id) } +// readyDMCall переводит звонок из хранилища в форму снапшота: та же форма, что +// у событий DM_CALL_* и REST-ручек звонков (AGENT.md 8.3). +func readyDMCall(call *store.DMCallWithParticipants) *ReadyDMCall { + payload := &ReadyDMCall{ + ID: formatID(call.ID), + ChannelID: formatID(call.ChannelID), + InitiatorID: formatID(call.InitiatorID), + Status: call.Status, + Media: call.Media, + CreatedAt: call.CreatedAt.UTC().Format(time.RFC3339), + Participants: make([]ReadyDMCallParticipant, 0, len(call.Participants)), + } + if call.StartedAt != nil { + payload.StartedAt = call.StartedAt.UTC().Format(time.RFC3339) + } + for _, participant := range call.Participants { + payload.Participants = append(payload.Participants, ReadyDMCallParticipant{ + UserID: formatID(participant.UserID), + State: participant.State, + SelfMute: participant.SelfMute, + SelfDeaf: participant.SelfDeaf, + Camera: participant.Camera, + Screen: participant.Screen, + }) + } + return payload +} + func formatID(id uint64) string { if id == 0 { return "" diff --git a/internal/server/api_dm_calls.go b/internal/server/api_dm_calls.go new file mode 100644 index 0000000..80d65db --- /dev/null +++ b/internal/server/api_dm_calls.go @@ -0,0 +1,702 @@ +package server + +import ( + "context" + "errors" + "log/slog" + "net/http" + "time" + + "github.com/danielgtaylor/huma/v2" + + "glchat/internal/store" + "glchat/internal/voice" +) + +// Звонки в личных и групповых беседах (AGENT.md 7.8, 7.14). Голос тот же, что +// в серверных комнатах — те же токены LiveKit и те же лимиты, — но серверных +// ролей здесь нет: позвонить может участник беседы, приглашение приходит +// событием DM_CALL_RING, а состояние звонка живёт в dm_calls. +const ( + // maxDMCallParticipants — предел участников звонка (AGENT.md 7.8). + maxDMCallParticipants = 10 + // dmCallRingTimeout — сколько звоним, прежде чем звонок станет пропущенным. + dmCallRingTimeout = time.Minute + // dmCallLimitWindow и dmCallLimitCount — не больше десяти новых звонков в + // минуту на пользователя: гудки не должны превращаться в спам. + dmCallLimitWindow = time.Minute + dmCallLimitCount = 10 +) + +type dmCallParticipantPayload struct { + UserID string `json:"user_id"` + State string `json:"state"` + SelfMute bool `json:"self_mute"` + SelfDeaf bool `json:"self_deaf"` + Camera bool `json:"camera"` + Screen bool `json:"screen"` + JoinedAt string `json:"joined_at,omitempty"` + LeftAt string `json:"left_at,omitempty"` +} + +// dmCallPayload — звонок для API и Gateway: состояние, состав и длительность. +type dmCallPayload struct { + ID string `json:"id"` + ChannelID string `json:"channel_id"` + InitiatorID string `json:"initiator_id"` + Status string `json:"status"` + Media string `json:"media"` + CreatedAt string `json:"created_at"` + StartedAt string `json:"started_at,omitempty"` + EndedAt string `json:"ended_at,omitempty"` + EndedBy string `json:"ended_by,omitempty"` + EndReason string `json:"end_reason,omitempty"` + DurationSeconds int `json:"duration_seconds"` + Participants []dmCallParticipantPayload `json:"participants"` +} + +type dmCallOutput struct { + Body struct { + Call dmCallPayload `json:"call"` + } +} + +// dmLiveCallOutput — текущий звонок беседы: null, если никто не звонит. +type dmLiveCallOutput struct { + Body struct { + Call *dmCallPayload `json:"call"` + } +} + +// dmCallJoinOutput — ответ на создание и принятие звонка: токен LiveKit и адрес +// сигналинга, как у входа в серверную голосовую комнату (AGENT.md 7.14). +type dmCallJoinOutput struct { + Body struct { + Call dmCallPayload `json:"call"` + Token string `json:"token"` + URL string `json:"url"` + Room string `json:"room"` + } +} + +type dmCallHistoryOutput struct { + Body struct { + Calls []dmCallPayload `json:"calls"` + } +} + +// registerDMCallRoutes описывает звонки бесед: начало, принятие, отклонение, +// выход, флаги микрофона и камеры, текущий звонок и историю. +func (s *Server) registerDMCallRoutes(api huma.API) { + security := []map[string][]string{{"sessionCookie": {}}, {"bearerAuth": {}}} + + huma.Register(api, huma.Operation{ + OperationID: "startDirectCall", + Method: http.MethodPost, + Path: "/channels/{channel_id}/call", + Summary: "Позвонить в личной беседе (выдаёт токен LiveKit)", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + Body struct { + Media string `json:"media,omitempty" enum:"audio,video"` + } + }, + ) (*dmCallJoinOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + // Новые звонки ограничены: гудки не должны превращаться в спам. + if !user.IsInstanceAdmin { + if allowed, retryAfter := s.dmCallLimiter.Allow(dmCallUserKey(user.ID)); !allowed { + return nil, rateLimitedError(retryAfter) + } + } + clientURL, err := s.requireVoice() + if err != nil { + return nil, err + } + media := input.Body.Media + if media == "" { + media = dmCallMediaAudio + } + if media != dmCallMediaAudio && media != dmCallMediaVideo { + return nil, humaErrorStatus(http.StatusUnprocessableEntity, "validation.failed", "unknown call media") + } + // Кто получит приглашение: все участники беседы, кроме звонящего. + // Заблокировавшие его не приглашаются — блокировка запрещает личные + // сообщения в обе стороны (AGENT.md 7.2). + participants, err := s.store.DMParticipantProfiles(ctx, channel.ID, user.ID) + if err != nil { + return nil, humaError(err) + } + if len(participants) > maxDMCallParticipants { + return nil, humaErrorStatus(http.StatusConflict, "dm.call_limit", "too many participants for a call") + } + invitees := make([]uint64, 0, len(participants)) + for _, participant := range participants { + if participant.UserID == user.ID { + continue + } + blocked, err := s.store.IsBlocked(ctx, user.ID, participant.UserID) + if err != nil { + return nil, humaError(err) + } + if blocked { + continue + } + invitees = append(invitees, participant.UserID) + } + call, err := s.store.CreateDMCall(ctx, channel.ID, user.ID, media, invitees) + if err != nil { + if errors.Is(err, store.ErrConflict) { + return nil, humaErrorStatus(http.StatusConflict, "dm.call_in_progress", "a call is already running") + } + return nil, humaError(err) + } + token, err := s.issueDMCallToken(user, call.ID) + if err != nil { + return nil, err + } + s.dispatchDMCallStart(call, user.ID) + output := &dmCallJoinOutput{} + output.Body.Call = s.dmCallPayload(call) + output.Body.Token = token + output.Body.URL = clientURL + output.Body.Room = voice.CallRoomName(call.ID) + return output, nil + }) + + huma.Register(api, huma.Operation{ + OperationID: "acceptDirectCall", + Method: http.MethodPost, + Path: "/channels/{channel_id}/call/accept", + Summary: "Принять звонок в личной беседе", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + }, + ) (*dmCallJoinOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + clientURL, err := s.requireVoice() + if err != nil { + return nil, err + } + call, err := s.store.GetLiveDMCall(ctx, channel.ID) + if err != nil { + return nil, humaErrorStatus(http.StatusNotFound, "dm.call_not_found", "no call is running") + } + invited, err := s.store.IsDMCallParticipant(ctx, call.ID, user.ID) + if err != nil { + return nil, humaError(err) + } + if !invited { + return nil, humaErrorStatus(http.StatusForbidden, "perm.denied", "you were not invited to this call") + } + // Лимит десяти участников: одиннадцатому в комнате места нет. + if joined, err := s.store.CountDMCallParticipants(ctx, call.ID, store.DMCallJoined); err != nil { + return nil, humaError(err) + } else if joined >= maxDMCallParticipants { + return nil, humaErrorStatus(http.StatusConflict, "dm.call_limit", "the call is full") + } + if _, err := s.store.JoinDMCallParticipant(ctx, call.ID, user.ID); err != nil { + return nil, humaError(err) + } + call, err = s.store.GetDMCall(ctx, call.ID) + if err != nil { + return nil, humaError(err) + } + token, err := s.issueDMCallToken(user, call.ID) + if err != nil { + return nil, err + } + s.dispatchDMCallUpdate(call) + output := &dmCallJoinOutput{} + output.Body.Call = s.dmCallPayload(call) + output.Body.Token = token + output.Body.URL = clientURL + output.Body.Room = voice.CallRoomName(call.ID) + return output, nil + }) + + huma.Register(api, huma.Operation{ + OperationID: "declineDirectCall", + Method: http.MethodPost, + Path: "/channels/{channel_id}/call/decline", + Summary: "Отклонить звонок в личной беседе", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + }, + ) (*dmCallOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + call, err := s.store.GetLiveDMCall(ctx, channel.ID) + if err != nil { + return nil, humaErrorStatus(http.StatusNotFound, "dm.call_not_found", "no call is running") + } + // Отклонение инициатора — это отмена звонка: остальные получают конец. + reason := store.DMCallReasonDeclined + if call.InitiatorID == user.ID { + reason = store.DMCallReasonCanceled + } + call, err = s.leaveDMCall(ctx, call, user.ID, store.DMCallDeclined, reason) + if err != nil { + return nil, err + } + output := &dmCallOutput{} + output.Body.Call = s.dmCallPayload(call) + return output, nil + }) + + huma.Register(api, huma.Operation{ + OperationID: "endDirectCall", + Method: http.MethodPost, + Path: "/channels/{channel_id}/call/end", + Summary: "Выйти из звонка или завершить его", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + }, + ) (*dmCallOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + call, err := s.store.GetLiveDMCall(ctx, channel.ID) + if err != nil { + return nil, humaErrorStatus(http.StatusNotFound, "dm.call_not_found", "no call is running") + } + reason := store.DMCallReasonCompleted + if call.Status == store.DMCallRinging { + reason = store.DMCallReasonCanceled + if call.InitiatorID != user.ID { + reason = store.DMCallReasonDeclined + } + } + call, err = s.leaveDMCall(ctx, call, user.ID, store.DMCallLeft, reason) + if err != nil { + return nil, err + } + output := &dmCallOutput{} + output.Body.Call = s.dmCallPayload(call) + return output, nil + }) + + huma.Register(api, huma.Operation{ + OperationID: "updateMyDirectCallState", + Method: http.MethodPatch, + Path: "/channels/{channel_id}/call/@me", + Summary: "Обновить свои флаги в звонке (микрофон, звук, камера, экран)", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + Body struct { + SelfMute *bool `json:"self_mute,omitempty"` + SelfDeaf *bool `json:"self_deaf,omitempty"` + Camera *bool `json:"camera,omitempty"` + Screen *bool `json:"screen,omitempty"` + } + }, + ) (*dmCallOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + call, err := s.store.GetLiveDMCall(ctx, channel.ID) + if err != nil { + return nil, humaErrorStatus(http.StatusNotFound, "dm.call_not_found", "no call is running") + } + current := dmCallParticipant(call, user.ID) + if current == nil || current.State != store.DMCallJoined { + return nil, humaErrorStatus(http.StatusForbidden, "perm.denied", "you are not in this call") + } + flags := store.DMCallFlags{ + SelfMute: current.SelfMute, + SelfDeaf: current.SelfDeaf, + Camera: current.Camera, + Screen: current.Screen, + } + if input.Body.SelfMute != nil { + flags.SelfMute = *input.Body.SelfMute + } + if input.Body.SelfDeaf != nil { + flags.SelfDeaf = *input.Body.SelfDeaf + } + if input.Body.Camera != nil { + flags.Camera = *input.Body.Camera + } + if input.Body.Screen != nil { + flags.Screen = *input.Body.Screen + } + if err := s.store.UpdateDMCallFlags(ctx, call.ID, user.ID, flags); err != nil { + return nil, humaError(err) + } + updated, err := s.store.GetDMCall(ctx, call.ID) + if err != nil { + return nil, humaError(err) + } + s.dispatchDMCallUpdate(updated) + output := &dmCallOutput{} + output.Body.Call = s.dmCallPayload(updated) + return output, nil + }) + + huma.Register(api, huma.Operation{ + OperationID: "getDirectCall", + Method: http.MethodGet, + Path: "/channels/{channel_id}/call", + Summary: "Текущий звонок беседы", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + }, + ) (*dmLiveCallOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + output := &dmLiveCallOutput{} + call, err := s.store.GetLiveDMCall(ctx, channel.ID) + if errors.Is(err, store.ErrNotFound) { + return output, nil + } + if err != nil { + return nil, humaError(err) + } + payload := s.dmCallPayload(call) + output.Body.Call = &payload + return output, nil + }) + + huma.Register(api, huma.Operation{ + OperationID: "listDirectCalls", + Method: http.MethodGet, + Path: "/channels/{channel_id}/calls", + Summary: "История звонков беседы", + Tags: []string{"DirectCalls"}, + Security: security, + }, func(ctx context.Context, input *struct { + ChannelID string `path:"channel_id"` + Limit int `query:"limit" default:"50" minimum:"1" maximum:"100"` + Before string `query:"before,omitempty"` + }, + ) (*dmCallHistoryOutput, error) { + user, _, err := requireUser(ctx) + if err != nil { + return nil, err + } + channel, err := s.groupDMForMember(ctx, input.ChannelID, user.ID) + if err != nil { + return nil, err + } + var beforeID uint64 + if input.Before != "" { + if beforeID, err = parseID("before", input.Before); err != nil { + return nil, err + } + } + calls, err := s.store.ListDMCalls(ctx, channel.ID, input.Limit, beforeID) + if err != nil { + return nil, humaError(err) + } + output := &dmCallHistoryOutput{} + output.Body.Calls = make([]dmCallPayload, 0, len(calls)) + for index := range calls { + output.Body.Calls = append(output.Body.Calls, s.dmCallPayload(&calls[index])) + } + return output, nil + }) +} + +// Принятые значения media. +const ( + dmCallMediaAudio = "audio" + dmCallMediaVideo = "video" +) + +// dmCallUserKey ключует лимит звонков по пользователю. +func dmCallUserKey(userID uint64) string { + return "user:" + formatSnowflake(userID) +} + +// requireVoice проверяет, что голос на инстансе включён, и отдаёт адрес +// сигналинга для клиента. Выключенный голос — понятный отказ voice.disabled, +// а не сломанный интерфейс звонка (AGENT.md 7.14). +func (s *Server) requireVoice() (string, error) { + if !s.voice.Enabled() { + return "", humaErrorStatus(http.StatusServiceUnavailable, "voice.disabled", "голос на инстансе не настроен") + } + clientURL := s.voice.VoiceURLForClient(s.cfg.LiveKitURL) + if clientURL == "" { + return "", humaErrorStatus(http.StatusServiceUnavailable, "voice.disabled", "голос на инстансе не настроен") + } + return clientURL, nil +} + +// issueDMCallToken выдаёт участнику токен LiveKit той же формы, что и вход в +// серверную голосовую комнату: серверных ролей в беседе нет, поэтому права на +// публикацию есть у каждого участника звонка. +func (s *Server) issueDMCallToken(user *store.User, callID uint64) (string, error) { + return s.voice.Issue(voice.CallRoomName(callID), formatSnowflake(user.ID), user.DisplayName, voice.Grants{ + RoomJoin: true, CanPublish: true, CanSubscribe: true, CanPublishData: true, + }) +} + +// dmCallParticipant ищет участника звонка в составе. +func dmCallParticipant(call *store.DMCallWithParticipants, userID uint64) *store.DMCallParticipant { + for index := range call.Participants { + if call.Participants[index].UserID == userID { + return &call.Participants[index] + } + } + return nil +} + +// leaveDMCall помечает участника вышедшим или отклонившим и завершает звонок, +// когда в комнате больше никого не осталось. В беседе вдвоём уход любого +// заканчивает звонок; в групповой — только когда вышли все. +func (s *Server) leaveDMCall(ctx context.Context, call *store.DMCallWithParticipants, userID uint64, state, reason string) (*store.DMCallWithParticipants, error) { + if dmCallParticipant(call, userID) == nil { + return nil, humaErrorStatus(http.StatusForbidden, "perm.denied", "you were not invited to this call") + } + if err := s.store.LeaveDMCallParticipant(ctx, call.ID, userID, state); err != nil { + return nil, humaError(err) + } + s.removeFromLiveKitRoom(ctx, voice.CallRoomName(call.ID), userID) + + joined, err := s.store.CountDMCallParticipants(ctx, call.ID, store.DMCallJoined) + if err != nil { + return nil, humaError(err) + } + lastOne := joined == 0 || len(call.Participants) <= 2 + // Гудки, на которые никто не ответил, тоже заканчиваются: ждать больше некого. + if call.Status == store.DMCallRinging { + waiting, err := s.store.CountDMCallParticipants(ctx, call.ID, store.DMCallInvited) + if err != nil { + return nil, humaError(err) + } + if waiting == 0 { + lastOne = true + } + } + if lastOne { + if err := s.store.EndDMCall(ctx, call.ID, &userID, reason); err != nil { + return nil, humaError(err) + } + call, err = s.store.GetDMCall(ctx, call.ID) + if err != nil { + return nil, humaError(err) + } + // Комнату закрываем со стороны SFU: иначе участники, чьи клиенты не + // успели получить событие, остались бы говорить. + s.removeCallParticipantsFromLiveKit(ctx, call) + s.dispatchDMCallEnd(call) + return call, nil + } + call, err = s.store.GetDMCall(ctx, call.ID) + if err != nil { + return nil, humaError(err) + } + s.dispatchDMCallUpdate(call) + return call, nil +} + +// finishDMCall завершает звонок целиком (сторож и вебхук LiveKit). +func (s *Server) finishDMCall(ctx context.Context, call *store.DMCallWithParticipants, reason string, endedBy *uint64) error { + if err := s.store.EndDMCall(ctx, call.ID, endedBy, reason); err != nil { + return err + } + finished, err := s.store.GetDMCall(ctx, call.ID) + if err != nil { + return err + } + s.removeCallParticipantsFromLiveKit(ctx, finished) + s.dispatchDMCallEnd(finished) + return nil +} + +// removeCallParticipantsFromLiveKit убирает из комнаты всех, кто в ней был. +func (s *Server) removeCallParticipantsFromLiveKit(ctx context.Context, call *store.DMCallWithParticipants) { + if s.voiceAdmin == nil || !s.voiceAdmin.Enabled() { + return + } + room := voice.CallRoomName(call.ID) + for _, participant := range call.Participants { + if participant.State != store.DMCallJoined { + continue + } + if err := s.voiceAdmin.RemoveParticipant(ctx, room, formatSnowflake(participant.UserID)); err != nil { + s.logger.DebugContext(ctx, "livekit call remove skipped", slogAnyError(err)) + } + } +} + +// removeFromLiveKitRoom убирает одного участника из произвольной комнаты. +func (s *Server) removeFromLiveKitRoom(ctx context.Context, room string, userID uint64) { + if s.voiceAdmin == nil || !s.voiceAdmin.Enabled() { + return + } + if err := s.voiceAdmin.RemoveParticipant(ctx, room, formatSnowflake(userID)); err != nil { + s.logger.DebugContext(ctx, "livekit remove skipped", slogAnyError(err)) + } +} + +// dispatchDMCallStart рассылает приглашение остальным участникам беседы, а +// инициатору — обновление: у него может быть открыто второе устройство. +func (s *Server) dispatchDMCallStart(call *store.DMCallWithParticipants, initiatorID uint64) { + if s.gateway == nil { + return + } + payload := s.dmCallPayload(call) + for _, participant := range call.Participants { + if participant.UserID == initiatorID { + s.gateway.SendToUser(participant.UserID, "DM_CALL_UPDATE", payload) + continue + } + s.gateway.SendToUser(participant.UserID, "DM_CALL_RING", payload) + } +} + +// dispatchDMCallUpdate сообщает участникам звонка об изменении состояния. +func (s *Server) dispatchDMCallUpdate(call *store.DMCallWithParticipants) { + if s.gateway == nil { + return + } + payload := s.dmCallPayload(call) + for _, participant := range call.Participants { + s.gateway.SendToUser(participant.UserID, "DM_CALL_UPDATE", payload) + } +} + +// dispatchDMCallEnd сообщает о завершении звонка: клиенты гасят рингтон, +// убирают маркер «идёт звонок» и показывают итог. +func (s *Server) dispatchDMCallEnd(call *store.DMCallWithParticipants) { + if s.gateway == nil { + return + } + payload := s.dmCallPayload(call) + for _, participant := range call.Participants { + s.gateway.SendToUser(participant.UserID, "DM_CALL_END", payload) + } +} + +// dmCallPayload собирает представление звонка для API и Gateway. +func (s *Server) dmCallPayload(call *store.DMCallWithParticipants) dmCallPayload { + payload := dmCallPayload{ + ID: formatSnowflake(call.ID), + ChannelID: formatSnowflake(call.ChannelID), + InitiatorID: formatSnowflake(call.InitiatorID), + Status: call.Status, + Media: call.Media, + CreatedAt: call.CreatedAt.UTC().Format(time.RFC3339), + EndReason: call.EndReason, + DurationSeconds: s.store.DMCallDurationSeconds(&call.DMCall), + Participants: make([]dmCallParticipantPayload, 0, len(call.Participants)), + } + if call.StartedAt != nil { + payload.StartedAt = call.StartedAt.UTC().Format(time.RFC3339) + } + if call.EndedAt != nil { + payload.EndedAt = call.EndedAt.UTC().Format(time.RFC3339) + } + if call.EndedBy != nil { + payload.EndedBy = formatSnowflake(*call.EndedBy) + } + for _, participant := range call.Participants { + entry := dmCallParticipantPayload{ + UserID: formatSnowflake(participant.UserID), + State: participant.State, + SelfMute: participant.SelfMute, + SelfDeaf: participant.SelfDeaf, + Camera: participant.Camera, + Screen: participant.Screen, + } + if participant.JoinedAt != nil { + entry.JoinedAt = participant.JoinedAt.UTC().Format(time.RFC3339) + } + if participant.LeftAt != nil { + entry.LeftAt = participant.LeftAt.UTC().Format(time.RFC3339) + } + payload.Participants = append(payload.Participants, entry) + } + return payload +} + +// liveCallPayload отдаёт текущий звонок беседы для READY-снапшота: nil, если +// никто не звонит (AGENT.md 8.3). +func (s *Server) liveCallPayload(ctx context.Context, channelID uint64) *dmCallPayload { + call, err := s.store.GetLiveDMCall(ctx, channelID) + if err != nil { + return nil + } + payload := s.dmCallPayload(call) + return &payload +} + +// expireStaleDMCalls завершает звонки, о которых уже никто не помнит: гудки, +// на которые минуту никто не ответил, и комнаты, из которых все вышли +// (вебхук LiveKit мог потеряться). Вызывается сторожем голосовых состояний. +func (s *Server) expireStaleDMCalls(ctx context.Context) { + now := time.Now().UTC() + calls, err := s.store.ListStaleDMCalls(ctx, now.Add(-dmCallRingTimeout), now.Add(-time.Minute)) + if err != nil { + s.logger.WarnContext(ctx, "failed to list stale dm calls", slog.Any("error", err)) + return + } + for index := range calls { + call, err := s.store.GetDMCall(ctx, calls[index].ID) + if err != nil { + continue + } + reason := store.DMCallReasonMissed + if calls[index].Status == store.DMCallActive { + reason = store.DMCallReasonCompleted + } + if err := s.finishDMCall(ctx, call, reason, nil); err != nil { + s.logger.WarnContext(ctx, "failed to finish stale dm call", slog.Any("error", err)) + continue + } + s.logger.InfoContext(ctx, "stale dm call finished", + slog.String("call_id", formatSnowflake(call.ID)), slog.String("reason", reason)) + } +} + +// slogAnyError прячет ошибку RoomService в атрибут лога. +func slogAnyError(err error) slog.Attr { return slog.Any("error", err) } diff --git a/internal/server/api_dm_icon.go b/internal/server/api_dm_icon.go index 30519fa..d8eae19 100644 --- a/internal/server/api_dm_icon.go +++ b/internal/server/api_dm_icon.go @@ -178,5 +178,5 @@ func (s *Server) writeDMChannel(ctx context.Context, w http.ResponseWriter, chan writeHumaAPIError(w, err) return } - httpxWriteJSON(w, http.StatusOK, map[string]any{"channel": s.dmChannelPayload(summary)}) + httpxWriteJSON(w, http.StatusOK, map[string]any{"channel": s.dmChannelPayload(ctx, summary)}) } diff --git a/internal/server/api_social.go b/internal/server/api_social.go index d907e7c..89642d2 100644 --- a/internal/server/api_social.go +++ b/internal/server/api_social.go @@ -78,6 +78,10 @@ type dmChannelPayload struct { MemberCount int `json:"member_count,omitempty"` OwnerID string `json:"owner_id,omitempty"` Recipients []dmRecipientPayload `json:"recipients,omitempty"` + // ActiveCall — идущий в беседе звонок (AGENT.md 7.8, Фаза 7): список + // бесед и шапка показывают маркер «идёт звонок» и после перезагрузки, не + // дожидаясь события DM_CALL_UPDATE. + ActiveCall *dmCallPayload `json:"active_call,omitempty"` } // dmRecipientPayload — участник групповой беседы. @@ -404,7 +408,7 @@ func (s *Server) registerSocialRoutes(api huma.API) { Timezone: recipient.Timezone, } output := &dmChannelOutput{} - output.Body.Channel = s.dmChannelPayload(summary) + output.Body.Channel = s.dmChannelPayload(ctx, summary) return output, nil }) @@ -427,7 +431,7 @@ func (s *Server) registerSocialRoutes(api huma.API) { output := &dmChannelListOutput{} output.Body.Channels = make([]dmChannelPayload, 0, len(channels)) for _, summary := range channels { - output.Body.Channels = append(output.Body.Channels, s.dmChannelPayload(summary)) + output.Body.Channels = append(output.Body.Channels, s.dmChannelPayload(ctx, summary)) } return output, nil }) @@ -481,7 +485,7 @@ func (s *Server) registerSocialRoutes(api huma.API) { return nil, err } output := &dmChannelOutput{} - output.Body.Channel = s.dmChannelPayload(summary) + output.Body.Channel = s.dmChannelPayload(ctx, summary) return output, nil }) @@ -544,7 +548,7 @@ func (s *Server) registerSocialRoutes(api huma.API) { return nil, err } output := &dmChannelOutput{} - output.Body.Channel = s.dmChannelPayload(summary) + output.Body.Channel = s.dmChannelPayload(ctx, summary) return output, nil }) @@ -639,7 +643,7 @@ func (s *Server) registerSocialRoutes(api huma.API) { return nil, err } output := &dmChannelOutput{} - output.Body.Channel = s.dmChannelPayload(summary) + output.Body.Channel = s.dmChannelPayload(ctx, summary) return output, nil }) } @@ -926,12 +930,15 @@ func (s *Server) relationshipPayload(_ uint64, profile *store.RelationshipProfil } // dmChannelPayload собирает личную беседу для API. -func (s *Server) dmChannelPayload(summary store.DMChannelSummary) dmChannelPayload { +func (s *Server) dmChannelPayload(ctx context.Context, summary store.DMChannelSummary) dmChannelPayload { payload := dmChannelPayload{ ID: formatSnowflake(summary.Channel.ID), Type: string(store.ChannelDM), CanSend: true, CanView: true, + // Звонок приходит вместе с беседой: клиенту не нужен отдельный запрос, + // чтобы показать «идёт звонок» в списке и шапке (AGENT.md 7.8). + ActiveCall: s.liveCallPayload(ctx, summary.Channel.ID), } if summary.IsGroup { payload.IsGroup = true diff --git a/internal/server/api_voice_webhook.go b/internal/server/api_voice_webhook.go index 8d9e1b6..fe17334 100644 --- a/internal/server/api_voice_webhook.go +++ b/internal/server/api_voice_webhook.go @@ -98,6 +98,12 @@ func (s *Server) markWebhookProcessed(id string) bool { // 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 @@ -194,6 +200,87 @@ func parseVoiceIdentity(identity string) uint64 { 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) { @@ -209,6 +296,7 @@ func (s *Server) startVoiceWatchdog(ctx context.Context, interval time.Duration) return case <-ticker.C: s.cleanupStaleVoiceStates(ctx) + s.expireStaleDMCalls(ctx) } } }() diff --git a/internal/server/dm_calls_test.go b/internal/server/dm_calls_test.go new file mode 100644 index 0000000..6ea46ef --- /dev/null +++ b/internal/server/dm_calls_test.go @@ -0,0 +1,515 @@ +package server + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" + + "glchat/internal/voice" +) + +// dmCallParticipantResponse — участник звонка в ответе API. +type dmCallParticipantResponse struct { + UserID string `json:"user_id"` + State string `json:"state"` + SelfMute bool `json:"self_mute"` + SelfDeaf bool `json:"self_deaf"` + Camera bool `json:"camera"` + Screen bool `json:"screen"` +} + +type dmCallResponse struct { + ID string `json:"id"` + ChannelID string `json:"channel_id"` + InitiatorID string `json:"initiator_id"` + Status string `json:"status"` + Media string `json:"media"` + EndReason string `json:"end_reason"` + DurationSeconds int `json:"duration_seconds"` + Participants []dmCallParticipantResponse `json:"participants"` +} + +// dmCallPayloadResponse — разбор ответа ручек звонка. +type dmCallPayloadResponse struct { + Call dmCallResponse `json:"call"` + Token string `json:"token"` + URL string `json:"url"` + Room string `json:"room"` +} + +// voiceTestServer включает голос в тестовом сервере: ручки звонков без ключей +// LiveKit отвечают voice.disabled (AGENT.md 7.14). +func voiceTestServer(t *testing.T) *Server { + t.Helper() + srv, _ := newTestServer(t) + srv.cfg.LiveKitAPIKey = "test-key" + srv.cfg.LiveKitAPISecret = "test-secret" + srv.cfg.LiveKitURL = "wss://gl.mhspx.su/rtc" + srv.voice = voice.NewIssuer(srv.cfg.LiveKitAPIKey, srv.cfg.LiveKitAPISecret, srv.cfg.LiveKitTokenTTL) + return srv +} + +// openTestDM заводит личную беседу между двумя пользователями и возвращает её id. +func openTestDM(t *testing.T, srv *Server, first *http.Cookie, secondID string) string { + t.Helper() + rec := doJSON(t, srv, http.MethodPost, "/api/v1/users/@me/channels", `{"recipient_id":"`+secondID+`"}`, first) + if rec.Code != http.StatusOK { + t.Fatalf("open direct channel = %d, body = %s", rec.Code, rec.Body.String()) + } + created := decodeResponse[struct { + Channel struct { + ID string `json:"id"` + } `json:"channel"` + }](t, rec) + if created.Channel.ID == "" { + t.Fatalf("direct channel id is empty: %s", rec.Body.String()) + } + return created.Channel.ID +} + +func userIDByEmail(t *testing.T, srv *Server, email string) string { + t.Helper() + user, err := srv.auth.UserByEmail(t.Context(), email) + if err != nil { + t.Fatalf("UserByEmail(%s): %v", email, err) + } + return formatSnowflake(user.ID) +} + +// registerManyUsers заводит несколько аккаунтов: лимит регистраций считается +// по IP (пять в минуту), поэтому часы лимитера сдвигаются вперёд. +func registerManyUsers(t *testing.T, srv *Server, count int, prefix string) []*http.Cookie { + t.Helper() + clock := time.Now() + srv.authLimiter.SetClock(func() time.Time { return clock }) + cookies := make([]*http.Cookie, 0, count) + for index := 0; index < count; index++ { + clock = clock.Add(30 * time.Second) + username := fmt.Sprintf("%s_%02d", prefix, index) + cookies = append(cookies, registerAndLogin(t, srv, username, username+"@example.com")) + } + return cookies +} + +func participantOf(t *testing.T, payload dmCallPayloadResponse, userID string) dmCallParticipantResponse { + t.Helper() + for _, participant := range payload.Call.Participants { + if participant.UserID == userID { + return participant + } + } + t.Fatalf("participant %s is missing in %+v", userID, payload.Call.Participants) + return dmCallParticipantResponse{} +} + +// TestDirectCallLifecycle — звонок 1:1 целиком (AGENT.md 7.8, Фаза 7): +// создание, гудки, принятие, флаги, завершение и история. +func TestDirectCallLifecycle(t *testing.T) { + srv := voiceTestServer(t) + aliceCookie := registerAndLogin(t, srv, "alice_call", "alice-call@example.com") + bobCookie := registerAndLogin(t, srv, "bob_call", "bob-call@example.com") + aliceID := userIDByEmail(t, srv, "alice-call@example.com") + bobID := userIDByEmail(t, srv, "bob-call@example.com") + channelID := openTestDM(t, srv, aliceCookie, bobID) + + // Звонок начинает Алиса: комната dm_call_*, токен LiveKit выдан, гудки идут. + rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{"media":"video"}`, aliceCookie) + if rec.Code != http.StatusOK { + t.Fatalf("start call = %d, body = %s", rec.Code, rec.Body.String()) + } + started := decodeResponse[dmCallPayloadResponse](t, rec) + if started.Call.Status != "ringing" || started.Call.Media != "video" { + t.Fatalf("call must start ringing: %s", rec.Body.String()) + } + if started.Call.InitiatorID != aliceID { + t.Fatalf("initiator = %s, want %s", started.Call.InitiatorID, aliceID) + } + if started.Token == "" || started.URL == "" { + t.Fatalf("call must return a LiveKit token and url: %s", rec.Body.String()) + } + if started.Room != voice.CallRoomName(mustParseID(t, started.Call.ID)) { + t.Fatalf("room = %q", started.Room) + } + if len(started.Call.Participants) != 2 { + t.Fatalf("participants = %+v", started.Call.Participants) + } + if state := participantOf(t, started, aliceID).State; state != "joined" { + t.Fatalf("initiator state = %q, want joined", state) + } + if state := participantOf(t, started, bobID).State; state != "invited" { + t.Fatalf("invitee state = %q, want invited", state) + } + + // Второй звонок в той же беседе начать нельзя. + rec = doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, bobCookie) + if rec.Code != http.StatusConflict { + t.Fatalf("second call = %d, want 409, body = %s", rec.Code, rec.Body.String()) + } + if code := errorCodeOf(t, rec); code != "dm.call_in_progress" { + t.Fatalf("error code = %q, want dm.call_in_progress", code) + } + + // Идущий звонок виден второму участнику вместе с беседой: список рисует + // маркер «идёт звонок» без отдельного запроса (AGENT.md 7.8). + list := decodeResponse[struct { + Channels []struct { + ID string `json:"id"` + ActiveCall *struct { + ID string `json:"id"` + Status string `json:"status"` + } `json:"active_call"` + } `json:"channels"` + }](t, doJSON(t, srv, http.MethodGet, "/api/v1/users/@me/channels", "", bobCookie)) + found := false + for _, channel := range list.Channels { + if channel.ID != channelID { + continue + } + found = true + if channel.ActiveCall == nil || channel.ActiveCall.Status != "ringing" { + t.Fatalf("active_call = %+v", channel.ActiveCall) + } + } + if !found { + t.Fatal("direct channel is missing in the list") + } + + // Принятие: звонок становится active, длительность считается с ответа. + rec = doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/accept", "", bobCookie) + if rec.Code != http.StatusOK { + t.Fatalf("accept call = %d, body = %s", rec.Code, rec.Body.String()) + } + accepted := decodeResponse[dmCallPayloadResponse](t, rec) + if accepted.Call.Status != "active" { + t.Fatalf("call status = %q, want active", accepted.Call.Status) + } + if accepted.Token == "" { + t.Fatal("accept must return a LiveKit token") + } + + // Флаги микрофона и камеры сохраняются у участника звонка. + rec = doJSON(t, srv, http.MethodPatch, "/api/v1/channels/"+channelID+"/call/@me", + `{"self_mute":true,"camera":true}`, bobCookie) + if rec.Code != http.StatusOK { + t.Fatalf("update flags = %d, body = %s", rec.Code, rec.Body.String()) + } + flags := decodeResponse[dmCallPayloadResponse](t, rec) + participant := participantOf(t, flags, bobID) + if !participant.SelfMute || !participant.Camera { + t.Fatalf("flags not saved: %+v", participant) + } + + // Завершает Алиса: в беседе вдвоём уход заканчивает разговор. + rec = doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/end", "", aliceCookie) + if rec.Code != http.StatusOK { + t.Fatalf("end call = %d, body = %s", rec.Code, rec.Body.String()) + } + ended := decodeResponse[dmCallPayloadResponse](t, rec) + if ended.Call.Status != "ended" || ended.Call.EndReason != "completed" { + t.Fatalf("call must be completed: %s", rec.Body.String()) + } + + // После завершения текущего звонка нет, а история содержит запись. + live := decodeResponse[struct { + Call *dmCallResponse `json:"call"` + }](t, doJSON(t, srv, http.MethodGet, "/api/v1/channels/"+channelID+"/call", "", bobCookie)) + if live.Call != nil { + t.Fatalf("live call = %+v, want null", live.Call) + } + history := decodeResponse[struct { + Calls []dmCallResponse `json:"calls"` + }](t, doJSON(t, srv, http.MethodGet, "/api/v1/channels/"+channelID+"/calls", "", bobCookie)) + if len(history.Calls) != 1 { + t.Fatalf("history = %+v", history.Calls) + } + entry := history.Calls[0] + if entry.Status != "ended" || entry.EndReason != "completed" || entry.InitiatorID != aliceID { + t.Fatalf("history entry = %+v", entry) + } + for _, participant := range entry.Participants { + if participant.State != "left" { + t.Fatalf("history participant = %+v, want left", participant) + } + } +} + +// TestDirectCallMissedHistory — звонок, на который не ответили: у приглашённого +// он остаётся отклонённым, а сам звонок закрывается отказом (AGENT.md 7.8). +func TestDirectCallMissedHistory(t *testing.T) { + srv := voiceTestServer(t) + aliceCookie := registerAndLogin(t, srv, "alice_miss", "alice-miss@example.com") + bobCookie := registerAndLogin(t, srv, "bob_miss", "bob-miss@example.com") + bobID := userIDByEmail(t, srv, "bob-miss@example.com") + channelID := openTestDM(t, srv, aliceCookie, bobID) + + if rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, aliceCookie); rec.Code != http.StatusOK { + t.Fatalf("start call = %d, body = %s", rec.Code, rec.Body.String()) + } + // Боб отклоняет — ждать больше некого, звонок завершается отказом. + rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/decline", "", bobCookie) + if rec.Code != http.StatusOK { + t.Fatalf("decline call = %d, body = %s", rec.Code, rec.Body.String()) + } + declined := decodeResponse[dmCallPayloadResponse](t, rec) + if declined.Call.Status != "ended" || declined.Call.EndReason != "declined" { + t.Fatalf("call must end as declined: %s", rec.Body.String()) + } + if state := participantOf(t, declined, bobID).State; state != "declined" { + t.Fatalf("bob state = %q, want declined", state) + } +} + +// TestDirectCallOutsiderDenied — посторонний не видит беседу, не звонит в неё +// и не принимает чужой звонок: 404, существование не подтверждаем (D-070). +func TestDirectCallOutsiderDenied(t *testing.T) { + srv := voiceTestServer(t) + aliceCookie := registerAndLogin(t, srv, "alice_out", "alice-out@example.com") + bobCookie := registerAndLogin(t, srv, "bob_out", "bob-out@example.com") + eveCookie := registerAndLogin(t, srv, "eve_out", "eve-out@example.com") + bobID := userIDByEmail(t, srv, "bob-out@example.com") + channelID := openTestDM(t, srv, aliceCookie, bobID) + + if rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, aliceCookie); rec.Code != http.StatusOK { + t.Fatalf("start call = %d, body = %s", rec.Code, rec.Body.String()) + } + for _, probe := range []struct { + method string + path string + body string + }{ + {http.MethodGet, "/call", ""}, + {http.MethodPost, "/call/accept", ""}, + {http.MethodPost, "/call/decline", ""}, + {http.MethodPost, "/call/end", ""}, + {http.MethodPatch, "/call/@me", `{"self_mute":true}`}, + {http.MethodGet, "/calls", ""}, + } { + rec := doJSON(t, srv, probe.method, "/api/v1/channels/"+channelID+probe.path, probe.body, eveCookie) + if rec.Code != http.StatusNotFound { + t.Fatalf("%s %s = %d, want 404, body = %s", probe.method, probe.path, rec.Code, rec.Body.String()) + } + } + // Чужой звонок всё ещё идёт: посторонний не смог его завершить. + live := decodeResponse[struct { + Call *dmCallResponse `json:"call"` + }](t, doJSON(t, srv, http.MethodGet, "/api/v1/channels/"+channelID+"/call", "", bobCookie)) + if live.Call == nil || live.Call.Status != "ringing" { + t.Fatalf("call must still be ringing: %+v", live.Call) + } +} + +// TestDirectCallGroupLimit — предел звонка в десять участников (AGENT.md 7.8): +// в беседе из одиннадцати человек позвонить нельзя. +func TestDirectCallGroupLimit(t *testing.T) { + srv := voiceTestServer(t) + cookies := registerManyUsers(t, srv, maxDMCallParticipants, "call_limit") + ids := make([]string, 0, maxDMCallParticipants) + for index := 0; index < maxDMCallParticipants; index++ { + ids = append(ids, userIDByEmail(t, srv, fmt.Sprintf("call_limit_%02d@example.com", index))) + } + owner := cookies[0] + body := `{"name":"Лимит","recipient_ids":[` + for index := 1; index < maxDMCallParticipants; index++ { + if index > 1 { + body += "," + } + body += `"` + ids[index] + `"` + } + body += `]}` + rec := doJSON(t, srv, http.MethodPost, "/api/v1/users/@me/channels/group", body, owner) + if rec.Code != http.StatusOK { + t.Fatalf("create group = %d, body = %s", rec.Code, rec.Body.String()) + } + channelID := decodeResponse[struct { + Channel struct { + ID string `json:"id"` + MemberCount int `json:"member_count"` + } `json:"channel"` + }](t, rec).Channel.ID + + // Ровно десять участников — звонок разрешён, состав весь в приглашении. + started := decodeResponse[dmCallPayloadResponse]( + t, doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, owner)) + if len(started.Call.Participants) != maxDMCallParticipants { + t.Fatalf("participants = %d, want %d", len(started.Call.Participants), maxDMCallParticipants) + } + // Все девять приглашённых принимают: в комнате десять человек. + for index := 1; index < maxDMCallParticipants; index++ { + rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/accept", "", cookies[index]) + if rec.Code != http.StatusOK { + t.Fatalf("accept #%d = %d, body = %s", index, rec.Code, rec.Body.String()) + } + } + // В групповой беседе уход одного разговор не заканчивает: звонок живёт, + // пока в комнате кто-то остаётся, и закрывается на последнем. + for index := 0; index < maxDMCallParticipants-1; index++ { + rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/end", "", cookies[index]) + if rec.Code != http.StatusOK { + t.Fatalf("leave #%d = %d, body = %s", index, rec.Code, rec.Body.String()) + } + if payload := decodeResponse[dmCallPayloadResponse](t, rec); payload.Call.Status != "active" { + t.Fatalf("call must stay active after leave #%d: %s", index, rec.Body.String()) + } + } + live := decodeResponse[dmCallPayloadResponse](t, doJSON(t, srv, http.MethodPost, + "/api/v1/channels/"+channelID+"/call/end", "", cookies[maxDMCallParticipants-1])) + if live.Call.Status != "ended" { + t.Fatalf("call status = %q, want ended", live.Call.Status) + } + + // Одиннадцатый участник: беседа больше предела — звонить нельзя. Состав + // расширяем напрямую в хранилище: ручка добавления держит тот же лимит. + extra := registerManyUsers(t, srv, 1, "call_extra") + _ = extra + if err := srv.store.AddDMParticipant(t.Context(), mustParseID(t, channelID), + mustParseID(t, userIDByEmail(t, srv, "call_extra_00@example.com"))); err != nil { + t.Fatalf("add extra participant: %v", err) + } + rec = doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, owner) + if rec.Code != http.StatusConflict { + t.Fatalf("call in a big conversation = %d, want 409, body = %s", rec.Code, rec.Body.String()) + } + if code := errorCodeOf(t, rec); code != "dm.call_limit" { + t.Fatalf("error code = %q, want dm.call_limit", code) + } +} + +// TestDirectCallVoiceDisabled — выключенный на инстансе голос даёт понятный +// отказ, а не сломанный звонок (AGENT.md 7.14). +func TestDirectCallVoiceDisabled(t *testing.T) { + srv, _ := newTestServer(t) + aliceCookie := registerAndLogin(t, srv, "alice_novoice", "alice-novoice@example.com") + bobCookie := registerAndLogin(t, srv, "bob_novoice", "bob-novoice@example.com") + bobID := userIDByEmail(t, srv, "bob-novoice@example.com") + channelID := openTestDM(t, srv, aliceCookie, bobID) + _ = bobCookie + + rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, aliceCookie) + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("start call without voice = %d, want 503, body = %s", rec.Code, rec.Body.String()) + } + if code := errorCodeOf(t, rec); code != "voice.disabled" { + t.Fatalf("error code = %q, want voice.disabled", code) + } + // Никакого звонка в базе не осталось. + if _, err := srv.store.GetLiveDMCall(t.Context(), mustParseID(t, channelID)); err == nil { + t.Fatal("no call must be stored when voice is disabled") + } +} + +// TestDirectCallGatewayEvents — приглашение приходит событием: участники +// получают DM_CALL_RING/DM_CALL_UPDATE/DM_CALL_END, посторонний — ничего. +func TestDirectCallGatewayEvents(t *testing.T) { + srv := voiceTestServer(t) + httpServer := httptest.NewServer(srv.Handler()) + t.Cleanup(httpServer.Close) + + aliceCookie := registerAndLogin(t, srv, "alice_ev", "alice-ev@example.com") + bobCookie := registerAndLogin(t, srv, "bob_ev", "bob-ev@example.com") + eveCookie := registerAndLogin(t, srv, "eve_ev", "eve-ev@example.com") + bobID := userIDByEmail(t, srv, "bob-ev@example.com") + channelID := openTestDM(t, srv, aliceCookie, bobID) + + aliceClient := dialGateway(t, httpServer, aliceCookie) + bobClient := dialGateway(t, httpServer, bobCookie) + eveClient := dialGateway(t, httpServer, eveCookie) + + if rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, aliceCookie); rec.Code != http.StatusOK { + t.Fatalf("start call = %d, body = %s", rec.Code, rec.Body.String()) + } + + // Боб получает приглашение, Алиса — обновление (у неё может быть второе + // устройство), посторонний не получает ничего. + ring := bobClient.expectEvent("DM_CALL_RING") + var ringPayload dmCallResponse + if err := json.Unmarshal(ring.D, &ringPayload); err != nil { + t.Fatalf("decode DM_CALL_RING: %v", err) + } + if ringPayload.Status != "ringing" || ringPayload.Media != "audio" || ringPayload.ChannelID != channelID { + t.Fatalf("ring payload = %+v", ringPayload) + } + aliceClient.expectEvent("DM_CALL_UPDATE") + + // Принятие видно инициатору как обновление звонка. + if rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/accept", "", bobCookie); rec.Code != http.StatusOK { + t.Fatalf("accept = %d, body = %s", rec.Code, rec.Body.String()) + } + update := aliceClient.expectEvent("DM_CALL_UPDATE") + var updatePayload dmCallResponse + if err := json.Unmarshal(update.D, &updatePayload); err != nil { + t.Fatalf("decode DM_CALL_UPDATE: %v", err) + } + if updatePayload.Status != "active" { + t.Fatalf("status = %q, want active", updatePayload.Status) + } + + // Завершение приходит отдельным событием с итогом. + if rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call/end", "", aliceCookie); rec.Code != http.StatusOK { + t.Fatalf("end = %d, body = %s", rec.Code, rec.Body.String()) + } + end := bobClient.expectEvent("DM_CALL_END") + var endPayload dmCallResponse + if err := json.Unmarshal(end.D, &endPayload); err != nil { + t.Fatalf("decode DM_CALL_END: %v", err) + } + if endPayload.Status != "ended" || endPayload.EndReason != "completed" { + t.Fatalf("end payload = %+v", endPayload) + } + + // Посторонний не получил события о звонке: беседа ему не видна. + eveClient.expectNoEvent("DM_CALL_RING", 300*time.Millisecond) +} + +// TestDirectCallInReadySnapshot — «идёт звонок» переживает перезагрузку: звонок +// приходит в READY вместе с беседой (AGENT.md 8.3). +func TestDirectCallInReadySnapshot(t *testing.T) { + srv := voiceTestServer(t) + httpServer := httptest.NewServer(srv.Handler()) + t.Cleanup(httpServer.Close) + + aliceCookie := registerAndLogin(t, srv, "alice_ready", "alice-ready@example.com") + bobCookie := registerAndLogin(t, srv, "bob_ready", "bob-ready@example.com") + bobID := userIDByEmail(t, srv, "bob-ready@example.com") + channelID := openTestDM(t, srv, aliceCookie, bobID) + + if rec := doJSON(t, srv, http.MethodPost, "/api/v1/channels/"+channelID+"/call", `{}`, aliceCookie); rec.Code != http.StatusOK { + t.Fatalf("start call = %d, body = %s", rec.Code, rec.Body.String()) + } + + client := dialGateway(t, httpServer, bobCookie) + var ready struct { + DMChannels []struct { + ID string `json:"id"` + ActiveCall *dmCallResponse `json:"active_call"` + } `json:"dm_channels"` + } + if err := json.Unmarshal(client.ready.D, &ready); err != nil { + t.Fatalf("decode READY: %v", err) + } + for _, channel := range ready.DMChannels { + if channel.ID != channelID { + continue + } + if channel.ActiveCall == nil || channel.ActiveCall.Status != "ringing" { + t.Fatalf("active_call = %+v", channel.ActiveCall) + } + if len(channel.ActiveCall.Participants) != 2 { + t.Fatalf("participants = %+v", channel.ActiveCall.Participants) + } + return + } + t.Fatal("direct channel is missing in READY") +} + +// mustParseID переводит строковый идентификатор в uint64. +func mustParseID(t *testing.T, value string) uint64 { + t.Helper() + id, err := parseID("id", value) + if err != nil { + t.Fatalf("parse id %q: %v", value, err) + } + return id +} diff --git a/internal/server/server.go b/internal/server/server.go index 27b96ad..f61c378 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -58,6 +58,9 @@ type Server struct { editLimiter *httpx.RateLimiter // soundboardLimiter — не чаще 3 звуков в 10 секунд на пользователя. soundboardLimiter *httpx.RateLimiter + // dmCallLimiter — 10 новых звонков в минуту на пользователя (Фаза 7): + // гудки в беседе не должны превращаться в спам. + dmCallLimiter *httpx.RateLimiter // inviteLimiter — 10 приглашений в сутки на пользователя (AGENT.md 8.6). inviteLimiter *httpx.RateLimiter // webhookLimiter — 30 сообщений в минуту на вебхук (AGENT.md 7.11, 8.6). @@ -136,6 +139,7 @@ func New(cfg config.Config, db *database.DB, logger *slog.Logger, deps Deps) *Se searchLimiter: httpx.NewRateLimiter(10, 10), editLimiter: httpx.NewRateLimiter(10, 10), soundboardLimiter: httpx.NewRateLimiterWindow(3, 10*time.Second, 3), + dmCallLimiter: httpx.NewRateLimiterWindow(dmCallLimitCount, dmCallLimitWindow, dmCallLimitCount), inviteLimiter: httpx.NewRateLimiterWindow(10, 24*time.Hour, 10), webhookLimiter: httpx.NewRateLimiterWindow(webhookRateLimit, time.Minute, webhookRateLimit), // Второй аргумент NewRateLimiter — запас (burst), а не окно. @@ -216,6 +220,7 @@ func New(cfg config.Config, db *database.DB, logger *slog.Logger, deps Deps) *Se s.registerDMIconRoutes(apiRouter) s.registerEmojiRoutes(s.api, apiRouter) s.registerVoiceRoutes(s.api) + s.registerDMCallRoutes(s.api) s.registerSoundsRoutes(s.api, apiRouter) s.registerWebhookRoutes(s.api, apiRouter) s.registerCosmeticRoutes(s.api, apiRouter) diff --git a/internal/voice/token.go b/internal/voice/token.go index 00c492e..b862d8b 100644 --- a/internal/voice/token.go +++ b/internal/voice/token.go @@ -121,6 +121,24 @@ func ParseRoomName(room string) (guildID, channelID uint64, ok bool) { return guildID, channelID, true } +// CallRoomName собирает имя комнаты звонка в личной беседе. Отдельный префикс +// нужен и вебхуку LiveKit, и сторожу: по имени комнаты видно, серверная это +// комната или звонок в DM (AGENT.md 7.8). +func CallRoomName(callID uint64) string { + return fmt.Sprintf("dm_call_%d", callID) +} + +// ParseCallRoomName разбирает имя комнаты звонка. +func ParseCallRoomName(room string) (callID uint64, ok bool) { + if !strings.HasPrefix(room, "dm_call_") { + return 0, false + } + if _, err := fmt.Sscanf(room, "dm_call_%d", &callID); err != nil || callID == 0 { + return 0, false + } + return callID, true +} + func encodeSegment(payload []byte) string { return base64.RawURLEncoding.EncodeToString(payload) }