Files
Edu/services/push-gateway/internal/hub/hub.go
SpecialX 0a71b02e04
Some checks failed
CI / quality-ts (push) Failing after 48s
CI / quality-go (push) Failing after 4s
CI / quality-proto (push) Failing after 2s
CI / deploy (push) Has been skipped
fix: code compliance audit and fix across all services
NestJS (6 services): implement @RequirePermission decorator with
SetMetadata+Reflector, register APP_GUARD globally, fix as assertions
to type guards, add explicit return types, fix import type for express,
fix /metrics implicit any, replace native Error with ApplicationError,
remove typeorm remnants, register LifecycleService.

teacher-bff: add logger, ApplicationError, GlobalErrorFilter, forward
real userId to downstream, log downstream failures, migrate health
controller to shared/health.

Go (2 services): interface to any, doc comments, CORS dev whitelist,
JWT secret fail-fast, push-gateway internal API auth, metrics and
readyz endpoints, remove dead code.

Python (2 services): lifespan return type, dev_mode to bool, data-ana
APIRouter, ai POST body model, ClickHouse async wrapping.
2026-07-09 17:28:27 +08:00

111 lines
2.2 KiB
Go

package hub
import (
"sync"
"github.com/gorilla/websocket"
)
// Connection 包装单个 WebSocket 连接及其异步发送通道
type Connection struct {
UserID string
conn *websocket.Conn
send chan []byte
}
// Send 将消息投递到该连接的发送通道(非阻塞,通道满则丢弃)
func (c *Connection) Send(message []byte) {
select {
case c.send <- message:
default:
// 通道满,丢弃消息避免阻塞 hub
}
}
// Outgoing 返回该连接的发送通道,供写协程消费
func (c *Connection) Outgoing() <-chan []byte {
return c.send
}
// Hub 管理所有在线 WebSocket 连接,按 userID 索引
type Hub struct {
mu sync.RWMutex
clients map[string]map[*Connection]bool
}
// NewHub 创建 Hub 实例
func NewHub() *Hub {
return &Hub{
clients: make(map[string]map[*Connection]bool),
}
}
// ClientCount 返回当前在线连接总数(供 /readyz 探针使用)
func (h *Hub) ClientCount() int {
h.mu.RLock()
defer h.mu.RUnlock()
count := 0
for _, conns := range h.clients {
count += len(conns)
}
return count
}
// Register 注册一个用户连接,返回 Connection 供调用方持有
func (h *Hub) Register(userID string, conn *websocket.Conn) *Connection {
c := &Connection{
UserID: userID,
conn: conn,
send: make(chan []byte, 64),
}
h.mu.Lock()
defer h.mu.Unlock()
if h.clients[userID] == nil {
h.clients[userID] = make(map[*Connection]bool)
}
h.clients[userID][c] = true
return c
}
// Unregister 注销一个用户连接并关闭其发送通道
func (h *Hub) Unregister(c *Connection) {
h.mu.Lock()
defer h.mu.Unlock()
conns, ok := h.clients[c.UserID]
if !ok {
return
}
if _, exists := conns[c]; exists {
delete(conns, c)
close(c.send)
}
if len(conns) == 0 {
delete(h.clients, c.UserID)
}
}
// SendToUser 向指定用户的所有在线连接推送消息
func (h *Hub) SendToUser(userID string, message []byte) error {
h.mu.RLock()
defer h.mu.RUnlock()
conns, ok := h.clients[userID]
if !ok {
return nil
}
for c := range conns {
c.Send(message)
}
return nil
}
// Broadcast 向所有在线连接广播消息
func (h *Hub) Broadcast(message []byte) {
h.mu.RLock()
defer h.mu.RUnlock()
for _, conns := range h.clients {
for c := range conns {
c.Send(message)
}
}
}