Files
glchat/internal/gateway/gateway.go
T

371 lines
12 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"`
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"`
}
// 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"`
// Поля личных бесед (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"`
}
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)
}
}
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})
}