mengstack-api/internal/kernel/websocket/hub.go
MengStack Dev df809f1045
Some checks failed
CI / Build & Test (push) Failing after 1m31s
feat: user CRUD API + dashboard stats + kernel infrastructure
- Add user management endpoints (list/create/get/update/delete) with pagination and search
- Add dashboard stats endpoint with tenant/user/online counts and growth metrics
- Add tenant resolver middleware for multi-tenant request scoping
- Add i18n kernel with zh/en message files and AcceptLanguage middleware
- Add WebSocket hub/handler for real-time communication
- Add job scheduler kernel with cron support
- Add plugin sandbox for isolated execution
- Add storage kernel (local filesystem)
- Add event bus kernel for pub/sub
- Add cache kernel abstraction
- Add database migration runner and version upgrade checker
- Add rate limiting middleware with Redis backend
- Add SQL migrations for rbac, audit_logs, settings, notifications, examples
- Extend user repository with list/delete/count operations
- Register all module routes with tenant resolver
2026-10-03 03:42:58 +08:00

148 lines
2.6 KiB
Go

package websocket
import (
"encoding/json"
"sync"
"time"
"go.uber.org/zap"
)
type Message struct {
Type string `json:"type"`
Payload interface{} `json:"payload"`
}
type Client struct {
hub *Hub
userID uint
tenantID string
send chan []byte
done chan struct{}
}
type Hub struct {
mu sync.RWMutex
clients map[uint]map[*Client]bool
log *zap.Logger
register chan *Client
unreg chan *Client
}
func NewHub(log *zap.Logger) *Hub {
h := &Hub{
clients: make(map[uint]map[*Client]bool),
log: log,
register: make(chan *Client, 64),
unreg: make(chan *Client, 64),
}
go h.run()
return h
}
func (h *Hub) run() {
for {
select {
case c := <-h.register:
h.mu.Lock()
if h.clients[c.userID] == nil {
h.clients[c.userID] = make(map[*Client]bool)
}
h.clients[c.userID][c] = true
h.mu.Unlock()
h.log.Debug("websocket client connected", zap.Uint("user_id", c.userID))
case c := <-h.unreg:
h.mu.Lock()
if m := h.clients[c.userID]; m != nil {
delete(m, c)
if len(m) == 0 {
delete(h.clients, c.userID)
}
}
h.mu.Unlock()
close(c.send)
h.log.Debug("websocket client disconnected", zap.Uint("user_id", c.userID))
}
}
}
func (h *Hub) Register(c *Client) {
h.register <- c
}
func (h *Hub) Unregister(c *Client) {
h.unreg <- c
}
func (h *Hub) SendToUser(userID uint, msg Message) {
data, err := json.Marshal(msg)
if err != nil {
h.log.Error("websocket marshal", zap.Error(err))
return
}
h.mu.RLock()
defer h.mu.RUnlock()
for c := range h.clients[userID] {
select {
case c.send <- data:
default:
go func(c *Client) {
h.unreg <- c
}(c)
}
}
}
func (h *Hub) Broadcast(msg Message) {
data, err := json.Marshal(msg)
if err != nil {
h.log.Error("websocket marshal", zap.Error(err))
return
}
h.mu.RLock()
defer h.mu.RUnlock()
for _, clients := range h.clients {
for c := range clients {
select {
case c.send <- data:
default:
go func(c *Client) {
h.unreg <- c
}(c)
}
}
}
}
func (h *Hub) ConnectedCount() int {
h.mu.RLock()
defer h.mu.RUnlock()
n := 0
for _, clients := range h.clients {
n += len(clients)
}
return n
}
func NewClient(hub *Hub, userID uint, tenantID string) *Client {
return &Client{
hub: hub,
userID: userID,
tenantID: tenantID,
send: make(chan []byte, 256),
done: make(chan struct{}),
}
}
func (c *Client) Done() <-chan struct{} {
return c.done
}
const (
writeWait = 10 * time.Second
pongWait = 60 * time.Second
pingPeriod = (pongWait * 9) / 10
maxMessageSize = 4096
)