Files
glchat/internal/gateway/gateway.go
T
grendervill 0b55fbdfe2 feat(dm): ручки и события звонков в беседах (Фаза 7)
Звонки 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.
2026-09-26 17:02:55 +03:00

506 lines
18 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package gateway реализует WebSocket Gateway glchat (AGENT.md 8.3):
// HELLO/IDENTIFY/READY/HEARTBEAT/RESUME, единый JSON-конверт и рассылку
// событий с фильтрацией по правам.
package gateway
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"net/http"
"sync"
"time"
"glchat/internal/auth"
"glchat/internal/store"
)
// Оп-коды протокола (AGENT.md 8.3).
const (
OpDispatch = 0
OpHeartbeat = 1
OpIdentify = 2
OpResume = 3
OpInvalidSess = 4
OpReconnect = 7
OpHello = 10
OpHeartbeatAck = 11
)
// Ограничения протокола (AGENT.md 8.3).
const (
HeartbeatInterval = 45 * time.Second
IdentifyRateLimit = 5 * time.Second
MaxIncomingFrame = 16 << 10
MaxOutgoingFrame = 64 << 10
ResumeBufferSize = 1000
ResumeBufferTTL = 5 * time.Minute
WriteTimeout = 10 * time.Second
)
// Envelope — единый конверт сообщений Gateway.
type Envelope struct {
Op int `json:"op"`
T string `json:"t,omitempty"`
D json.RawMessage `json:"d,omitempty"`
S int64 `json:"s,omitempty"`
}
// Ready — снапшот состояния пользователя при подключении (AGENT.md 8.3).
type Ready struct {
User ReadyUser `json:"user"`
Guilds []ReadyGuild `json:"guilds"`
DMChannels []ReadyChannel `json:"dm_channels"`
ReadStates []ReadyRead `json:"read_states"`
SessionID string `json:"session_id"`
HeartbeatMS int `json:"heartbeat_interval_ms"`
}
type ReadyUser struct {
ID string `json:"id"`
Username string `json:"username"`
DisplayName string `json:"display_name"`
AvatarFileID string `json:"avatar_file_id,omitempty"`
IsInstanceAdmin bool `json:"is_instance_admin"`
Badges []string `json:"badges"`
}
type ReadyGuild struct {
ID string `json:"id"`
Name string `json:"name"`
OwnerID string `json:"owner_id"`
IsMain bool `json:"is_main"`
// Оформление сервера (AGENT.md 7.4): иконка, баннер, splash и акцент.
IconFileID string `json:"icon_file_id,omitempty"`
BannerFileID string `json:"banner_file_id,omitempty"`
SplashFileID string `json:"splash_file_id,omitempty"`
AccentColor int64 `json:"accent_color,omitempty"`
Channels []ReadyChannel `json:"channels"`
Roles []ReadyRole `json:"roles"`
MemberIDs []string `json:"member_ids"`
MyRoles []string `json:"my_role_ids"`
MyNickname string `json:"my_nickname,omitempty"`
MyPerms []string `json:"my_permissions"`
// Emojis — кастомные эмодзи сервера (AGENT.md 7.12).
Emojis []ReadyEmoji `json:"emojis"`
// VoiceStates — кто находится в голосовых комнатах (AGENT.md 7.14).
VoiceStates []ReadyVoiceState `json:"voice_states"`
// Sounds — саундборд и звуковая палитра сервера (AGENT.md 7.13).
Sounds []ReadySound `json:"sounds"`
}
// ReadySound — звук сервера в снапшоте.
type ReadySound struct {
ID string `json:"id"`
Name string `json:"name"`
FileID string `json:"file_id"`
Kind string `json:"kind"`
Event string `json:"event,omitempty"`
Emoji string `json:"emoji,omitempty"`
}
// ReadyVoiceState — участник голосовой комнаты в снапшоте.
type ReadyVoiceState struct {
UserID string `json:"user_id"`
ChannelID string `json:"channel_id"`
SelfMute bool `json:"self_mute"`
SelfDeaf bool `json:"self_deaf"`
ServerMute bool `json:"server_mute"`
ServerDeaf bool `json:"server_deaf"`
Camera bool `json:"camera"`
Screen bool `json:"screen"`
}
// ReadyEmoji — эмодзи сервера в снапшоте: имя, файл и токен для вставки.
type ReadyEmoji struct {
ID string `json:"id"`
Name string `json:"name"`
FileID string `json:"file_id"`
Animated bool `json:"animated"`
Token string `json:"token"`
}
type ReadyChannel struct {
ID string `json:"id"`
GuildID string `json:"guild_id,omitempty"`
Type string `json:"type"`
Name string `json:"name"`
Position int `json:"position"`
ParentID string `json:"parent_id,omitempty"`
UserLimit int `json:"user_limit,omitempty"`
Slowmode int `json:"slowmode_seconds,omitempty"`
CanSend bool `json:"can_send"`
CanConnect bool `json:"can_connect"`
// Description — статус комнаты: шапка показывает его рядом с названием
// (AGENT.md 7.5, 11.6), поэтому он приходит и в READY, и в событиях комнат.
Description string `json:"description,omitempty"`
// BackgroundFileID — фон комнаты: приходит и в READY, и в CHANNEL_UPDATE
// (AGENT.md 7.5), иначе после перезагрузки фон терялся бы до REST-запроса.
BackgroundFileID string `json:"background_file_id,omitempty"`
// Поля личных бесед (type = dm): собеседник и последнее сообщение.
RecipientID string `json:"recipient_id,omitempty"`
IconFileID string `json:"icon_file_id,omitempty"`
Status string `json:"status,omitempty"`
LastMessageID string `json:"last_message_id,omitempty"`
LastMessageAt string `json:"last_message_at,omitempty"`
// Групповая личная беседа (AGENT.md 7.8, Фаза 7): имя, состав и владелец.
// У 1:1 is_group = false, а собеседник лежит в recipient_id.
IsGroup bool `json:"is_group,omitempty"`
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 {
ID string `json:"id"`
Name string `json:"name"`
Color int64 `json:"color"`
Position int `json:"position"`
Permissions string `json:"permissions"`
IsDefault bool `json:"is_default"`
Mentionable bool `json:"mentionable"`
Hoist bool `json:"hoist"`
}
type ReadyRead struct {
ChannelID string `json:"channel_id"`
LastMessageID string `json:"last_message_id,omitempty"`
MentionCount int `json:"mention_count"`
}
// serverFrame — входящее сообщение от клиента.
type serverFrame struct {
Op int `json:"op"`
D json.RawMessage `json:"d"`
T string `json:"t"`
S int64 `json:"s"`
}
type identifyPayload struct {
Token string `json:"token"`
ResumeSeq int64 `json:"resume_seq"`
SessionID string `json:"session_id"`
}
// Visibility отвечает, видит ли пользователь комнату: события комнат получают
// только те сессии, у которых есть VIEW_CHANNEL (AGENT.md 8.3, 9.7).
type Visibility interface {
CanViewChannel(ctx context.Context, guildID, channelID, userID uint64, instanceAdmin bool) bool
}
// Service — Gateway: управляет подключениями и рассылкой событий.
type Service struct {
store *store.Store
auth *auth.Service
readiness SnapshotBuilder
visibility Visibility
logger *slog.Logger
// allowedOrigins — домены, с которых разрешено подключаться (AGENT.md 9.7).
allowedOrigins []string
mu sync.RWMutex
sessions map[string]*clientSession
buffers map[string]*resumeBuffer
seq int64
}
// SnapshotBuilder собирает READY-снапшот (вынесено для тестируемости).
type SnapshotBuilder interface {
Build(ctx context.Context, user *store.User) (*Ready, error)
}
func New(st *store.Store, authService *auth.Service, builder SnapshotBuilder, logger *slog.Logger, allowedOrigins []string) *Service {
service := &Service{
store: st,
auth: authService,
readiness: builder,
logger: logger,
allowedOrigins: allowedOrigins,
sessions: map[string]*clientSession{},
buffers: map[string]*resumeBuffer{},
}
// Сборщик READY умеет считать права: переиспользуем его для фильтрации.
if visibility, ok := builder.(Visibility); ok {
service.visibility = visibility
}
return service
}
// Handler отдаёт http.Handler для маршрута /gateway.
func (s *Service) Handler() http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
s.ServeHTTP(w, r)
})
}
// Dispatch рассылает событие всем сессиям пользователя (AGENT.md 8.3).
func (s *Service) Dispatch(event string, payload any) {
raw, err := json.Marshal(payload)
if err != nil {
s.logger.Error("marshal gateway payload", slog.String("event", event), slog.Any("error", err))
return
}
s.mu.Lock()
s.seq++
seq := s.seq
envelope := Envelope{Op: OpDispatch, T: event, D: raw, S: seq}
encoded, err := json.Marshal(envelope)
if err != nil {
s.mu.Unlock()
s.logger.Error("marshal gateway envelope", slog.Any("error", err))
return
}
sessions := make([]*clientSession, 0, len(s.sessions))
for _, session := range s.sessions {
sessions = append(sessions, session)
}
// Буфер RESUME наполняем для всех известных сессий, включая недавно
// отключённые: иначе события, пришедшие в разрыв связи, потеряются
// (AGENT.md 8.3). Буферы ограничены по размеру и времени жизни.
for _, buffer := range s.buffers {
buffer.append(seq, encoded)
}
s.mu.Unlock()
for _, session := range sessions {
if err := session.write(encoded); err != nil {
s.logger.Debug("gateway write failed", slog.Any("error", err))
}
}
if s.seq%512 == 0 {
s.cleanupBuffers()
}
}
func (s *Service) cleanupBuffers() {
cutoff := time.Now().Add(-ResumeBufferTTL)
s.mu.Lock()
defer s.mu.Unlock()
for id, buffer := range s.buffers {
if buffer.lastSeen().Before(cutoff) {
delete(s.buffers, id)
}
}
}
// ActiveSessions возвращает число активных сессий (метрики и тесты).
func (s *Service) ActiveSessions() int {
s.mu.RLock()
defer s.mu.RUnlock()
return len(s.sessions)
}
// DispatchToChannel рассылает событие комнаты только тем сессиям, которые
// видят эту комнату (AGENT.md 8.3: права фильтруются на сервере).
func (s *Service) DispatchToChannel(ctx context.Context, channelID uint64, event string, payload any) {
s.dispatchToChannel(ctx, channelID, 0, event, payload)
}
// DispatchToChannelExcept рассылает событие комнаты всем, кроме указанного
// пользователя: так работает typing (AGENT.md 7.6).
func (s *Service) DispatchToChannelExcept(ctx context.Context, channelID, exceptUserID uint64, event string, payload any) {
s.dispatchToChannel(ctx, channelID, exceptUserID, event, payload)
}
func (s *Service) dispatchToChannel(ctx context.Context, channelID, exceptUserID uint64, event string, payload any) {
channel, err := s.store.GetChannel(ctx, channelID)
if err != nil {
s.logger.DebugContext(ctx, "gateway channel event skipped", slog.Any("error", err))
return
}
guildID := uint64(0)
if channel.GuildID != nil {
guildID = *channel.GuildID
}
raw, err := json.Marshal(payload)
if err != nil {
s.logger.ErrorContext(ctx, "marshal channel event", slog.String("event", event), slog.Any("error", err))
return
}
s.mu.Lock()
s.seq++
envelope := Envelope{Op: OpDispatch, T: event, D: raw, S: s.seq}
encoded, err := json.Marshal(envelope)
if err != nil {
s.mu.Unlock()
s.logger.ErrorContext(ctx, "marshal channel envelope", slog.Any("error", err))
return
}
targets := make([]*clientSession, 0, len(s.sessions))
for _, session := range s.sessions {
targets = append(targets, session)
}
for _, buffer := range s.buffers {
buffer.append(envelope.S, encoded)
}
s.mu.Unlock()
for _, session := range targets {
if exceptUserID != 0 && session.userID == exceptUserID {
continue
}
if !s.canViewChannel(ctx, guildID, channelID, session) {
continue
}
if err := session.write(encoded); err != nil {
s.logger.DebugContext(ctx, "gateway write failed", slog.Any("error", err))
}
}
}
// canViewChannel проверяет видимость комнаты для сессии. Для личных бесед
// (без сервера) событие получают только участники (AGENT.md 7.8, 9.7).
func (s *Service) canViewChannel(ctx context.Context, guildID, channelID uint64, session *clientSession) bool {
if guildID == 0 {
participant, err := s.store.IsDMParticipant(ctx, channelID, session.userID)
if err != nil {
return false
}
return participant
}
if s.visibility == nil {
return true
}
instanceAdmin := session.user != nil && session.user.IsInstanceAdmin
return s.visibility.CanViewChannel(ctx, guildID, channelID, session.userID, instanceAdmin)
}
// SendToUser доставляет событие конкретному пользователю.
func (s *Service) SendToUser(userID uint64, event string, payload any) {
raw, err := json.Marshal(payload)
if err != nil {
return
}
s.mu.Lock()
s.seq++
envelope := Envelope{Op: OpDispatch, T: event, D: raw, S: s.seq}
encoded, err := json.Marshal(envelope)
if err != nil {
s.mu.Unlock()
return
}
targets := make([]*clientSession, 0, 2)
for _, session := range s.sessions {
if session.userID == userID {
targets = append(targets, session)
}
}
// Событие попадает и в буфер отключённого пользователя: при RESUME он
// получит пропущенное (AGENT.md 8.3).
for _, buffer := range s.buffers {
if buffer.owner == userID {
buffer.append(envelope.S, encoded)
}
}
s.mu.Unlock()
for _, session := range targets {
_ = session.write(encoded)
}
}
// InvalidateUser закрывает все соединения пользователя: клиент получает
// INVALID_SESSION и уходит на экран входа. Нужно там, где сессии отозваны
// (logout-all, смена пароля, действия администратора) — иначе второе
// устройство остаётся «в приложении» до следующего запроса (AGENT.md 11.6).
func (s *Service) InvalidateUser(userID uint64, reason string) {
s.invalidateUser(userID, "", reason)
}
// InvalidateUserExcept закрывает соединения пользователя, кроме текущего:
// смена пароля отзывает остальные сессии, но устройство, с которого её
// сменили, продолжает работать (AGENT.md 7.1, 11.6).
func (s *Service) InvalidateUserExcept(userID uint64, keepTokenHash, reason string) {
s.invalidateUser(userID, keepTokenHash, reason)
}
// InvalidateSession закрывает соединение одной сессии пользователя: панель
// инстанса отзывает конкретное устройство (AGENT.md 7.18), остальные
// устройства продолжают работу. Соединение опознаётся по хэшу токена сессии —
// как и в InvalidateUserExcept, только условие обратное.
func (s *Service) InvalidateSession(userID uint64, tokenHash, reason string) {
if tokenHash == "" {
return
}
s.mu.Lock()
targets := make([]*clientSession, 0, 1)
for _, session := range s.sessions {
if session.userID != userID {
continue
}
if session.buffer == nil || session.buffer.tokenHash != tokenHash {
continue
}
targets = append(targets, session)
}
s.mu.Unlock()
for _, session := range targets {
// Сначала кадр (writeSync дожидается записи), потом закрытие.
s.sendInvalidSession(session, reason)
session.close()
}
}
func (s *Service) invalidateUser(userID uint64, keepTokenHash, reason string) {
s.mu.Lock()
targets := make([]*clientSession, 0, 2)
for _, session := range s.sessions {
if session.userID != userID {
continue
}
// Соединение опознаём по хэшу токена сессии: у каждой сессии свой буфер
// RESUME, он же хранит владельца (AGENT.md 8.3).
if keepTokenHash != "" && session.buffer != nil && session.buffer.tokenHash == keepTokenHash {
continue
}
targets = append(targets, session)
}
s.mu.Unlock()
for _, session := range targets {
// Сначала кадр (writeSync дожидается записи), потом закрытие.
s.sendInvalidSession(session, reason)
session.close()
}
}
func frame(op int, payload any) ([]byte, error) {
var encoded json.RawMessage
if payload != nil {
raw, err := json.Marshal(payload)
if err != nil {
return nil, fmt.Errorf("marshal payload: %w", err)
}
encoded = raw
}
return json.Marshal(Envelope{Op: op, D: encoded})
}