Files
glchat/internal/gateway/gateway.go
T
grendervill 3dc200c196 feat(api): расширенная админ-панель инстанса (Фаза 7)
Статистика, карточка пользователя и работа с журналом для администратора
инстанса (AGENT.md 3.2, 7.18):

- GET /instance/stats: рост пользователей и сообщений по дням, активность,
  размеры базы и файлов, топы серверов по участникам и сообщениям;
- GET /instance/users/{id}: профиль, активные устройства, серверы с ролями,
  события безопасности и аудит по пользователю;
- DELETE /instance/users/{id}/sessions/{sid} и POST .../reset-2fa: отзыв
  одного устройства и сброс второго фактора со step-up, аудитом и записью
  в события безопасности; чужой ключ администратора не сбрасывается;
- журнал инстанса: фильтры по действию, актору, цели, серверу и датам,
  пагинация с общим числом, список действий и выгрузка CSV (лимит 5/мин);
- список серверов: поиск по названию, владелец, главный сервер, пагинация;
- миграция 00025: нормализованное название сервера `name_lower` — SQLite
  lower() не знает кириллицу, поэтому регистр приводит приложение (как для
  текста сообщений), старые записи дополняются backfill'ом при старте.
2026-09-26 16:44:49 +03:00

478 lines
17 KiB
Go
Raw 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"`
}
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})
}