3dc200c196
Статистика, карточка пользователя и работа с журналом для администратора
инстанса (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'ом при старте.
478 lines
17 KiB
Go
478 lines
17 KiB
Go
// 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})
|
||
}
|