package service import ( "context" "crypto/rand" "database/sql" "encoding/hex" "encoding/json" "errors" "fmt" "strings" "time" "cmautobuy/admin/model" "cmautobuy/admin/repository" "cmautobuy/admin/syb" ) // InnerCodeWriter 是安全回写所需的最小顺运宝接口。 type InnerCodeWriter interface { InnerCodeDetailReader DeleteInnerCode(context.Context, int64) error CreateInnerCodeDetail(context.Context, int64, string) (int64, error) UpdateDetailCode(context.Context, int64, int64, string) error } // InnerCodeApplyResult 是一次用户确认操作的统计。 type InnerCodeApplyResult struct { Requested int Updated int AlreadyFilled int Skipped int Failed int NeedsCheck int } // InnerCodeApplyBatch 是一次已经落库的后台回写请求。 type InnerCodeApplyBatch struct { ID string Count int } const innerCodeApplyChunkSize = 20 type innerCodeApplyOutcome struct { Status model.InnerCodeStatus Message string RemoteCode string } type innerCodeRemoteItem struct { Code string `json:"code"` DetailID int64 `json:"detail_id,omitempty"` Source string `json:"source,omitempty"` Title string `json:"title,omitempty"` Status string `json:"status"` } type innerCodeCheckpoint func([]innerCodeRemoteItem, string) error // QueueInnerCodeApplyBatch 把任意数量的已选记录原子排队,远端请求由后台执行器处理。 func QueueInnerCodeApplyBatch(db *sql.DB, ids []int64, actorUserID string) (*InnerCodeApplyBatch, error) { ids = uniquePositiveInnerCodeIDs(ids) if len(ids) == 0 { return nil, fmt.Errorf("没有选择可回写记录") } batchID, err := newInnerCodeApplyBatchID() if err != nil { return nil, err } count, err := repository.QueueInnerCodeApplyBatch(db, ids, batchID, actorUserID, model.NowISO()) if err != nil { if errors.Is(err, repository.ErrInnerCodeApplyConflict) { return nil, fmt.Errorf("所选记录中有记录已不再可回写;本批没有部分入队,请刷新后重新选择") } return nil, err } return &InnerCodeApplyBatch{ID: batchID, Count: count}, nil } // RunInnerCodeApplyBatch 分组读取后台队列,再逐条执行原有安全回写门禁。 // 分组大小只用于控制一次数据库读取,远端请求始终逐条发送且不会自动重试。 func RunInnerCodeApplyBatch(ctx context.Context, db *sql.DB, writer InnerCodeWriter, batchID, actorUserID string) (*InnerCodeApplyResult, error) { result := &InnerCodeApplyResult{} for { if err := ctx.Err(); err != nil { return result, err } ids, err := repository.ListQueuedInnerCodeIDs(db, batchID, innerCodeApplyChunkSize) if err != nil { return result, err } if len(ids) == 0 { return result, nil } result.Requested += len(ids) for _, id := range ids { if err := ctx.Err(); err != nil { return result, err } now := model.NowISO() record, claimed, err := repository.ClaimQueuedInnerCodeForApply(db, id, batchID, actorUserID, now) if err != nil { return result, err } if !claimed { result.Skipped++ continue } checkpoint := func(items []innerCodeRemoteItem, message string) error { raw, err := json.Marshal(items) if err != nil { return fmt.Errorf("序列化档口入库码逐件检查点失败: %w", err) } return repository.SaveInnerCodeRemoteItems(db, id, string(raw), compactInnerCodeMessage(message), model.NowISO()) } outcome := applyClaimedInnerCode(ctx, writer, *record, checkpoint) finishedAt := model.NowISO() if err := repository.FinishInnerCodeApply(db, id, outcome.Status, compactInnerCodeMessage(outcome.Message), outcome.RemoteCode, finishedAt); err != nil { return result, err } switch outcome.Status { case model.InnerCodeUpdated: result.Updated++ case model.InnerCodeAlreadyFilled: result.AlreadyFilled++ case model.InnerCodeNeedsCheck: result.NeedsCheck++ case model.InnerCodeFailed: result.Failed++ default: result.Skipped++ } } } } // InterruptInnerCodeApplyBatch 在后台异常时收敛当前批次,不自动重试任何远端请求。 func InterruptInnerCodeApplyBatch(db *sql.DB, batchID string) (int, int, error) { return repository.InterruptInnerCodeApplyBatch(db, strings.TrimSpace(batchID), model.NowISO()) } // GetInnerCodeApplyBatchProgress 返回页面使用的单表聚合进度。 func GetInnerCodeApplyBatchProgress(db *sql.DB, batchID string) (*repository.InnerCodeApplyBatchProgress, error) { batchID = strings.TrimSpace(batchID) if batchID == "" || len(batchID) > 191 { return nil, nil } return repository.GetInnerCodeApplyBatchProgress(db, batchID) } func newInnerCodeApplyBatchID() (string, error) { random := make([]byte, 8) if _, err := rand.Read(random); err != nil { return "", fmt.Errorf("生成档口入库码后台批次失败: %w", err) } return "ICB-" + time.Now().UTC().Format("20060102T150405") + "-" + hex.EncodeToString(random), nil } func applyClaimedInnerCode(ctx context.Context, writer InnerCodeWriter, record model.InnerCodeRecord, checkpoint innerCodeCheckpoint) innerCodeApplyOutcome { codes, err := splitInnerCodes(record.InnerCode) if err != nil || len(codes) != record.SourceDuplicateCount { return innerCodeApplyOutcome{Status: model.InnerCodeFailed, Message: fmt.Sprintf("导入单件码数量与源行数不一致(入库码 %d 个,源行 %d 行),没有发送远端写请求", len(codes), record.SourceDuplicateCount), RemoteCode: record.RemoteInnerCode} } stock, original, err := readCurrentInnerCodeStock(ctx, writer, record) if err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: "写入前重新读取失败,没有发送远端写请求:" + err.Error(), RemoteCode: record.RemoteInnerCode} } if original.ProductQty != len(codes) { return innerCodeApplyOutcome{Status: model.InnerCodeFailed, Message: fmt.Sprintf("顺运宝商品数量为 %d,单件入库码为 %d 个,数量不一致;没有发送远端写请求", original.ProductQty, len(codes)), RemoteCode: innerCodeRawText(original.Raw["innerExpCode"])} } if !innerCodeDetailIdentityMatches(record, *original) { return innerCodeApplyOutcome{Status: model.InnerCodeSkipped, Message: "顺运宝商品规格或档口身份已变化,停止回写,请重新匹配", RemoteCode: innerCodeRawText(original.Raw["innerExpCode"])} } if innerCodeRawText(original.Raw["purchasePlatform"]) != "" || innerCodeRawText(original.Raw["purchaseCode"]) != "" { return innerCodeApplyOutcome{Status: model.InnerCodeSkipped, Message: "顺运宝商品已有采购平台或采购单号,停止回写", RemoteCode: innerCodeRawText(original.Raw["innerExpCode"])} } items, missing, prepareErr := planInnerCodeRemoteItems(record, stock, codes) if prepareErr != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: prepareErr.Error(), RemoteCode: record.RemoteInnerCode} } if checkpoint == nil { checkpoint = func([]innerCodeRemoteItem, string) error { return nil } } if len(missing) == 0 { if err := checkpoint(items, fmt.Sprintf("远端已存在全部 %d 个单件入库码", len(codes))); err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: err.Error(), RemoteCode: strings.Join(codes, ",")} } return innerCodeApplyOutcome{Status: model.InnerCodeAlreadyFilled, Message: fmt.Sprintf("远端已存在全部 %d 个单件入库码,无需重复写入", len(codes)), RemoteCode: strings.Join(codes, ",")} } currentCode := innerCodeRawText(original.Raw["innerExpCode"]) if currentCode != "" && !containsInnerCode(codes, currentCode) { if currentCode != record.RemoteInnerCode { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: "原商品快递单号在规划后发生变化,已停止回写,请人工核对", RemoteCode: currentCode} } if err := checkpoint(items, "准备清除原商品上的旧快递单号"); err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: err.Error(), RemoteCode: currentCode} } if err := writer.DeleteInnerCode(ctx, record.DetailID); err != nil { return withInnerCodeProgress(unknownOrFailedInnerCodeOutcome(err, "删除旧快递单号", currentCode), items) } currentCode = "" } for _, codeIndex := range missing { item := &items[codeIndex] if item.DetailID == 0 && currentCode == "" { item.DetailID = record.DetailID item.Source = "original" currentCode = "reserved" } if item.DetailID == 0 { item.Title = innerCodePlaceholderTitle(record.ID, codeIndex+1) item.Status = "creating" if err := checkpoint(items, fmt.Sprintf("准备为第 %d 件创建零价明细", codeIndex+1)); err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: err.Error(), RemoteCode: record.RemoteInnerCode} } detailID, err := writer.CreateInnerCodeDetail(ctx, record.StockID, item.Title) if err != nil { return withInnerCodeProgress(unknownOrFailedInnerCodeOutcome(err, fmt.Sprintf("创建第 %d 件明细", codeIndex+1), record.RemoteInnerCode), items) } item.DetailID = detailID item.Source = "created" item.Status = "created" if err := checkpoint(items, fmt.Sprintf("第 %d 件明细已创建,准备写入入库码", codeIndex+1)); err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: "新增明细后保存检查点失败,禁止继续写入:" + err.Error(), RemoteCode: record.RemoteInnerCode} } } item.Status = "writing" if err := checkpoint(items, fmt.Sprintf("准备写入第 %d 个单件入库码", codeIndex+1)); err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: err.Error(), RemoteCode: record.RemoteInnerCode} } if err := writer.UpdateDetailCode(ctx, record.StockID, item.DetailID, item.Code); err != nil { return withInnerCodeProgress(unknownOrFailedInnerCodeOutcome(err, fmt.Sprintf("写入第 %d 个单件入库码", codeIndex+1), record.RemoteInnerCode), items) } verifiedStock, _, err := readCurrentInnerCodeStock(ctx, writer, record) if err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: "写入成功响应后重新读取失败,禁止自动重试", RemoteCode: record.RemoteInnerCode} } if countInnerCodeInStock(verifiedStock, item.Code) != 1 || innerCodeDetailCode(verifiedStock, item.DetailID) != item.Code { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: "写入后单件入库码未唯一出现在预期明细,禁止自动重试", RemoteCode: record.RemoteInnerCode} } item.Status = "confirmed" if err := checkpoint(items, fmt.Sprintf("第 %d 个单件入库码已回读确认", codeIndex+1)); err != nil { return innerCodeApplyOutcome{Status: model.InnerCodeNeedsCheck, Message: "写入确认后保存检查点失败,请人工核对:" + err.Error(), RemoteCode: record.RemoteInnerCode} } } return innerCodeApplyOutcome{Status: model.InnerCodeUpdated, Message: fmt.Sprintf("回写完成:%d 个单件入库码均已逐件回读确认", len(codes)), RemoteCode: strings.Join(codes, ",")} } func readCurrentInnerCodeDetail(ctx context.Context, reader InnerCodeDetailReader, record model.InnerCodeRecord) (*syb.DetailItem, error) { _, found, err := readCurrentInnerCodeStock(ctx, reader, record) return found, err } func readCurrentInnerCodeStock(ctx context.Context, reader InnerCodeDetailReader, record model.InnerCodeRecord) (syb.StockDetail, *syb.DetailItem, error) { if record.StockID <= 0 || record.DetailID <= 0 { return syb.StockDetail{}, nil, fmt.Errorf("记录缺少有效的货运单或商品明细 ID") } stocks, err := reader.DetailListByStock(ctx, []int64{record.StockID}) if err != nil { return syb.StockDetail{}, nil, err } if len(stocks) != 1 || stocks[0].ID != record.StockID { return syb.StockDetail{}, nil, fmt.Errorf("顺运宝没有唯一返回货运单 id=%d", record.StockID) } var found *syb.DetailItem for index := range stocks[0].Details { if stocks[0].Details[index].ID != record.DetailID { continue } if found != nil { return syb.StockDetail{}, nil, fmt.Errorf("顺运宝重复返回商品明细 id=%d", record.DetailID) } item := stocks[0].Details[index] found = &item } if found == nil { return syb.StockDetail{}, nil, fmt.Errorf("顺运宝未返回商品明细 id=%d", record.DetailID) } return stocks[0], found, nil } func splitInnerCodes(raw string) ([]string, error) { parts := strings.Split(raw, ",") codes := make([]string, 0, len(parts)) seen := make(map[string]bool, len(parts)) for _, part := range parts { code := strings.TrimSpace(part) if code == "" { return nil, fmt.Errorf("单件入库码不能为空") } if seen[code] { return nil, fmt.Errorf("单件入库码 %q 重复", code) } seen[code] = true codes = append(codes, code) } return codes, nil } func planInnerCodeRemoteItems(record model.InnerCodeRecord, stock syb.StockDetail, codes []string) ([]innerCodeRemoteItem, []int, error) { items := make([]innerCodeRemoteItem, len(codes)) indexByCode := make(map[string]int, len(codes)) for index, code := range codes { items[index] = innerCodeRemoteItem{Code: code, Status: "planned"} indexByCode[code] = index } for _, detail := range stock.Details { code := innerCodeRawText(detail.Raw["innerExpCode"]) if index, ok := indexByCode[code]; ok { if items[index].DetailID != 0 { return nil, nil, fmt.Errorf("单件入库码 %s 在货运单中出现多次,停止自动回写", code) } source := "existing" if detail.ID == record.DetailID { source = "original" } items[index].DetailID = detail.ID items[index].Source = source items[index].Status = "confirmed" continue } for index := range items { title := innerCodePlaceholderTitle(record.ID, index+1) if detail.ProductTitle != title { continue } if items[index].DetailID != 0 { return nil, nil, fmt.Errorf("第 %d 件占位明细在货运单中出现多次,停止自动回写", index+1) } if code != "" { return nil, nil, fmt.Errorf("第 %d 件占位明细已有非目标快递单号,停止自动回写", index+1) } items[index].DetailID = detail.ID items[index].Source = "existing_placeholder" items[index].Title = title items[index].Status = "created" } } missing := make([]int, 0, len(items)) for index := range items { if items[index].Status != "confirmed" { missing = append(missing, index) } } return items, missing, nil } func innerCodePlaceholderTitle(recordID int64, ordinal int) string { return fmt.Sprintf("档口第%d件-IC%d", ordinal, recordID) } func containsInnerCode(codes []string, value string) bool { for _, code := range codes { if code == value { return true } } return false } func countInnerCodeInStock(stock syb.StockDetail, code string) int { count := 0 for _, detail := range stock.Details { if innerCodeRawText(detail.Raw["innerExpCode"]) == code { count++ } } return count } func innerCodeDetailCode(stock syb.StockDetail, detailID int64) string { for _, detail := range stock.Details { if detail.ID == detailID { return innerCodeRawText(detail.Raw["innerExpCode"]) } } return "" } func unknownOrFailedInnerCodeOutcome(err error, action, remoteCode string) innerCodeApplyOutcome { status := model.InnerCodeFailed message := action + "失败,系统不会自动重试:" + err.Error() if errors.Is(err, syb.ErrWriteResultUnknown) { status = model.InnerCodeNeedsCheck message = action + "结果未知,禁止自动重试,请重新核对" } return innerCodeApplyOutcome{Status: status, Message: message, RemoteCode: remoteCode} } func withInnerCodeProgress(outcome innerCodeApplyOutcome, items []innerCodeRemoteItem) innerCodeApplyOutcome { confirmed := 0 for _, item := range items { if item.Status == "confirmed" { confirmed++ } } outcome.Message = fmt.Sprintf("已确认 %d/%d 件;%s", confirmed, len(items), outcome.Message) return outcome } // RecheckInnerCode 只重新读取一条 needs_check 记录,不发送任何写请求。 func RecheckInnerCode(ctx context.Context, db *sql.DB, reader InnerCodeDetailReader, id int64) (model.InnerCodeStatus, string, error) { record, err := repository.GetInnerCodeForRecheck(db, id) if err != nil { return "", "", err } if record == nil { return "", "", fmt.Errorf("档口入库码记录不存在") } if record.Status != model.InnerCodeNeedsCheck { return record.Status, "当前记录不需要核对", nil } stock, item, readErr := readCurrentInnerCodeStock(ctx, reader, *record) status := model.InnerCodeNeedsCheck remoteCode := record.RemoteInnerCode remoteItemsJSON := record.RemoteItemsJSON message := "重新读取失败,仍需人工核对:" + errorText(readErr) if readErr == nil { codes, splitErr := splitInnerCodes(record.InnerCode) items, missing, planErr := planInnerCodeRemoteItems(*record, stock, codes) if splitErr == nil && planErr == nil { if raw, err := json.Marshal(items); err == nil { remoteItemsJSON = string(raw) } } remoteCode = innerCodeRawText(item.Raw["innerExpCode"]) if !innerCodeDetailIdentityMatches(*record, *item) { message = "重新读取到的商品身份与规划不一致;保持需核对,系统没有写入" } else if splitErr != nil || planErr != nil { message = "重新读取到的逐件入库码存在冲突;保持需核对,系统没有写入" } else if len(missing) == 0 { status = model.InnerCodeUpdated remoteCode = strings.Join(codes, ",") message = fmt.Sprintf("重新读取确认远端已存在全部 %d 个单件入库码;没有重复写入", len(codes)) } else { message = fmt.Sprintf("重新读取后仍缺少 %d 个单件入库码;保持需核对,系统没有写入", len(missing)) } } checkedAt := model.NowISO() message = compactInnerCodeMessage(message) if err := repository.SaveInnerCodeRecheck(db, id, status, message, remoteCode, remoteItemsJSON, checkedAt); err != nil { return "", "", err } return status, message, nil } func innerCodeDetailIdentityMatches(record model.InnerCodeRecord, item syb.DetailItem) bool { return item.ProductSpec == record.SybSpec && innerCodeRawText(item.Raw["sku"]) == record.SybSKU && innerCodeRawText(item.Raw["variationSku"]) == record.SybVariationSKU } func errorText(err error) string { if err == nil { return "" } return err.Error() } func uniquePositiveInnerCodeIDs(ids []int64) []int64 { seen := make(map[int64]bool, len(ids)) result := make([]int64, 0, len(ids)) for _, id := range ids { if id > 0 && !seen[id] { seen[id] = true result = append(result, id) } } return result } // InterruptApplyingInnerCodes 在 Admin 启动时收敛未确认的远端写结果。 func InterruptApplyingInnerCodes(db *sql.DB, now time.Time) (int, error) { return repository.InterruptApplyingInnerCodes(db, now.UTC().Format(model.TimeLayout)) } func compactInnerCodeMessage(message string) string { message = strings.TrimSpace(message) if len([]rune(message)) <= 500 { return message } return string([]rune(message)[:500]) }