Files
rikkei_simple_care/server/internal/syncjobs/jobs.go
2026-06-30 09:31:33 +07:00

583 lines
14 KiB
Go

package syncjobs
import (
"context"
"encoding/json"
"errors"
"log"
"sync"
"time"
"server/internal/models"
internalDb "server/internal/db"
"server/internal/qldt"
"server/internal/util"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
// Trạng thái đồng bộ sinh viên toàn hệ thống
type StudentsSyncStatus struct {
Running bool `json:"running"`
Done bool `json:"done"`
Error string `json:"error,omitempty"`
Page int `json:"page"`
PageSize int `json:"pageSize"`
Total int `json:"total"`
Synced int `json:"synced"`
UpdatedAt int64 `json:"updatedAt"`
}
type StudentsSyncJob struct {
mu sync.RWMutex
status StudentsSyncStatus
}
func NewStudentsSyncJob() *StudentsSyncJob {
return &StudentsSyncJob{
status: StudentsSyncStatus{UpdatedAt: time.Now().Unix()},
}
}
func (j *StudentsSyncJob) Status() StudentsSyncStatus {
j.mu.RLock()
defer j.mu.RUnlock()
return j.status
}
func (j *StudentsSyncJob) Start(db *gorm.DB, client *qldt.Client, token string) bool {
j.mu.Lock()
if j.status.Running {
j.mu.Unlock()
return false
}
j.status = StudentsSyncStatus{
Running: true,
Done: false,
Page: 0,
PageSize: 100,
Total: 0,
Synced: 0,
UpdatedAt: time.Now().Unix(),
}
j.mu.Unlock()
go j.run(db, client, token)
return true
}
func (j *StudentsSyncJob) run(db *gorm.DB, client *qldt.Client, token string) {
defer func() {
j.mu.Lock()
j.status.Running = false
j.status.Done = j.status.Error == ""
j.status.UpdatedAt = time.Now().Unix()
j.mu.Unlock()
}()
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Minute)
defer cancel()
pageSize := 100
page := 1
total := 0
synced := 0
for {
res, err := client.GetStudents(ctx, token, page, pageSize)
if err != nil {
j.setError(err.Error())
return
}
if total == 0 {
total = res.Total
}
if len(res.Data) == 0 {
break
}
batch := make([]models.Student, 0, len(res.Data))
for _, s := range res.Data {
var dob *time.Time
if s.DateOfBirth != nil && *s.DateOfBirth != "" {
if t, err := time.Parse("2006-01-02", *s.DateOfBirth); err == nil {
dob = &t
}
}
avatarVal := ""
if s.Avatar != nil {
avatarVal = util.NormalizeAvatarURL(*s.Avatar)
}
var avatar *string
if avatarVal != "" {
avatar = &avatarVal
}
st := models.Student{
RkID: s.ID,
StudentCode: s.StudentCode,
FullName: s.FullName,
Phone: s.Phone,
Email: s.Email,
DateOfBirth: dob,
Gender: s.Gender,
Status: s.Status,
Location: s.Location,
Avatar: avatar,
}
if s.System != nil {
sysID := s.System.ID
sysName := s.System.Name
st.SystemID = &sysID
st.SystemName = &sysName
}
batch = append(batch, st)
}
if err := db.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "rk_id"}},
DoUpdates: clause.AssignmentColumns([]string{
"student_code",
"full_name",
"phone",
"email",
"date_of_birth",
"gender",
"status",
"location",
"avatar",
"system_id",
"system_name",
"updated_at",
}),
}).CreateInBatches(batch, 100).Error; err != nil {
j.setError("Failed to save students to DB: " + err.Error())
return
}
synced += len(batch)
j.mu.Lock()
j.status.Page = page
j.status.PageSize = pageSize
j.status.Total = total
j.status.Synced = synced
j.status.UpdatedAt = time.Now().Unix()
j.mu.Unlock()
if synced >= total || len(res.Data) < pageSize {
break
}
page++
time.Sleep(50 * time.Millisecond) // Tránh spam API quá tải
}
}
func (j *StudentsSyncJob) setError(msg string) {
j.mu.Lock()
defer j.mu.Unlock()
j.status.Error = msg
j.status.UpdatedAt = time.Now().Unix()
}
// Trạng thái đồng bộ lớp học và sinh viên của lớp
type ClassesSyncStatus struct {
Running bool `json:"running"`
Done bool `json:"done"`
Error string `json:"error,omitempty"`
Total int `json:"total"`
Synced int `json:"synced"`
UpdatedAt int64 `json:"updatedAt"`
}
type ClassesSyncJob struct {
mu sync.RWMutex
status ClassesSyncStatus
}
func NewClassesSyncJob() *ClassesSyncJob {
return &ClassesSyncJob{
status: ClassesSyncStatus{UpdatedAt: time.Now().Unix()},
}
}
func (j *ClassesSyncJob) Status() ClassesSyncStatus {
j.mu.RLock()
defer j.mu.RUnlock()
return j.status
}
func (j *ClassesSyncJob) Start(db *gorm.DB, client *qldt.Client, token string) bool {
j.mu.Lock()
if j.status.Running {
j.mu.Unlock()
return false
}
j.status = ClassesSyncStatus{
Running: true,
Done: false,
Total: 0,
Synced: 0,
UpdatedAt: time.Now().Unix(),
}
j.mu.Unlock()
go j.run(db, client, token)
return true
}
func (j *ClassesSyncJob) run(db *gorm.DB, client *qldt.Client, token string) {
defer func() {
j.mu.Lock()
j.status.Running = false
j.status.Done = j.status.Error == ""
j.status.UpdatedAt = time.Now().Unix()
j.mu.Unlock()
}()
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Minute)
defer cancel()
// 1. Lấy toàn bộ lớp học từ hệ thống chính (System ID = 0 đại diện cho tất cả)
res, err := client.GetClassesBySystem(ctx, token, 0)
if err != nil {
j.setError("Failed to fetch classes: " + err.Error())
return
}
j.mu.Lock()
j.status.Total = len(res.Data)
j.status.UpdatedAt = time.Now().Unix()
j.mu.Unlock()
var syncedClassIDs []int64
for _, c := range res.Data {
syncedClassIDs = append(syncedClassIDs, c.ID)
}
// Đồng bộ từng lớp
for idx, cl := range res.Data {
tx := db.Begin()
if tx.Error != nil {
j.setError("Failed to start transaction: " + tx.Error.Error())
return
}
// Map thông tin class
classRow := models.Class{
RkID: cl.ID,
Name: cl.Name,
ClassCode: cl.ClassCode,
Type: cl.Type,
StudentCount: cl.StudentCount,
}
if cl.Specializes != nil {
specID := cl.Specializes.ID
specName := cl.Specializes.Name
classRow.SpecializeRkID = &specID
classRow.SpecializeName = &specName
if cl.Specializes.Systems != nil {
sysID := cl.Specializes.Systems.ID
sysCode := cl.Specializes.Systems.SystemCode
sysName := cl.Specializes.Systems.Name
classRow.SystemRkID = &sysID
classRow.SystemCode = &sysCode
classRow.SystemName = &sysName
}
}
// Xem lớp học hiện tại trong DB
var existing models.Class
err := tx.Where("rk_id = ?", cl.ID).First(&existing).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
classRow.IsStudying = false // Lớp mới mặc định chưa học
if err := tx.Create(&classRow).Error; err != nil {
tx.Rollback()
j.setError("Failed to insert class: " + err.Error())
return
}
} else if err != nil {
tx.Rollback()
j.setError("Lookup class failed: " + err.Error())
return
} else {
// Giữ trạng thái IsStudying cũ
classRow.ID = existing.ID
classRow.IsStudying = existing.IsStudying
classRow.CreatedAt = existing.CreatedAt
if err := tx.Save(&classRow).Error; err != nil {
tx.Rollback()
j.setError("Failed to update class: " + err.Error())
return
}
}
// Lưu danh sách môn học của lớp (từ API lớp + portal courses)
courseDTOs := make([]internalDb.ClassCourseDTO, 0, len(cl.Courses))
for _, co := range cl.Courses {
if co.ID > 0 {
courseDTOs = append(courseDTOs, internalDb.ClassCourseDTO{
ID: co.ID,
Name: co.Name,
})
}
}
if portalRes, err := client.GetClassCourses(ctx, token, cl.ID); err == nil && len(portalRes.Data) > 0 {
portalDTOs := make([]internalDb.ClassCourseDTO, 0, len(portalRes.Data))
for _, co := range portalRes.Data {
if co.ID > 0 {
portalDTOs = append(portalDTOs, internalDb.ClassCourseDTO{
ID: co.ID,
Name: co.Name,
CourseCode: co.CourseCode,
})
}
}
courseDTOs = internalDb.MergeClassCourses(courseDTOs, portalDTOs)
}
if len(courseDTOs) > 0 {
if err := internalDb.UpsertClassCourses(tx, cl.ID, courseDTOs); err != nil {
tx.Rollback()
j.setError("Failed to save class courses: " + err.Error())
return
}
}
// Đồng bộ sinh viên của lớp: Gọi class dashboard để trích xuất sinh viên
// Ta chọn course đầu tiên của lớp để kéo thông tin sinh viên
var classStudentsToMap []models.ClassStudent
if len(cl.Courses) > 0 && cl.Courses[0].ID > 0 {
dashboardJSON, err := client.GetClassDashboard(ctx, token, cl.ID, cl.Courses[0].ID)
if err != nil {
// Nếu lấy dashboard lỗi thì log lại nhưng không dừng tiến trình sync lớp
log.Printf("Warning: failed to get dashboard for class ID %d, course ID %d: %v", cl.ID, cl.Courses[0].ID, err)
} else {
stList, err := extractStudentsFromDashboard(dashboardJSON)
if err != nil {
log.Printf("Warning: failed to parse students from dashboard class %d: %v", cl.ID, err)
} else if len(stList) > 0 {
// Upsert sinh viên vào table students
if err := tx.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "rk_id"}},
DoUpdates: clause.AssignmentColumns([]string{
"student_code",
"full_name",
"phone",
"email",
"date_of_birth",
"gender",
"status",
"location",
"avatar",
"system_id",
"system_name",
"updated_at",
}),
}).Create(&stList).Error; err != nil {
tx.Rollback()
j.setError("Failed to save class students: " + err.Error())
return
}
// Lưu map Class-Student
for _, st := range stList {
classStudentsToMap = append(classStudentsToMap, models.ClassStudent{
ClassRkID: cl.ID,
StudentRkID: st.RkID,
})
}
}
}
}
// Cập nhật map Class-Student trong DB (xóa map cũ, insert map mới)
if err := tx.Where("class_rk_id = ?", cl.ID).Delete(&models.ClassStudent{}).Error; err != nil {
tx.Rollback()
j.setError("Failed to clear old class students mapping: " + err.Error())
return
}
if len(classStudentsToMap) > 0 {
if err := tx.Clauses(clause.OnConflict{
DoNothing: true, // Tránh trùng lặp do lỗi nào đó
}).Create(&classStudentsToMap).Error; err != nil {
tx.Rollback()
j.setError("Failed to map class students: " + err.Error())
return
}
}
if err := tx.Commit().Error; err != nil {
j.setError("Transaction commit failed: " + err.Error())
return
}
j.mu.Lock()
j.status.Synced = idx + 1
j.status.UpdatedAt = time.Now().Unix()
j.mu.Unlock()
time.Sleep(50 * time.Millisecond) // Thở tí tránh spam API
}
// 2. Xóa các lớp học không còn tồn tại trên QLĐT
if len(syncedClassIDs) > 0 {
db.Where("rk_id NOT IN ?", syncedClassIDs).Delete(&models.Class{})
} else {
db.Session(&gorm.Session{AllowGlobalUpdate: true}).Delete(&models.Class{})
}
}
func (j *ClassesSyncJob) setError(msg string) {
j.mu.Lock()
defer j.mu.Unlock()
j.status.Error = msg
j.status.UpdatedAt = time.Now().Unix()
}
// Helper giải nén danh sách sinh viên từ Dashboard JSON
func extractStudentsFromDashboard(payloadJSON []byte) ([]models.Student, error) {
var root interface{}
if err := json.Unmarshal(payloadJSON, &root); err != nil {
return nil, err
}
var arr []interface{}
if m, ok := root.(map[string]interface{}); ok {
if dataVal, exists := m["data"]; exists {
if dArr, ok := dataVal.([]interface{}); ok {
arr = dArr
} else if dMap, ok := dataVal.(map[string]interface{}); ok {
for _, k := range []string{"studentsInClass", "students", "rows", "items"} {
if val, exists := dMap[k]; exists {
if innerArr, ok := val.([]interface{}); ok {
arr = innerArr
break
}
}
}
}
}
if len(arr) == 0 {
for _, k := range []string{"studentsInClass", "students", "rows", "items"} {
if val, exists := m[k]; exists {
if dArr, ok := val.([]interface{}); ok {
arr = dArr
break
}
}
}
}
} else if sliceVal, ok := root.([]interface{}); ok {
arr = sliceVal
}
var students []models.Student
for _, item := range arr {
m, ok := item.(map[string]interface{})
if !ok || m == nil {
continue
}
studentObj, ok := m["student"].(map[string]interface{})
if !ok || studentObj == nil {
continue
}
var rkID int64
if idVal, exists := studentObj["id"]; exists {
switch v := idVal.(type) {
case float64:
rkID = int64(v)
case json.Number:
if id64, err := v.Int64(); err == nil {
rkID = id64
}
case int:
rkID = int64(v)
case int64:
rkID = v
}
}
if rkID == 0 {
continue
}
stCode, _ := studentObj["studentCode"].(string)
fullName, _ := studentObj["fullName"].(string)
email, _ := studentObj["email"].(string)
var phone *string
if p, ok := studentObj["phone"].(string); ok && p != "" {
phone = &p
}
var dob *time.Time
if dobStr, ok := studentObj["dateOfBirth"].(string); ok && dobStr != "" {
if t, err := time.Parse("2006-01-02", dobStr); err == nil {
dob = &t
}
}
var gender *int
if gVal, exists := studentObj["gender"]; exists {
var g int
switch v := gVal.(type) {
case float64:
g = int(v)
case json.Number:
if i64, err := v.Int64(); err == nil {
g = int(i64)
}
case int:
g = v
}
gender = &g
}
status, _ := studentObj["status"].(string)
location, _ := studentObj["location"].(string)
var avatar *string
if av := util.PickAvatarFromMap(studentObj); av != "" {
avatar = &av
}
student := models.Student{
RkID: rkID,
StudentCode: stCode,
FullName: fullName,
Phone: phone,
Email: email,
DateOfBirth: dob,
Gender: gender,
Status: &status,
Location: &location,
Avatar: avatar,
}
if sysVal, ok := studentObj["system"].(map[string]interface{}); ok && sysVal != nil {
var sysID int64
if sID, exists := sysVal["id"]; exists {
switch v := sID.(type) {
case float64:
sysID = int64(v)
case int64:
sysID = v
}
}
sysName, _ := sysVal["name"].(string)
if sysID > 0 {
student.SystemID = &sysID
student.SystemName = &sysName
}
}
students = append(students, student)
}
return students, nil
}