299 lines
11 KiB
Go
299 lines
11 KiB
Go
package task
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"time"
|
|
|
|
"go-admin/app/goauto/device"
|
|
"go-admin/app/goauto/models"
|
|
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
type DeleteResponse struct {
|
|
TaskID uint64 `json:"taskId"`
|
|
Deleted bool `json:"deleted"`
|
|
Replayed bool `json:"replayed,omitempty"`
|
|
}
|
|
|
|
type AgentResetResponse struct {
|
|
TaskID uint64 `json:"taskId"`
|
|
AttemptNumber int `json:"attemptNumber"`
|
|
Status string `json:"status"`
|
|
Replayed bool `json:"replayed,omitempty"`
|
|
}
|
|
|
|
func (service *Service) Reset(ctx context.Context, taskID uint64, request ActionRequest) (DetailResponse, error) {
|
|
return service.reset(ctx, taskID, request, nil)
|
|
}
|
|
|
|
func (service *Service) ResetForDevice(ctx context.Context, taskID uint64, request ActionRequest, token string) (AgentResetResponse, error) {
|
|
detail, err := service.reset(ctx, taskID, request, &token)
|
|
if err != nil {
|
|
return AgentResetResponse{}, err
|
|
}
|
|
return AgentResetResponse{TaskID: detail.Task.ID, AttemptNumber: detail.Task.AttemptNumber, Status: detail.Task.Status, Replayed: detail.Replayed}, nil
|
|
}
|
|
|
|
func (service *Service) reset(ctx context.Context, taskID uint64, request ActionRequest, deviceToken *string) (DetailResponse, error) {
|
|
if err := validateAction(taskID, request); err != nil {
|
|
return DetailResponse{}, err
|
|
}
|
|
var replayed bool
|
|
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
var deviceRecord *models.AgentDevice
|
|
if deviceToken != nil {
|
|
authenticated, err := device.NewService(tx).Authenticate(ctx, *deviceToken)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&authenticated, authenticated.ID).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
if authenticated.Status != models.DeviceStatusOnline {
|
|
return serviceError(CodeDeviceOffline, "设备离线,不能重新采集")
|
|
}
|
|
deviceRecord = &authenticated
|
|
}
|
|
var record models.CollectionTask
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&record, taskID).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return serviceError(CodeTaskNotFound, "任务不存在")
|
|
}
|
|
return internalError(err)
|
|
}
|
|
if deviceRecord != nil && (record.DeviceID == nil || *record.DeviceID != deviceRecord.ID) {
|
|
return serviceError(CodeTaskNotFound, "采集任务不存在")
|
|
}
|
|
if deviceRecord == nil && record.DeviceID != nil {
|
|
var assigned models.AgentDevice
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&assigned, *record.DeviceID).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
if assigned.Status != models.DeviceStatusOnline {
|
|
return serviceError(CodeDeviceOffline, "设备离线,不能重新采集")
|
|
}
|
|
deviceRecord = &assigned
|
|
}
|
|
if record.ResetRequestID != nil && *record.ResetRequestID == request.RequestID {
|
|
replayed = true
|
|
return nil
|
|
}
|
|
var archivedReplay models.CollectionTaskAttempt
|
|
if err := tx.Where("archived_by_reset_request_id = ?", request.RequestID).First(&archivedReplay).Error; err == nil {
|
|
if archivedReplay.TaskID != record.ID {
|
|
return serviceError(CodeTaskStateConflict, "requestId 已被其他任务使用")
|
|
}
|
|
replayed = true
|
|
return nil
|
|
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return internalError(err)
|
|
}
|
|
if record.Status == models.TaskStatusRunning {
|
|
return serviceError(CodeTaskStateConflict, "执行中的任务不能重置")
|
|
}
|
|
if record.Source == models.CollectionTaskSourceAgentCurrentPage && record.Status != models.TaskStatusFailed {
|
|
return serviceError(CodeTaskStateConflict, "当前页面采集只允许失败任务重新采集")
|
|
}
|
|
if record.Status != models.TaskStatusCompleted && record.Status != models.TaskStatusCompletedPartial && record.Status != models.TaskStatusFailed {
|
|
return serviceError(CodeTaskStateConflict, "只有终态任务可以重置")
|
|
}
|
|
if record.PDDProductID != nil {
|
|
var active int64
|
|
if err := tx.Model(&models.CollectionTask{}).Where("id <> ? AND pdd_product_id = ? AND status IN ?", record.ID, record.PDDProductID,
|
|
[]string{models.TaskStatusPending, models.TaskStatusRunning}).Count(&active).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
if active > 0 {
|
|
return serviceError(CodeProductTaskActive, "该商品已有待执行或执行中的任务")
|
|
}
|
|
}
|
|
if deviceRecord != nil {
|
|
if err := ensureDeviceIdleForReset(tx, deviceRecord.ID, record.ID, service.Now()); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
rule, err := currentResetRule(tx, record)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if deviceRecord != nil {
|
|
if err := ensureRuleCompatible(*deviceRecord, rule.ContentJSON); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
now := service.Now()
|
|
if err := archiveCollectionAttempt(tx, record, request.RequestID, now); err != nil {
|
|
return err
|
|
}
|
|
if err := deleteResultChildren(tx, record.ID); err != nil {
|
|
return err
|
|
}
|
|
one := uint8(1)
|
|
updates := map[string]any{
|
|
"status": models.TaskStatusPending, "active_slot": one, "device_run_slot": nil,
|
|
"attempt_number": record.AttemptNumber + 1, "rule_id": rule.ID, "rule_snapshot": rule.ContentJSON,
|
|
"lease_expires_at": nil, "claim_request_id": nil, "start_request_id": nil,
|
|
"result_request_id": nil, "fail_request_id": nil, "identify_request_id": nil, "reset_request_id": request.RequestID,
|
|
"title": nil, "shop_name": nil, "sales_text": nil, "review_count": nil, "missing_json": nil,
|
|
"error_code": nil, "error_message": nil, "started_at": nil, "finished_at": nil, "identity_resolved_at": nil,
|
|
}
|
|
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).Where("id = ?", record.ID).Updates(updates).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return DetailResponse{}, err
|
|
}
|
|
detail, err := service.Detail(ctx, taskID)
|
|
detail.Replayed = replayed
|
|
return detail, err
|
|
}
|
|
|
|
func currentResetRule(tx *gorm.DB, record models.CollectionTask) (models.CollectionRule, error) {
|
|
var rule models.CollectionRule
|
|
if record.Source == models.CollectionTaskSourceAgentCurrentPage {
|
|
var setting models.AgentManualCollectionSetting
|
|
if err := tx.First(&setting, 1).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return rule, serviceError(CodeAgentManualRuleNotConfigured, "请先配置 Agent 手动采集规则")
|
|
}
|
|
return rule, internalError(err)
|
|
}
|
|
if err := tx.First(&rule, setting.RuleID).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return rule, serviceError(CodeAgentManualRuleNotConfigured, "请先配置 Agent 手动采集规则")
|
|
}
|
|
return rule, internalError(err)
|
|
}
|
|
if err := ensureCurrentPageRule(rule.ContentJSON); err != nil {
|
|
return rule, err
|
|
}
|
|
return rule, nil
|
|
}
|
|
if err := tx.First(&rule, record.RuleID).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return rule, serviceError(CodeRuleNotFound, "采集规则不存在或已删除")
|
|
}
|
|
return rule, internalError(err)
|
|
}
|
|
return rule, nil
|
|
}
|
|
|
|
func archiveCollectionAttempt(tx *gorm.DB, record models.CollectionTask, resetRequestID string, archivedAt time.Time) error {
|
|
counts := struct {
|
|
Dimensions int64 `json:"dimensions"`
|
|
ColorPrices int64 `json:"colorPrices"`
|
|
SKUs int64 `json:"skus"`
|
|
}{}
|
|
queries := []struct {
|
|
Model any
|
|
Count *int64
|
|
}{
|
|
{Model: &models.CollectionDimension{}, Count: &counts.Dimensions},
|
|
{Model: &models.CollectionColorPrice{}, Count: &counts.ColorPrices},
|
|
{Model: &models.CollectionSKU{}, Count: &counts.SKUs},
|
|
}
|
|
for _, query := range queries {
|
|
if err := tx.Model(query.Model).Where("task_id = ?", record.ID).Count(query.Count).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
}
|
|
summary := struct {
|
|
Title *string `json:"title,omitempty"`
|
|
ShopName *string `json:"shopName,omitempty"`
|
|
MissingJSON *string `json:"missing,omitempty"`
|
|
DimensionCount int64 `json:"dimensionCount"`
|
|
ColorPriceCount int64 `json:"colorPriceCount"`
|
|
SKUCount int64 `json:"skuCount"`
|
|
}{
|
|
Title: record.Title, ShopName: record.ShopName, MissingJSON: record.MissingJSON,
|
|
DimensionCount: counts.Dimensions, ColorPriceCount: counts.ColorPrices, SKUCount: counts.SKUs,
|
|
}
|
|
raw, err := json.Marshal(summary)
|
|
if err != nil {
|
|
return internalError(err)
|
|
}
|
|
attemptNumber := record.AttemptNumber
|
|
if attemptNumber < 1 {
|
|
attemptNumber = 1
|
|
}
|
|
attempt := models.CollectionTaskAttempt{
|
|
TaskID: record.ID, AttemptNumber: attemptNumber, Source: record.Source, DeviceID: record.DeviceID,
|
|
RuleID: record.RuleID, RuleSnapshot: record.RuleSnapshot, Status: record.Status,
|
|
ErrorCode: record.ErrorCode, ErrorMessage: record.ErrorMessage, ResultSummaryJSON: string(raw),
|
|
StartedAt: record.StartedAt, FinishedAt: record.FinishedAt,
|
|
ArchivedByResetRequestID: resetRequestID, ArchivedAt: archivedAt,
|
|
}
|
|
if err := tx.Create(&attempt).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func ensureDeviceIdleForReset(tx *gorm.DB, deviceID, taskID uint64, now time.Time) error {
|
|
var busy int64
|
|
if err := tx.Model(&models.CollectionTask{}).
|
|
Where("id <> ? AND device_id = ? AND (status = ? OR (status = ? AND lease_expires_at > ?))", taskID, deviceID, models.TaskStatusRunning, models.TaskStatusPending, now).
|
|
Count(&busy).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
if busy > 0 {
|
|
return serviceError(CodeDeviceBusy, "设备正在执行其他任务")
|
|
}
|
|
if err := tx.Model(&models.PurchaseTask{}).
|
|
Where("device_id = ? AND (status IN ? OR (status = ? AND lease_expires_at > ?))", deviceID,
|
|
[]string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted, models.PurchaseTaskStatusSpecProbePending},
|
|
models.PurchaseTaskStatusPending, now).
|
|
Count(&busy).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
if busy > 0 {
|
|
return serviceError(CodeDeviceBusy, "设备正在执行其他任务")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (service *Service) DeleteFailed(ctx context.Context, taskID uint64, request ActionRequest) (DeleteResponse, error) {
|
|
if err := validateAction(taskID, request); err != nil {
|
|
return DeleteResponse{}, err
|
|
}
|
|
response := DeleteResponse{TaskID: taskID}
|
|
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
var record models.CollectionTask
|
|
if err := tx.Unscoped().Clauses(clause.Locking{Strength: "UPDATE"}).First(&record, taskID).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return serviceError(CodeTaskNotFound, "任务不存在")
|
|
}
|
|
return internalError(err)
|
|
}
|
|
if record.DeletedAt.Valid {
|
|
if record.DeleteRequestID != nil && *record.DeleteRequestID == request.RequestID {
|
|
response.Deleted = true
|
|
response.Replayed = true
|
|
return nil
|
|
}
|
|
return serviceError(CodeTaskNotFound, "任务不存在")
|
|
}
|
|
if record.Status == models.TaskStatusRunning {
|
|
return serviceError(CodeTaskStateConflict, "执行中的任务不能删除")
|
|
}
|
|
if record.Status != models.TaskStatusFailed {
|
|
return serviceError(CodeTaskStateConflict, "只允许删除失败任务")
|
|
}
|
|
if err := tx.Model(&record).Update("delete_request_id", request.RequestID).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
if err := tx.Delete(&record).Error; err != nil {
|
|
return internalError(err)
|
|
}
|
|
response.Deleted = true
|
|
return nil
|
|
})
|
|
return response, err
|
|
}
|