Files
rikkei_simple_care/server/internal/websocket/websocket.go
2026-06-30 10:39:29 +07:00

342 lines
8.5 KiB
Go

package websocket
import (
"encoding/json"
"log"
"strconv"
"sync"
"github.com/gofiber/websocket/v2"
"gorm.io/gorm"
internalDb "server/internal/db"
)
type SocketMsg struct {
Event string `json:"event"`
Data map[string]any `json:"data"`
}
type SocketClient struct {
Conn *websocket.Conn
StudentID int64
ClassID int64
StaffID uint
Role string // "student" | "teacher"
Addr string
}
type WsHub struct {
mu sync.RWMutex
students map[int64]*SocketClient
teachers map[string]*SocketClient
teachersByStaff map[uint][]string // staffId -> teacher connection addresses
subscribers map[int64][]string // studentId -> list of teacher connection addresses
}
var Hub = &WsHub{
students: make(map[int64]*SocketClient),
teachers: make(map[string]*SocketClient),
teachersByStaff: make(map[uint][]string),
subscribers: make(map[int64][]string),
}
func (h *WsHub) IsStudentOnline(studentRkID int64) bool {
h.mu.RLock()
defer h.mu.RUnlock()
_, ok := h.students[studentRkID]
return ok
}
func (h *WsHub) PushChatToStudent(studentRkID int64, data map[string]any) {
h.mu.RLock()
client, ok := h.students[studentRkID]
h.mu.RUnlock()
if !ok || client == nil {
return
}
_ = client.Conn.WriteJSON(SocketMsg{Event: "chat:message", Data: data})
}
func (h *WsHub) PushChatToStaff(staffID uint, data map[string]any) {
h.mu.RLock()
addrs := append([]string(nil), h.teachersByStaff[staffID]...)
h.mu.RUnlock()
msg := SocketMsg{Event: "chat:message", Data: data}
for _, addr := range addrs {
h.mu.RLock()
t, found := h.teachers[addr]
h.mu.RUnlock()
if found && t != nil {
_ = t.Conn.WriteJSON(msg)
}
}
}
func (h *WsHub) GetOnlineStudentIDs(classID int64) []int64 {
h.mu.RLock()
defer h.mu.RUnlock()
var ids []int64
for _, client := range h.students {
if client.ClassID == classID {
ids = append(ids, client.StudentID)
}
}
if ids == nil {
return []int64{}
}
return ids
}
func (h *WsHub) Register(c *SocketClient) {
h.mu.Lock()
defer h.mu.Unlock()
c.Addr = c.Conn.RemoteAddr().String()
if c.Role == "student" {
h.students[c.StudentID] = c
log.Printf("[WS] Student %d registered (Address: %s, Class: %d)", c.StudentID, c.Addr, c.ClassID)
} else if c.Role == "teacher" {
h.teachers[c.Addr] = c
if c.StaffID > 0 {
h.teachersByStaff[c.StaffID] = append(h.teachersByStaff[c.StaffID], c.Addr)
}
log.Printf("[WS] Teacher registered (Address: %s, Staff: %d)", c.Addr, c.StaffID)
}
}
func (h *WsHub) Unregister(c *SocketClient) {
h.mu.Lock()
defer h.mu.Unlock()
if c.Role == "student" {
delete(h.students, c.StudentID)
log.Printf("[WS] Student %d disconnected", c.StudentID)
// Báo cho các giáo viên đang xem là stream của học sinh đã dừng
if teachers, exists := h.subscribers[c.StudentID]; exists {
for _, tAddr := range teachers {
if t, found := h.teachers[tAddr]; found {
_ = t.Conn.WriteJSON(SocketMsg{
Event: "teacher:stream-stopped",
Data: map[string]any{"studentId": c.StudentID},
})
}
}
delete(h.subscribers, c.StudentID)
}
} else if c.Role == "teacher" {
delete(h.teachers, c.Addr)
if c.StaffID > 0 {
list := h.teachersByStaff[c.StaffID]
next := list[:0]
for _, addr := range list {
if addr != c.Addr {
next = append(next, addr)
}
}
if len(next) == 0 {
delete(h.teachersByStaff, c.StaffID)
} else {
h.teachersByStaff[c.StaffID] = next
}
}
log.Printf("[WS] Teacher %s disconnected", c.Addr)
// Dọn dẹp subscriptions của giáo viên này
for sID, teachersList := range h.subscribers {
newList := []string{}
for _, addr := range teachersList {
if addr != c.Addr {
newList = append(newList, addr)
}
}
if len(newList) == 0 {
delete(h.subscribers, sID)
// Nếu không còn ai xem học sinh này, gửi lệnh tắt camera/screen cho client học sinh
if student, exists := h.students[sID]; exists {
_ = student.Conn.WriteJSON(SocketMsg{Event: "stop_screenshot_stream"})
_ = student.Conn.WriteJSON(SocketMsg{Event: "stop_webcam_stream"})
}
} else {
h.subscribers[sID] = newList
}
}
}
}
// Teacher bắt đầu xem stream của Student
func (h *WsHub) Subscribe(teacherAddr string, studentID int64) {
h.mu.Lock()
defer h.mu.Unlock()
// Thêm giáo viên vào list người xem của học sinh
teachersList := h.subscribers[studentID]
alreadySubscribed := false
for _, addr := range teachersList {
if addr == teacherAddr {
alreadySubscribed = true
break
}
}
if !alreadySubscribed {
h.subscribers[studentID] = append(teachersList, teacherAddr)
log.Printf("[WS] Teacher %s subscribed to student %d stream", teacherAddr, studentID)
}
// Phát lệnh cho máy học sinh bật stream (nếu học sinh đang online)
if student, exists := h.students[studentID]; exists {
_ = student.Conn.WriteJSON(SocketMsg{Event: "start_screenshot_stream"})
_ = student.Conn.WriteJSON(SocketMsg{Event: "start_webcam_stream"})
}
}
// Teacher dừng xem stream của Student
func (h *WsHub) Unsubscribe(teacherAddr string, studentID int64) {
h.mu.Lock()
defer h.mu.Unlock()
teachersList, exists := h.subscribers[studentID]
if !exists {
return
}
newList := []string{}
for _, addr := range teachersList {
if addr != teacherAddr {
newList = append(newList, addr)
}
}
if len(newList) == 0 {
delete(h.subscribers, studentID)
log.Printf("[WS] Student %d has no more proctor subscribers. Stopping streams.", studentID)
// Báo học sinh tắt camera & screen stream để tiết kiệm mạng và CPU
if student, exists := h.students[studentID]; exists {
_ = student.Conn.WriteJSON(SocketMsg{Event: "stop_screenshot_stream"})
_ = student.Conn.WriteJSON(SocketMsg{Event: "stop_webcam_stream"})
}
} else {
h.subscribers[studentID] = newList
}
}
// Chuyển tiếp frame ảnh từ Student đến các Teacher đã subscribe
func (h *WsHub) RelayFrame(studentID int64, event string, data map[string]any) {
h.mu.RLock()
defer h.mu.RUnlock()
teachersList, exists := h.subscribers[studentID]
if !exists || len(teachersList) == 0 {
return
}
relayEvent := "teacher:screenshot-stream-frame"
if event == "webcam_stream_frame" {
relayEvent = "teacher:webcam-stream-frame"
}
msg := SocketMsg{
Event: relayEvent,
Data: map[string]any{
"studentId": studentID,
"imageBuffer": data["imageBuffer"],
},
}
for _, addr := range teachersList {
if t, found := h.teachers[addr]; found {
_ = t.Conn.WriteJSON(msg)
}
}
}
// WebSocket handler cho Fiber route
func WebSocketHandler(db *gorm.DB) func(*websocket.Conn) {
return func(c *websocket.Conn) {
role := c.Query("role", "student")
studentIDStr := c.Query("studentId", "0")
classIDStr := c.Query("classId", "0")
staffIDStr := c.Query("staffId", "0")
studentID, _ := strconv.ParseInt(studentIDStr, 10, 64)
classID, _ := strconv.ParseInt(classIDStr, 10, 64)
staffID64, _ := strconv.ParseUint(staffIDStr, 10, 64)
if role == "student" && db != nil {
resolvedClassID := internalDb.FindActiveClassForStudent(db, studentID)
if resolvedClassID > 0 {
classID = resolvedClassID
}
}
client := &SocketClient{
Conn: c,
StudentID: studentID,
ClassID: classID,
StaffID: uint(staffID64),
Role: role,
}
Hub.Register(client)
defer func() {
Hub.Unregister(client)
c.Close()
}()
for {
_, msgBytes, err := c.ReadMessage()
if err != nil {
break
}
var msg SocketMsg
if err := json.Unmarshal(msgBytes, &msg); err != nil {
continue
}
// Xử lý các sự kiện
switch msg.Event {
case "screenshot_stream_frame", "webcam_stream_frame":
// Nhận frame từ học sinh, chuyển tiếp về các thầy cô
Hub.RelayFrame(client.StudentID, msg.Event, msg.Data)
case "teacher:subscribe":
// Giáo viên đăng ký xem học sinh cụ thể
if client.Role == "teacher" {
if sIDVal, ok := msg.Data["studentId"]; ok {
var sID int64
switch v := sIDVal.(type) {
case float64:
sID = int64(v)
case string:
sID, _ = strconv.ParseInt(v, 10, 64)
}
if sID > 0 {
Hub.Subscribe(client.Addr, sID)
}
}
}
case "teacher:unsubscribe":
// Giáo viên hủy đăng ký
if client.Role == "teacher" {
if sIDVal, ok := msg.Data["studentId"]; ok {
var sID int64
switch v := sIDVal.(type) {
case float64:
sID = int64(v)
case string:
sID, _ = strconv.ParseInt(v, 10, 64)
}
if sID > 0 {
Hub.Unsubscribe(client.Addr, sID)
}
}
}
}
}
}
}