2026-09-19 21:42:41 +03:00
|
|
|
// 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"`
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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 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"`
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-19 23:32:11 +03:00
|
|
|
// Visibility отвечает, видит ли пользователь комнату: события комнат получают
|
|
|
|
|
// только те сессии, у которых есть VIEW_CHANNEL (AGENT.md 8.3, 9.7).
|
|
|
|
|
type Visibility interface {
|
|
|
|
|
CanViewChannel(ctx context.Context, guildID, channelID, userID uint64, instanceAdmin bool) bool
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-19 21:42:41 +03:00
|
|
|
// Service — Gateway: управляет подключениями и рассылкой событий.
|
|
|
|
|
type Service struct {
|
2026-09-19 23:32:11 +03:00
|
|
|
store *store.Store
|
|
|
|
|
auth *auth.Service
|
|
|
|
|
readiness SnapshotBuilder
|
|
|
|
|
visibility Visibility
|
|
|
|
|
logger *slog.Logger
|
2026-09-19 21:42:41 +03:00
|
|
|
// 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 {
|
2026-09-19 23:32:11 +03:00
|
|
|
service := &Service{
|
2026-09-19 21:42:41 +03:00
|
|
|
store: st,
|
|
|
|
|
auth: authService,
|
|
|
|
|
readiness: builder,
|
|
|
|
|
logger: logger,
|
|
|
|
|
allowedOrigins: allowedOrigins,
|
|
|
|
|
sessions: map[string]*clientSession{},
|
|
|
|
|
buffers: map[string]*resumeBuffer{},
|
|
|
|
|
}
|
2026-09-19 23:32:11 +03:00
|
|
|
// Сборщик READY умеет считать права: переиспользуем его для фильтрации.
|
|
|
|
|
if visibility, ok := builder.(Visibility); ok {
|
|
|
|
|
service.visibility = visibility
|
|
|
|
|
}
|
|
|
|
|
return service
|
2026-09-19 21:42:41 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// 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)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-19 23:32:11 +03:00
|
|
|
// 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 проверяет видимость комнаты для сессии. Без калькулятора
|
|
|
|
|
// прав событие доставляется всем: так работают тесты и режим без БД.
|
|
|
|
|
func (s *Service) canViewChannel(ctx context.Context, guildID, channelID uint64, session *clientSession) bool {
|
|
|
|
|
if s.visibility == nil || guildID == 0 {
|
|
|
|
|
return true
|
|
|
|
|
}
|
|
|
|
|
instanceAdmin := session.user != nil && session.user.IsInstanceAdmin
|
|
|
|
|
return s.visibility.CanViewChannel(ctx, guildID, channelID, session.userID, instanceAdmin)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-19 21:42:41 +03:00
|
|
|
// 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})
|
|
|
|
|
}
|