Browse Source

fix:任务执行超时修复

master
刘浩东 4 weeks ago
parent
commit
d58102289d
  1. 1
      model/task.go
  2. 20
      mqtt/mqtt.go
  3. 2
      service/account_resource_service.go
  4. 2
      service/alarm_service.go
  5. 8
      service/billing_service.go
  6. 65
      service/command_service.go
  7. 2
      service/dock_service.go
  8. 2
      service/drone_service.go
  9. 2
      service/execution_service.go
  10. 4
      service/firmware_service.go
  11. 6
      service/invoice_service.go
  12. 8
      service/live_service.go
  13. 2
      service/operation_log_service.go
  14. 2
      service/route_service.go
  15. 2
      service/system_service.go
  16. 2
      service/task_service.go
  17. 2
      service/user_service.go
  18. 2
      service/video_service.go
  19. 2
      service/workflow_service.go
  20. 1
      sql/001_schema.sql
  21. 2
      sql/010_task_workflow_timeout.sql
  22. 54
      tool/snowflake.go

1
model/task.go

@ -35,6 +35,7 @@ type TaskExecution struct {
StartTime *time.Time `gorm:"column:start_time" json:"startTime"` StartTime *time.Time `gorm:"column:start_time" json:"startTime"`
EndTime *time.Time `gorm:"column:end_time" json:"endTime"` EndTime *time.Time `gorm:"column:end_time" json:"endTime"`
Status string `gorm:"column:status;type:VARCHAR(16);default:pending" json:"status"` Status string `gorm:"column:status;type:VARCHAR(16);default:pending" json:"status"`
ResultCode string `gorm:"column:result_code;type:VARCHAR(64)" json:"resultCode"`
TrajectoryJSON json.RawMessage `gorm:"column:trajectory_json;type:JSON" json:"trajectoryJson"` TrajectoryJSON json.RawMessage `gorm:"column:trajectory_json;type:JSON" json:"trajectoryJson"`
CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"`
} }

20
mqtt/mqtt.go

@ -55,17 +55,35 @@ func Subscribe(topics map[string]byte, handler mqtt.MessageHandler) {
func resubscribe(c mqtt.Client) { func resubscribe(c mqtt.Client) {
for topic, qos := range subTopics { for topic, qos := range subTopics {
if token := c.Subscribe(topic, qos, subHandler); token.Wait() && token.Error() != nil {
if token := c.Subscribe(topic, qos, inboundHandler); token.Wait() && token.Error() != nil {
logger.ERROR("订阅失败 topic="+topic, token.Error()) logger.ERROR("订阅失败 topic="+topic, token.Error())
} }
} }
} }
func inboundHandler(c mqtt.Client, msg mqtt.Message) {
//logger.INFO("MQTT 收到消息",
// "topic="+msg.Topic(),
// "qos=", msg.Qos(),
// "retained=", msg.Retained(),
// "payload="+string(msg.Payload()),
//)
if subHandler != nil {
subHandler(c, msg)
}
}
// Publish 发布消息到指定 topic // Publish 发布消息到指定 topic
func Publish(topic string, qos byte, retained bool, payload []byte) error { func Publish(topic string, qos byte, retained bool, payload []byte) error {
if client == nil || !client.IsConnected() { if client == nil || !client.IsConnected() {
return errors.New("mqtt client not connected") return errors.New("mqtt client not connected")
} }
logger.INFO("MQTT 发送消息",
"topic="+topic,
"qos=", qos,
"retained=", retained,
"payload="+string(payload),
)
token := client.Publish(topic, qos, retained, payload) token := client.Publish(topic, qos, retained, payload)
token.Wait() token.Wait()
return token.Error() return token.Error()

2
service/account_resource_service.go

@ -24,7 +24,7 @@ func (s *AccountService) GetDownloadPage(userID int64, req *vo.UsagePageReq) (*c
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var rows []row var rows []row
if err := db.Scopes(req.Paginate).Order("download_log.id DESC").Find(&rows).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("download_log.created_at DESC, download_log.id DESC").Find(&rows).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
list := make([]vo.DownloadRecordVO, 0, len(rows)) list := make([]vo.DownloadRecordVO, 0, len(rows))

2
service/alarm_service.go

@ -44,7 +44,7 @@ func (s *AlarmService) GetPage(userID int64, isAdmin bool, req *vo.AlarmPageReq)
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.Alarm var list []model.Alarm
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("triggered_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询告警列表失败", err) logger.ERROR("查询告警列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

8
service/billing_service.go

@ -173,7 +173,7 @@ func (b *BillingService) GetUsagePage(userID int64, req *vo.UsagePageReq) (*comm
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.TrafficUsageLog var list []model.TrafficUsageLog
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return common.Page(req.Pagination, total, list), nil return common.Page(req.Pagination, total, list), nil
@ -190,7 +190,7 @@ func (b *BillingService) GetOrderPage(userID int64, req *vo.OrderPageReq) (*comm
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.TrafficOrder var list []model.TrafficOrder
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return common.Page(req.Pagination, total, list), nil return common.Page(req.Pagination, total, list), nil
@ -263,7 +263,7 @@ func (b *BillingService) ListSimCards(userID int64, req *vo.SimCardPageReq) (*co
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.SimCard var list []model.SimCard
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("updated_at DESC, id DESC").Find(&list).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return common.Page(req.Pagination, total, list), nil return common.Page(req.Pagination, total, list), nil
@ -283,7 +283,7 @@ func (b *BillingService) GetSimRechargeLogPage(userID int64, req *vo.SimRecharge
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.SimRechargeLog var list []model.SimRechargeLog
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return common.Page(req.Pagination, total, list), nil return common.Page(req.Pagination, total, list), nil

65
service/command_service.go

@ -24,7 +24,11 @@ type CommandService struct{}
var DefaultCommandService = &CommandService{} var DefaultCommandService = &CommandService{}
const maxRetryCount = 2
const (
maxRetryCount = 2
workflowStartTimeout = 2 * time.Minute
workflowProgressTimeout = 5 * time.Minute
)
// Dispatch 下发指令:校验在线 → 落库 → MQTT 发布 // Dispatch 下发指令:校验在线 → 落库 → MQTT 发布
func (s *CommandService) Dispatch(userID int64, isAdmin bool, id int64, req *vo.CommandReq) (*model.DeviceCommandLog, *common.BusiError) { func (s *CommandService) Dispatch(userID int64, isAdmin bool, id int64, req *vo.CommandReq) (*model.DeviceCommandLog, *common.BusiError) {
@ -133,9 +137,16 @@ func failTaskExecution(cmd model.DeviceCommandLog, reason string) {
if reason == "" { if reason == "" {
reason = "COMMAND_REJECTED" reason = "COMMAND_REJECTED"
} }
_ = common.DB.Model(&model.TaskExecution{}).
result := common.DB.Model(&model.TaskExecution{}).
Where("dock_id = ? AND command_id = ? AND status IN ?", cmd.DockID, strconv.FormatInt(cmd.ID, 10), []string{"pending", "running"}). Where("dock_id = ? AND command_id = ? AND status IN ?", cmd.DockID, strconv.FormatInt(cmd.ID, 10), []string{"pending", "running"}).
Updates(map[string]any{"status": "failed", "end_time": time.Now()})
Updates(map[string]any{"status": "failed", "result_code": reason, "end_time": time.Now()})
if result.Error != nil {
logger.ERROR("更新任务执行失败状态失败", result.Error)
return
}
if result.RowsAffected > 0 {
DefaultTrajectoryStore.Finalize(cmd.DockID, "")
}
} }
// DispatchToDock 持久化并下发已完成权限与在线校验的设备指令。 // DispatchToDock 持久化并下发已完成权限与在线校验的设备指令。
@ -232,6 +243,10 @@ func (s *CommandService) retryTimeout() {
if cmd.AckedAt == nil { if cmd.AckedAt == nil {
continue continue
} }
if cmd.CommandType == "workflow.start_task" {
s.checkWorkflowTimeout(cmd, now)
continue
}
ttl := int64(cmd.TTLMs) ttl := int64(cmd.TTLMs)
if ttl <= 0 { if ttl <= 0 {
ttl = 30000 ttl = 30000
@ -242,6 +257,48 @@ func (s *CommandService) retryTimeout() {
} }
} }
func (s *CommandService) checkWorkflowTimeout(cmd model.DeviceCommandLog, now time.Time) {
commandID := strconv.FormatInt(cmd.ID, 10)
var executions []model.TaskExecution
executionQuery := common.DB.Where("dock_id = ? AND command_id = ? AND status IN ?", cmd.DockID, commandID, []string{"pending", "running"}).Find(&executions)
if executionQuery.Error != nil {
logger.ERROR("查询工作流执行记录失败", executionQuery.Error)
return
}
if executionQuery.RowsAffected == 0 {
return
}
var workflow model.WorkflowState
err := common.DB.Where("dock_id = ? AND command_id = ?", cmd.DockID, commandID).First(&workflow).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
if now.Sub(*cmd.AckedAt) > workflowStartTimeout {
s.timeoutWorkflow(cmd, "WORKFLOW_START_TIMEOUT")
}
return
}
if err != nil {
logger.ERROR("查询工作流状态失败", err)
return
}
if workflow.State == "running" && now.Sub(workflow.UpdatedAt) > workflowProgressTimeout {
s.timeoutWorkflow(cmd, "WORKFLOW_PROGRESS_TIMEOUT")
}
}
func (s *CommandService) timeoutWorkflow(cmd model.DeviceCommandLog, reason string) {
result := common.DB.Model(&model.DeviceCommandLog{}).
Where("id = ? AND status = ?", cmd.ID, "acked").
Updates(map[string]any{"status": "terminal", "ack_result_code": reason})
if result.Error != nil {
logger.ERROR("更新工作流超时指令失败", result.Error)
return
}
if result.RowsAffected > 0 {
failTaskExecution(cmd, reason)
}
}
// commandEnvelope 构造下行指令通用外层(文档 §5.1) // commandEnvelope 构造下行指令通用外层(文档 §5.1)
func (s *CommandService) commandEnvelope(cmd *model.DeviceCommandLog, requestID string, params map[string]any) *mqtt.Envelope { func (s *CommandService) commandEnvelope(cmd *model.DeviceCommandLog, requestID string, params map[string]any) *mqtt.Envelope {
return mqtt.NewEnvelope(requestID, cmd.DockID, cmd.DroneSN, map[string]any{ return mqtt.NewEnvelope(requestID, cmd.DockID, cmd.DroneSN, map[string]any{
@ -284,7 +341,7 @@ func (s *CommandService) GetPage(userID int64, isAdmin bool, req *vo.CommandPage
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.DeviceCommandLog var list []model.DeviceCommandLog
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("sent_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询指令日志失败", err) logger.ERROR("查询指令日志失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/dock_service.go

@ -53,7 +53,7 @@ func (s *DockService) GetPage(userID int64, isAdmin bool, req *vo.DockPageReq) (
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.Dock var list []model.Dock
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询机巢列表失败", err) logger.ERROR("查询机巢列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/drone_service.go

@ -32,7 +32,7 @@ func (s *DroneService) GetPage(userID int64, isAdmin bool, req *vo.DronePageReq)
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.Drone var list []model.Drone
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询无人机列表失败", err) logger.ERROR("查询无人机列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/execution_service.go

@ -35,7 +35,7 @@ func (s *ExecutionService) GetPage(userID int64, isAdmin bool, req *vo.Execution
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.TaskExecution var list []model.TaskExecution
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询执行记录失败", err) logger.ERROR("查询执行记录失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

4
service/firmware_service.go

@ -44,7 +44,7 @@ func (s *FirmwareService) GetPage(req *vo.FirmwarePageReq) (*common.PageResponse
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.Firmware var list []model.Firmware
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询固件列表失败", err) logger.ERROR("查询固件列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }
@ -54,7 +54,7 @@ func (s *FirmwareService) GetPage(req *vo.FirmwarePageReq) (*common.PageResponse
// GetReleasedList 已发布固件列表(供普通用户选择升级) // GetReleasedList 已发布固件列表(供普通用户选择升级)
func (s *FirmwareService) GetReleasedList() ([]model.Firmware, *common.BusiError) { func (s *FirmwareService) GetReleasedList() ([]model.Firmware, *common.BusiError) {
var list []model.Firmware var list []model.Firmware
if err := common.DB.Where("status = ?", "released").Order("id DESC").Find(&list).Error; err != nil {
if err := common.DB.Where("status = ?", "released").Order("updated_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询已发布固件失败", err) logger.ERROR("查询已发布固件失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

6
service/invoice_service.go

@ -18,7 +18,7 @@ var DefaultInvoiceService = &InvoiceService{}
func (s *InvoiceService) ListProfiles(userID int64) ([]model.InvoiceProfile, *common.BusiError) { func (s *InvoiceService) ListProfiles(userID int64) ([]model.InvoiceProfile, *common.BusiError) {
var profiles []model.InvoiceProfile var profiles []model.InvoiceProfile
if err := common.DB.Where("user_id = ?", userID).Order("is_default DESC, id DESC").Find(&profiles).Error; err != nil {
if err := common.DB.Where("user_id = ?", userID).Order("is_default DESC, updated_at DESC, id DESC").Find(&profiles).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return profiles, nil return profiles, nil
@ -67,7 +67,7 @@ func (s *InvoiceService) DeleteProfile(userID, profileID int64) *common.BusiErro
func (s *InvoiceService) ListEligibleOrders(userID int64) ([]model.TrafficOrder, *common.BusiError) { func (s *InvoiceService) ListEligibleOrders(userID int64) ([]model.TrafficOrder, *common.BusiError) {
var orders []model.TrafficOrder var orders []model.TrafficOrder
if err := common.DB.Where("user_id = ? AND pay_status = ? AND NOT EXISTS (SELECT 1 FROM invoice_request WHERE invoice_request.order_id = traffic_order.id)", userID, "paid").Order("paid_at DESC").Find(&orders).Error; err != nil {
if err := common.DB.Where("user_id = ? AND pay_status = ? AND NOT EXISTS (SELECT 1 FROM invoice_request WHERE invoice_request.order_id = traffic_order.id)", userID, "paid").Order("paid_at DESC, id DESC").Find(&orders).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return orders, nil return orders, nil
@ -104,7 +104,7 @@ func (s *InvoiceService) CreateRequest(userID int64, req *vo.InvoiceRequestCreat
func (s *InvoiceService) ListRequests(userID int64) ([]model.InvoiceRequest, *common.BusiError) { func (s *InvoiceService) ListRequests(userID int64) ([]model.InvoiceRequest, *common.BusiError) {
var requests []model.InvoiceRequest var requests []model.InvoiceRequest
if err := common.DB.Where("user_id = ?", userID).Order("requested_at DESC").Find(&requests).Error; err != nil {
if err := common.DB.Where("user_id = ?", userID).Order("requested_at DESC, id DESC").Find(&requests).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return requests, nil return requests, nil

8
service/live_service.go

@ -63,7 +63,7 @@ func (s *LiveService) Join(userID int64, isAdmin bool, dockID string, req *vo.Li
if err := tx.Set("gorm:query_option", "FOR UPDATE").Where("dock_id = ?", dockID).First(&model.Dock{}).Error; err != nil { if err := tx.Set("gorm:query_option", "FOR UPDATE").Where("dock_id = ?", dockID).First(&model.Dock{}).Error; err != nil {
return err return err
} }
if err := tx.Where("dock_id = ? AND phase IN ('starting','streaming','reconnecting','stopping')", dockID).Order("created_at DESC").First(&session).Error; err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
if err := tx.Where("dock_id = ? AND phase IN ('starting','streaming','reconnecting','stopping')", dockID).Order("created_at DESC, id DESC").First(&session).Error; err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
return err return err
} }
if session.ID == "" { if session.ID == "" {
@ -230,7 +230,7 @@ func (s *LiveService) Stop(userID int64, isAdmin bool, dockID string) *common.Bu
return common.ErrInternal return common.ErrInternal
} }
var active model.LiveSession var active model.LiveSession
if err := common.DB.Where("dock_id = ? AND phase IN ('starting','streaming','reconnecting','stopping')", dockID).Order("id DESC").First(&active).Error; err != nil {
if err := common.DB.Where("dock_id = ? AND phase IN ('starting','streaming','reconnecting','stopping')", dockID).Order("created_at DESC, id DESC").First(&active).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) { if errors.Is(err, gorm.ErrRecordNotFound) {
return common.ErrLiveNotFound return common.ErrLiveNotFound
} }
@ -274,7 +274,7 @@ func (s *LiveService) GetPage(userID int64, isAdmin bool, req *vo.LivePageReq) (
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.LiveSession var list []model.LiveSession
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return common.Page(req.Pagination, total, list), nil return common.Page(req.Pagination, total, list), nil
@ -287,7 +287,7 @@ func (s *LiveService) GetPlayURL(userID int64, isAdmin bool, dockID string) (*vo
return nil, busiErr return nil, busiErr
} }
var session model.LiveSession var session model.LiveSession
if err := common.DB.Scopes(withDockFilter(userID, isAdmin)).Where("dock_id = ? AND phase = 'streaming'", dockID).Order("id DESC").First(&session).Error; err != nil {
if err := common.DB.Scopes(withDockFilter(userID, isAdmin)).Where("dock_id = ? AND phase = 'streaming'", dockID).Order("created_at DESC, id DESC").First(&session).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) { if errors.Is(err, gorm.ErrRecordNotFound) {
return nil, common.ErrLiveNotFound return nil, common.ErrLiveNotFound
} }

2
service/operation_log_service.go

@ -65,7 +65,7 @@ func (s *OperationLogService) GetPage(req *vo.OperationLogPageReq) (*common.Page
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.OperationLog var list []model.OperationLog
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询操作日志失败", err) logger.ERROR("查询操作日志失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/route_service.go

@ -31,7 +31,7 @@ func (s *RouteService) GetPage(userID int64, isAdmin bool, req *vo.RoutePageReq)
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.Route var list []model.Route
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询航线列表失败", err) logger.ERROR("查询航线列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/system_service.go

@ -38,7 +38,7 @@ func (s *SystemService) GetUserPage(req *vo.UserPageReq) (*common.PageResponse[m
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.User var list []model.User
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
logger.ERROR("查询用户列表失败", err) logger.ERROR("查询用户列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/task_service.go

@ -43,7 +43,7 @@ func (s *TaskService) GetPage(userID int64, isAdmin bool, req *vo.TaskPageReq) (
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var tasks []model.TaskPlan var tasks []model.TaskPlan
if err := db.Scopes(req.Paginate).Order("created_at DESC").Find(&tasks).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&tasks).Error; err != nil {
logger.ERROR("查询任务列表失败", err) logger.ERROR("查询任务列表失败", err)
return nil, common.ErrInternal return nil, common.ErrInternal
} }

2
service/user_service.go

@ -144,7 +144,7 @@ func (s *UserService) ChangePassword(userID int64, req *vo.ChangePasswordReq) *c
func (s *UserService) ListSessions(userID, currentSessionID int64) ([]vo.SessionVO, *common.BusiError) { func (s *UserService) ListSessions(userID, currentSessionID int64) ([]vo.SessionVO, *common.BusiError) {
var sessions []model.UserLoginSession var sessions []model.UserLoginSession
if err := common.DB.Where("user_id = ? AND revoked_at IS NULL", userID).Order("last_active_at DESC").Find(&sessions).Error; err != nil {
if err := common.DB.Where("user_id = ? AND revoked_at IS NULL", userID).Order("last_active_at DESC, id DESC").Find(&sessions).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
result := make([]vo.SessionVO, 0, len(sessions)) result := make([]vo.SessionVO, 0, len(sessions))

2
service/video_service.go

@ -186,7 +186,7 @@ func (s *VideoService) GetPage(userID int64, isAdmin bool, req *vo.VideoPageReq)
return nil, common.ErrInternal return nil, common.ErrInternal
} }
var list []model.Video var list []model.Video
if err := db.Scopes(req.Paginate).Order("id DESC").Find(&list).Error; err != nil {
if err := db.Scopes(req.Paginate).Order("created_at DESC, id DESC").Find(&list).Error; err != nil {
return nil, common.ErrInternal return nil, common.ErrInternal
} }
return common.Page(req.Pagination, total, list), nil return common.Page(req.Pagination, total, list), nil

2
service/workflow_service.go

@ -82,7 +82,7 @@ func (s *WorkflowService) Upsert(dockID, requestID string, in *WorkflowStateIn)
// 回写 task_execution,优先按 command_id 精确关联,兼容旧设备时只回退到最新未终态记录。 // 回写 task_execution,优先按 command_id 精确关联,兼容旧设备时只回退到最新未终态记录。
if (in.CommandID != "" || in.TaskID != "") && (in.State == "running" || isTerminalWorkflowState(in.State)) { if (in.CommandID != "" || in.TaskID != "") && (in.State == "running" || isTerminalWorkflowState(in.State)) {
updates := map[string]any{"status": in.State}
updates := map[string]any{"status": in.State, "result_code": in.ResultCode}
if in.State == "running" { if in.State == "running" {
updates["start_time"] = now updates["start_time"] = now
} else { } else {

1
sql/001_schema.sql

@ -176,6 +176,7 @@ CREATE TABLE IF NOT EXISTS task_execution (
start_time DATETIME, start_time DATETIME,
end_time DATETIME, end_time DATETIME,
status VARCHAR(16) DEFAULT 'pending', status VARCHAR(16) DEFAULT 'pending',
result_code VARCHAR(64),
trajectory_json JSON, trajectory_json JSON,
created_at DATETIME created_at DATETIME
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

2
sql/010_task_workflow_timeout.sql

@ -0,0 +1,2 @@
ALTER TABLE task_execution
ADD COLUMN result_code VARCHAR(64) NULL;

54
tool/snowflake.go

@ -10,7 +10,7 @@ const (
// 起始时间戳 (2023-01-01 00:00:00 UTC) // 起始时间戳 (2023-01-01 00:00:00 UTC)
epoch int64 = 1672531200000 epoch int64 = 1672531200000
timestampBits = 28 // 时间戳位数(约17年)
timestampBits = 41 // 时间戳位数(约69年)
workerIDBits = 5 // 工作机器ID所占位数 workerIDBits = 5 // 工作机器ID所占位数
sequenceBits = 12 // 序列号所占位数 sequenceBits = 12 // 序列号所占位数
@ -37,7 +37,7 @@ func init() {
func NewSnowflake(workerID int64) *Snowflake { func NewSnowflake(workerID int64) *Snowflake {
if workerID < 0 || workerID > maxWorkerID { if workerID < 0 || workerID > maxWorkerID {
panic(errors.New("worker ID must be between 0 and 1023"))
panic(errors.New("worker ID must be between 0 and 31"))
} }
return &Snowflake{ return &Snowflake{
timestamp: 0, timestamp: 0,
@ -52,42 +52,38 @@ func (s *Snowflake) NextID() (int64, error) {
s.mu.Lock() s.mu.Lock()
defer s.mu.Unlock() defer s.mu.Unlock()
now := time.Now().UnixMilli()
now = (now - epoch) & (-1 ^ (-1 << timestampBits))
maxTimestamp := int64(1<<timestampBits) - 1
now := time.Now().UnixMilli() - epoch
if now <= 0 { if now <= 0 {
now = 1
return -1, errors.New("timestamp is before snowflake epoch")
}
if now > maxTimestamp {
return -1, errors.New("timestamp overflow")
} }
if now < s.lastTime { if now < s.lastTime {
waitTime := s.lastTime - now waitTime := s.lastTime - now
if waitTime > 100 { if waitTime > 100 {
now = s.lastTime + 1
if now > (1<<timestampBits)-1 {
return -1, errors.New("timestamp overflow due to clock moved backwards")
return -1, errors.New("clock moved backwards")
} }
} else {
time.Sleep(time.Duration(waitTime) * time.Millisecond) time.Sleep(time.Duration(waitTime) * time.Millisecond)
now = (time.Now().UnixMilli() - epoch) & (-1 ^ (-1 << timestampBits))
if now <= 0 {
now = 1
}
now = time.Now().UnixMilli() - epoch
if now < s.lastTime { if now < s.lastTime {
now = s.lastTime + 1
if now > (1<<timestampBits)-1 {
return -1, errors.New("timestamp overflow due to clock moved backwards")
}
return -1, errors.New("clock moved backwards")
} }
if now > maxTimestamp {
return -1, errors.New("timestamp overflow")
} }
} }
if s.lastTime == now {
if now == s.lastTime {
s.sequence = (s.sequence + 1) & maxSequence s.sequence = (s.sequence + 1) & maxSequence
if s.sequence == 0 { if s.sequence == 0 {
for now <= s.lastTime { for now <= s.lastTime {
now = (time.Now().UnixMilli() - epoch) & (-1 ^ (-1 << timestampBits))
if now <= 0 {
now = s.lastTime + 1
break
time.Sleep(time.Millisecond)
now = time.Now().UnixMilli() - epoch
if now > maxTimestamp {
return -1, errors.New("timestamp overflow")
} }
} }
} }
@ -96,19 +92,7 @@ func (s *Snowflake) NextID() (int64, error) {
} }
s.lastTime = now s.lastTime = now
id := (now << timestampShift) |
(s.workerID << workerIDShift) |
s.sequence
if id == 0 {
s.sequence = 1
id = (now << timestampShift) |
(s.workerID << workerIDShift) |
s.sequence
}
return id, nil
return (now << timestampShift) | (s.workerID << workerIDShift) | s.sequence, nil
} }
// NextID 全局函数,使用默认实例生成ID // NextID 全局函数,使用默认实例生成ID

Loading…
Cancel
Save