Files
glchat/internal/gateway/gateway.go
T

506 lines
18 KiB
Go
Raw Normal View History

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