2026-08-15 08:47:22 +08:00
|
|
|
|
package repository
|
|
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
|
"database/sql"
|
|
|
|
|
|
"errors"
|
|
|
|
|
|
"fmt"
|
2026-08-15 09:41:45 +08:00
|
|
|
|
"sort"
|
2026-08-15 08:51:38 +08:00
|
|
|
|
"strings"
|
2026-08-15 08:47:22 +08:00
|
|
|
|
|
2026-08-15 09:12:23 +08:00
|
|
|
|
"github.com/go-sql-driver/mysql"
|
|
|
|
|
|
|
2026-08-15 08:47:22 +08:00
|
|
|
|
"cmautobuy/admin/model"
|
|
|
|
|
|
)
|
|
|
|
|
|
|
2026-08-15 09:12:23 +08:00
|
|
|
|
// ErrInnerCodeUniqueConflict 表示同一日期的入库码已属于另一个业务键。
|
|
|
|
|
|
var ErrInnerCodeUniqueConflict = errors.New("同一业务日期的档口入库码已被其他记录使用")
|
|
|
|
|
|
|
2026-08-15 10:49:28 +08:00
|
|
|
|
// ErrInnerCodeRestoreConflict 表示已完成或结果未知的软删除记录不能换入库码后直接恢复。
|
|
|
|
|
|
var ErrInnerCodeRestoreConflict = errors.New("已回写或需核对的删除记录不能用不同入库码恢复")
|
|
|
|
|
|
|
|
|
|
|
|
// ErrInnerCodeDeleteConflict 表示批量删除时记录已经不可见或不存在,整批不会部分删除。
|
|
|
|
|
|
var ErrInnerCodeDeleteConflict = errors.New("部分档口入库码记录已删除或不存在")
|
|
|
|
|
|
|
2026-08-15 17:08:17 +08:00
|
|
|
|
// ErrInnerCodeApplyConflict 表示所选记录有一条已不再可回写,整批不会部分入队。
|
|
|
|
|
|
var ErrInnerCodeApplyConflict = errors.New("部分档口入库码记录已不再可回写")
|
|
|
|
|
|
|
2026-08-15 08:59:40 +08:00
|
|
|
|
// InnerCodeListFilter 是独立页面可组合的查询条件。
|
|
|
|
|
|
type InnerCodeListFilter struct {
|
|
|
|
|
|
BusinessDate string
|
|
|
|
|
|
Status string
|
|
|
|
|
|
Keyword string
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// InnerCodeStatusCounts 是当前业务日期的底栏摘要。
|
|
|
|
|
|
type InnerCodeStatusCounts struct {
|
|
|
|
|
|
Total int
|
|
|
|
|
|
Ready int
|
2026-08-15 17:08:17 +08:00
|
|
|
|
Queued int
|
|
|
|
|
|
Applying int
|
2026-08-15 08:59:40 +08:00
|
|
|
|
NeedsCheck int
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 17:08:17 +08:00
|
|
|
|
// InnerCodeApplyBatchProgress 是从单表聚合得到的后台回写进度。
|
|
|
|
|
|
type InnerCodeApplyBatchProgress struct {
|
|
|
|
|
|
BatchID string
|
|
|
|
|
|
Total int
|
|
|
|
|
|
Queued int
|
|
|
|
|
|
Applying int
|
|
|
|
|
|
Updated int
|
|
|
|
|
|
AlreadyFilled int
|
|
|
|
|
|
Skipped int
|
|
|
|
|
|
Failed int
|
|
|
|
|
|
NeedsCheck int
|
|
|
|
|
|
Ready int
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Processed 返回已经结束远端处理的数量;恢复为 ready 的未开始记录不算已处理。
|
|
|
|
|
|
func (p InnerCodeApplyBatchProgress) Processed() int {
|
|
|
|
|
|
return p.Updated + p.AlreadyFilled + p.Skipped + p.Failed + p.NeedsCheck
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 10:49:28 +08:00
|
|
|
|
// InnerCodeImportOutcome 说明幂等导入是新增、更新还是恢复软删除记录。
|
2026-08-15 08:47:22 +08:00
|
|
|
|
type InnerCodeImportOutcome string
|
|
|
|
|
|
|
|
|
|
|
|
const (
|
2026-08-15 10:49:28 +08:00
|
|
|
|
InnerCodeImportCreated InnerCodeImportOutcome = "created"
|
|
|
|
|
|
InnerCodeImportUpdated InnerCodeImportOutcome = "updated"
|
|
|
|
|
|
InnerCodeImportRestored InnerCodeImportOutcome = "restored"
|
2026-08-15 08:47:22 +08:00
|
|
|
|
)
|
|
|
|
|
|
|
2026-08-15 14:27:31 +08:00
|
|
|
|
// InnerCodeCodeOwner 是导入事务中核对单件码归属所需的最小现有记录。
|
|
|
|
|
|
type InnerCodeCodeOwner struct {
|
|
|
|
|
|
ID int64
|
|
|
|
|
|
OrderNumber string
|
|
|
|
|
|
Stall string
|
|
|
|
|
|
SpecKey string
|
|
|
|
|
|
InnerCode string
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// LockInnerCodeCodeOwnersByDate 锁定同日现有记录,供 service 拆分聚合值并防止单件码串单。
|
|
|
|
|
|
// 查询包含软删除记录,因为软删除记录仍保留业务唯一键并可能被重新导入恢复。
|
|
|
|
|
|
func LockInnerCodeCodeOwnersByDate(tx *sql.Tx, businessDate string) ([]InnerCodeCodeOwner, error) {
|
|
|
|
|
|
rows, err := tx.Query(`SELECT id,order_number,stall,spec_key,inner_code
|
|
|
|
|
|
FROM syb_inner_code_records WHERE business_date=? ORDER BY id FOR UPDATE`, businessDate)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("锁定同日档口入库码记录失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
|
owners := make([]InnerCodeCodeOwner, 0)
|
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
|
var owner InnerCodeCodeOwner
|
|
|
|
|
|
if err := rows.Scan(&owner.ID, &owner.OrderNumber, &owner.Stall, &owner.SpecKey, &owner.InnerCode); err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("读取同日档口入库码记录失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
owners = append(owners, owner)
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("遍历同日档口入库码记录失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return owners, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 08:47:22 +08:00
|
|
|
|
// UpsertInnerCodeImportRow 按已确认业务键写入一行。
|
|
|
|
|
|
// 必须在事务中调用;先锁定业务键,避免另一个唯一键冲突时更新错行。
|
|
|
|
|
|
func UpsertInnerCodeImportRow(tx *sql.Tx, row model.InnerCodeImportRow, now string) (InnerCodeImportOutcome, error) {
|
2026-08-15 09:12:23 +08:00
|
|
|
|
var codeRecordID int64
|
|
|
|
|
|
var codeOrder, codeStall, codeSpecKey string
|
|
|
|
|
|
codeErr := tx.QueryRow(`SELECT id,order_number,stall,spec_key FROM syb_inner_code_records
|
|
|
|
|
|
WHERE business_date=? AND inner_code=? FOR UPDATE`, row.BusinessDate, row.InnerCode).
|
|
|
|
|
|
Scan(&codeRecordID, &codeOrder, &codeStall, &codeSpecKey)
|
|
|
|
|
|
if codeErr != nil && !errors.Is(codeErr, sql.ErrNoRows) {
|
|
|
|
|
|
return "", fmt.Errorf("核对档口入库码唯一性失败: %w", codeErr)
|
|
|
|
|
|
}
|
|
|
|
|
|
if codeErr == nil && (codeOrder != row.OrderNumber || codeStall != row.Stall || codeSpecKey != row.SpecKey) {
|
|
|
|
|
|
return "", fmt.Errorf("%w(记录 %d)", ErrInnerCodeUniqueConflict, codeRecordID)
|
|
|
|
|
|
}
|
2026-08-15 08:47:22 +08:00
|
|
|
|
var id int64
|
2026-08-15 10:49:28 +08:00
|
|
|
|
var status model.InnerCodeStatus
|
|
|
|
|
|
var currentInnerCode, deletedAt string
|
2026-08-15 08:47:22 +08:00
|
|
|
|
err := tx.QueryRow(`
|
2026-08-15 10:49:28 +08:00
|
|
|
|
SELECT id,status,inner_code,COALESCE(deleted_at,'') FROM syb_inner_code_records
|
2026-08-15 08:47:22 +08:00
|
|
|
|
WHERE business_date=? AND order_number=? AND stall=? AND spec_key=?
|
2026-08-15 10:49:28 +08:00
|
|
|
|
FOR UPDATE`, row.BusinessDate, row.OrderNumber, row.Stall, row.SpecKey).
|
|
|
|
|
|
Scan(&id, &status, ¤tInnerCode, &deletedAt)
|
2026-08-15 08:47:22 +08:00
|
|
|
|
if err != nil && !errors.Is(err, sql.ErrNoRows) {
|
|
|
|
|
|
return "", fmt.Errorf("锁定档口入库码业务键失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
|
|
|
|
_, err = tx.Exec(`
|
|
|
|
|
|
INSERT INTO syb_inner_code_records (
|
2026-08-18 15:25:58 +08:00
|
|
|
|
business_date,source_row,print_sequence,order_number,shop_name,stall,source_sku_raw,
|
2026-08-15 08:47:22 +08:00
|
|
|
|
spec_raw,spec_key,inner_code,source_duplicate_count,status,
|
|
|
|
|
|
created_by_user_id,created_at,updated_at
|
2026-08-18 15:25:58 +08:00
|
|
|
|
) VALUES (?,?,?,?,?,?,?,?,?,?,?,'pending',?,?,?)`,
|
2026-08-15 08:47:22 +08:00
|
|
|
|
row.BusinessDate, row.SourceRow, nullablePositiveInt(row.PrintSequence), row.OrderNumber,
|
2026-08-18 15:25:58 +08:00
|
|
|
|
nullableString(row.ShopName), row.Stall, nullableStringPreserveSpace(row.SourceSKURaw), row.SpecRaw, row.SpecKey, row.InnerCode,
|
2026-08-15 08:47:22 +08:00
|
|
|
|
row.SourceDuplicateCount, row.CreatedByUserID, now, now)
|
|
|
|
|
|
if err != nil {
|
2026-08-15 09:12:23 +08:00
|
|
|
|
return "", innerCodeImportWriteError("新增档口入库码记录失败", err)
|
2026-08-15 08:47:22 +08:00
|
|
|
|
}
|
|
|
|
|
|
return InnerCodeImportCreated, nil
|
|
|
|
|
|
}
|
2026-08-15 10:49:28 +08:00
|
|
|
|
if deletedAt != "" {
|
|
|
|
|
|
if innerCodeStatusPreservesImportResult(status) {
|
|
|
|
|
|
if currentInnerCode != row.InnerCode {
|
|
|
|
|
|
return "", fmt.Errorf("%w(记录 %d,原入库码 %q,新入库码 %q)",
|
|
|
|
|
|
ErrInnerCodeRestoreConflict, id, currentInnerCode, row.InnerCode)
|
|
|
|
|
|
}
|
|
|
|
|
|
_, err = tx.Exec(`UPDATE syb_inner_code_records
|
2026-08-18 15:25:58 +08:00
|
|
|
|
SET source_row=?,print_sequence=?,shop_name=?,source_sku_raw=?,spec_raw=?,source_duplicate_count=?,
|
2026-08-15 10:49:28 +08:00
|
|
|
|
deleted_at=NULL,deleted_by_user_id=NULL,updated_at=?
|
|
|
|
|
|
WHERE id=?`, row.SourceRow, nullablePositiveInt(row.PrintSequence), nullableString(row.ShopName),
|
2026-08-18 15:25:58 +08:00
|
|
|
|
nullableStringPreserveSpace(row.SourceSKURaw), row.SpecRaw, row.SourceDuplicateCount, now, id)
|
2026-08-15 10:49:28 +08:00
|
|
|
|
} else {
|
|
|
|
|
|
_, err = tx.Exec(`UPDATE syb_inner_code_records
|
2026-08-18 15:25:58 +08:00
|
|
|
|
SET source_row=?,print_sequence=?,shop_name=?,source_sku_raw=?,spec_raw=?,inner_code=?,source_duplicate_count=?,
|
2026-08-15 10:49:28 +08:00
|
|
|
|
status='pending',stock_id=NULL,detail_id=NULL,syb_spec=NULL,syb_sku=NULL,
|
|
|
|
|
|
syb_variation_sku=NULL,purchase_platform=NULL,purchase_code=NULL,
|
2026-08-18 15:18:35 +08:00
|
|
|
|
remote_inner_code=NULL,remote_items_json=NULL,result_message=NULL,planned_at=NULL,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
apply_batch_id=NULL,apply_queued_at=NULL,apply_started_at=NULL,
|
2026-08-15 10:49:28 +08:00
|
|
|
|
deleted_at=NULL,deleted_by_user_id=NULL,updated_at=?
|
|
|
|
|
|
WHERE id=?`, row.SourceRow, nullablePositiveInt(row.PrintSequence), nullableString(row.ShopName),
|
2026-08-18 15:25:58 +08:00
|
|
|
|
nullableStringPreserveSpace(row.SourceSKURaw), row.SpecRaw, row.InnerCode, row.SourceDuplicateCount, now, id)
|
2026-08-15 10:49:28 +08:00
|
|
|
|
}
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return "", innerCodeImportWriteError("恢复档口入库码导入记录失败", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return InnerCodeImportRestored, nil
|
|
|
|
|
|
}
|
2026-08-15 08:47:22 +08:00
|
|
|
|
|
|
|
|
|
|
_, err = tx.Exec(`
|
|
|
|
|
|
UPDATE syb_inner_code_records
|
2026-08-18 15:25:58 +08:00
|
|
|
|
SET source_row=?,print_sequence=?,shop_name=?,source_sku_raw=?,spec_raw=?,inner_code=?,
|
2026-08-15 08:47:22 +08:00
|
|
|
|
source_duplicate_count=?,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
status=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check')
|
2026-08-15 08:47:22 +08:00
|
|
|
|
THEN status ELSE 'pending' END,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
stock_id=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN stock_id ELSE NULL END,
|
|
|
|
|
|
detail_id=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN detail_id ELSE NULL END,
|
|
|
|
|
|
syb_spec=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN syb_spec ELSE NULL END,
|
|
|
|
|
|
syb_sku=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN syb_sku ELSE NULL END,
|
|
|
|
|
|
syb_variation_sku=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN syb_variation_sku ELSE NULL END,
|
|
|
|
|
|
purchase_platform=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN purchase_platform ELSE NULL END,
|
|
|
|
|
|
purchase_code=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN purchase_code ELSE NULL END,
|
|
|
|
|
|
remote_inner_code=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN remote_inner_code ELSE NULL END,
|
2026-08-18 15:18:35 +08:00
|
|
|
|
remote_items_json=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN remote_items_json ELSE NULL END,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
result_message=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN result_message ELSE NULL END,
|
|
|
|
|
|
planned_at=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN planned_at ELSE NULL END,
|
|
|
|
|
|
apply_batch_id=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN apply_batch_id ELSE NULL END,
|
|
|
|
|
|
apply_queued_at=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN apply_queued_at ELSE NULL END,
|
|
|
|
|
|
apply_started_at=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check') THEN apply_started_at ELSE NULL END,
|
2026-08-15 08:47:22 +08:00
|
|
|
|
updated_at=?
|
|
|
|
|
|
WHERE id=?`,
|
2026-08-18 15:25:58 +08:00
|
|
|
|
row.SourceRow, nullablePositiveInt(row.PrintSequence), nullableString(row.ShopName), nullableStringPreserveSpace(row.SourceSKURaw), row.SpecRaw,
|
2026-08-15 08:47:22 +08:00
|
|
|
|
row.InnerCode, row.SourceDuplicateCount, now, id)
|
|
|
|
|
|
if err != nil {
|
2026-08-15 09:12:23 +08:00
|
|
|
|
return "", innerCodeImportWriteError("更新档口入库码导入记录失败", err)
|
2026-08-15 08:47:22 +08:00
|
|
|
|
}
|
|
|
|
|
|
return InnerCodeImportUpdated, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 10:49:28 +08:00
|
|
|
|
func innerCodeStatusPreservesImportResult(status model.InnerCodeStatus) bool {
|
|
|
|
|
|
switch status {
|
2026-08-15 17:08:17 +08:00
|
|
|
|
case model.InnerCodeQueued, model.InnerCodeApplying, model.InnerCodeUpdated,
|
|
|
|
|
|
model.InnerCodeAlreadyFilled, model.InnerCodeNeedsCheck:
|
2026-08-15 10:49:28 +08:00
|
|
|
|
return true
|
|
|
|
|
|
default:
|
|
|
|
|
|
return false
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 09:12:23 +08:00
|
|
|
|
func innerCodeImportWriteError(action string, err error) error {
|
|
|
|
|
|
// 并发导入可能都在 SELECT 时看不到对方,最终仍由唯一索引裁决。
|
|
|
|
|
|
// 不把 MySQL 索引名或 SQL 原文显示给操作员。
|
|
|
|
|
|
var mysqlError *mysql.MySQLError
|
|
|
|
|
|
if errors.As(err, &mysqlError) && mysqlError.Number == 1062 {
|
|
|
|
|
|
return fmt.Errorf("%s: %w", action, ErrInnerCodeUniqueConflict)
|
|
|
|
|
|
}
|
|
|
|
|
|
return fmt.Errorf("%s: %w", action, err)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 09:41:45 +08:00
|
|
|
|
// ListInnerCodeRecordsForPlanning 返回某日选中且允许重新规划的记录。
|
|
|
|
|
|
func ListInnerCodeRecordsForPlanning(q Execer, businessDate string, ids []int64) ([]model.InnerCodeRecord, error) {
|
|
|
|
|
|
if len(ids) == 0 {
|
|
|
|
|
|
return nil, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
placeholders := make([]string, len(ids))
|
|
|
|
|
|
args := make([]any, 0, len(ids)+1)
|
|
|
|
|
|
args = append(args, businessDate)
|
|
|
|
|
|
for index, id := range ids {
|
|
|
|
|
|
placeholders[index] = "?"
|
|
|
|
|
|
args = append(args, id)
|
|
|
|
|
|
}
|
|
|
|
|
|
rows, err := q.Query(`SELECT `+innerCodeListColumns+`
|
2026-08-15 08:51:38 +08:00
|
|
|
|
FROM syb_inner_code_records
|
2026-08-15 09:41:45 +08:00
|
|
|
|
WHERE business_date=? AND id IN (`+strings.Join(placeholders, ",")+`)
|
2026-08-15 10:49:28 +08:00
|
|
|
|
AND deleted_at IS NULL
|
2026-08-15 09:41:45 +08:00
|
|
|
|
AND status IN ('pending','ready','skipped','failed')
|
|
|
|
|
|
ORDER BY source_row,id`, args...)
|
2026-08-15 08:51:38 +08:00
|
|
|
|
if err != nil {
|
2026-08-15 09:41:45 +08:00
|
|
|
|
return nil, fmt.Errorf("查询选中档口入库码失败: %w", err)
|
2026-08-15 08:51:38 +08:00
|
|
|
|
}
|
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
|
result := make([]model.InnerCodeRecord, 0)
|
|
|
|
|
|
for rows.Next() {
|
2026-08-15 09:41:45 +08:00
|
|
|
|
row, err := scanInnerCodeRecord(rows)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, err
|
2026-08-15 08:51:38 +08:00
|
|
|
|
}
|
|
|
|
|
|
result = append(result, row)
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := rows.Err(); err != nil {
|
2026-08-15 09:41:45 +08:00
|
|
|
|
return nil, fmt.Errorf("遍历选中档口入库码失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return result, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// ListInnerCodePlanningContext 返回选中订单的全部同日记录,用于预留未选记录已经占用的远端明细。
|
|
|
|
|
|
func ListInnerCodePlanningContext(q Execer, businessDate string, orderNumbers []string) ([]model.InnerCodeRecord, error) {
|
|
|
|
|
|
if len(orderNumbers) == 0 {
|
|
|
|
|
|
return nil, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
placeholders := make([]string, len(orderNumbers))
|
|
|
|
|
|
args := make([]any, 0, len(orderNumbers)+1)
|
|
|
|
|
|
args = append(args, businessDate)
|
|
|
|
|
|
for index, orderNumber := range orderNumbers {
|
|
|
|
|
|
placeholders[index] = "?"
|
|
|
|
|
|
args = append(args, orderNumber)
|
|
|
|
|
|
}
|
|
|
|
|
|
rows, err := q.Query(`SELECT `+innerCodeListColumns+`
|
|
|
|
|
|
FROM syb_inner_code_records
|
|
|
|
|
|
WHERE business_date=? AND order_number IN (`+strings.Join(placeholders, ",")+`)
|
2026-08-15 17:08:17 +08:00
|
|
|
|
AND (deleted_at IS NULL OR status IN ('queued','applying','updated','already_filled','needs_check'))
|
2026-08-15 09:41:45 +08:00
|
|
|
|
ORDER BY source_row,id`, args...)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("查询档口入库码匹配上下文失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
|
result := make([]model.InnerCodeRecord, 0)
|
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
|
row, err := scanInnerCodeRecord(rows)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, err
|
|
|
|
|
|
}
|
|
|
|
|
|
result = append(result, row)
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("遍历档口入库码匹配上下文失败: %w", err)
|
2026-08-15 08:51:38 +08:00
|
|
|
|
}
|
|
|
|
|
|
return result, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// SaveInnerCodePlans 在同一事务中保存一批只读规划结果。
|
|
|
|
|
|
func SaveInnerCodePlans(db *sql.DB, plans []model.InnerCodeRecord, plannedAt string) error {
|
|
|
|
|
|
tx, err := db.Begin()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("开始保存档口入库码规划事务失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer tx.Rollback()
|
2026-08-15 09:41:45 +08:00
|
|
|
|
if err := lockAndValidateInnerCodePlanClaims(tx, plans); err != nil {
|
|
|
|
|
|
return err
|
|
|
|
|
|
}
|
2026-08-15 08:51:38 +08:00
|
|
|
|
for _, plan := range plans {
|
|
|
|
|
|
result, err := tx.Exec(`
|
|
|
|
|
|
UPDATE syb_inner_code_records
|
|
|
|
|
|
SET stock_id=?,detail_id=?,syb_spec=?,syb_sku=?,syb_variation_sku=?,
|
|
|
|
|
|
purchase_platform=?,purchase_code=?,remote_inner_code=?,status=?,
|
2026-08-18 15:18:35 +08:00
|
|
|
|
remote_items_json=NULL,result_message=?,planned_at=?,apply_batch_id=NULL,apply_queued_at=NULL,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
apply_started_at=NULL,updated_at=?
|
2026-08-15 10:49:28 +08:00
|
|
|
|
WHERE id=? AND deleted_at IS NULL AND status IN ('pending','ready','skipped','failed')`,
|
2026-08-15 08:51:38 +08:00
|
|
|
|
nullablePositiveInt64(plan.StockID), nullablePositiveInt64(plan.DetailID),
|
|
|
|
|
|
nullableString(plan.SybSpec), nullableString(plan.SybSKU), nullableString(plan.SybVariationSKU),
|
|
|
|
|
|
nullableString(plan.PurchasePlatform), nullableString(plan.PurchaseCode),
|
|
|
|
|
|
nullableString(plan.RemoteInnerCode), plan.Status, nullableString(plan.ResultMessage),
|
|
|
|
|
|
plannedAt, plannedAt, plan.ID)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("保存档口入库码记录 %d 规划失败: %w", plan.ID, err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
|
|
|
|
|
|
return fmt.Errorf("档口入库码记录 %d 状态已变化,规划整体未保存", plan.ID)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := tx.Commit(); err != nil {
|
|
|
|
|
|
return fmt.Errorf("提交档口入库码规划失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 09:41:45 +08:00
|
|
|
|
type innerCodePlanOrderKey struct {
|
|
|
|
|
|
BusinessDate string
|
|
|
|
|
|
OrderNumber string
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// lockAndValidateInnerCodePlanClaims 串行化同订单的部分匹配,并拒绝覆盖未选记录已有的明细占用。
|
|
|
|
|
|
func lockAndValidateInnerCodePlanClaims(tx *sql.Tx, plans []model.InnerCodeRecord) error {
|
|
|
|
|
|
selected := make(map[int64]bool, len(plans))
|
|
|
|
|
|
keySet := make(map[innerCodePlanOrderKey]bool)
|
|
|
|
|
|
for _, plan := range plans {
|
|
|
|
|
|
selected[plan.ID] = true
|
|
|
|
|
|
keySet[innerCodePlanOrderKey{BusinessDate: plan.BusinessDate, OrderNumber: plan.OrderNumber}] = true
|
|
|
|
|
|
}
|
|
|
|
|
|
keys := make([]innerCodePlanOrderKey, 0, len(keySet))
|
|
|
|
|
|
for key := range keySet {
|
|
|
|
|
|
keys = append(keys, key)
|
|
|
|
|
|
}
|
|
|
|
|
|
sort.Slice(keys, func(i, j int) bool {
|
|
|
|
|
|
if keys[i].BusinessDate == keys[j].BusinessDate {
|
|
|
|
|
|
return keys[i].OrderNumber < keys[j].OrderNumber
|
|
|
|
|
|
}
|
|
|
|
|
|
return keys[i].BusinessDate < keys[j].BusinessDate
|
|
|
|
|
|
})
|
|
|
|
|
|
reserved := make(map[innerCodePlanOrderKey]map[int64]int64, len(keys))
|
|
|
|
|
|
for _, key := range keys {
|
|
|
|
|
|
rows, err := tx.Query(`SELECT id,COALESCE(detail_id,0),status
|
|
|
|
|
|
FROM syb_inner_code_records
|
2026-08-15 10:49:28 +08:00
|
|
|
|
WHERE business_date=? AND order_number=?
|
2026-08-15 17:08:17 +08:00
|
|
|
|
AND (deleted_at IS NULL OR status IN ('queued','applying','updated','already_filled','needs_check'))
|
2026-08-15 10:49:28 +08:00
|
|
|
|
ORDER BY id FOR UPDATE`, key.BusinessDate, key.OrderNumber)
|
2026-08-15 09:41:45 +08:00
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("锁定档口入库码订单匹配上下文失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
|
var id, detailID int64
|
|
|
|
|
|
var status model.InnerCodeStatus
|
|
|
|
|
|
if err := rows.Scan(&id, &detailID, &status); err != nil {
|
|
|
|
|
|
rows.Close()
|
|
|
|
|
|
return fmt.Errorf("读取档口入库码订单匹配上下文失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if selected[id] || detailID <= 0 || !innerCodeStatusHoldsDetail(status) {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
|
|
|
|
|
if reserved[key] == nil {
|
|
|
|
|
|
reserved[key] = make(map[int64]int64)
|
|
|
|
|
|
}
|
|
|
|
|
|
reserved[key][detailID] = id
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
|
rows.Close()
|
|
|
|
|
|
return fmt.Errorf("遍历档口入库码订单匹配上下文失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
rows.Close()
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, plan := range plans {
|
|
|
|
|
|
if plan.DetailID <= 0 || !innerCodeStatusHoldsDetail(plan.Status) {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
|
|
|
|
|
key := innerCodePlanOrderKey{BusinessDate: plan.BusinessDate, OrderNumber: plan.OrderNumber}
|
|
|
|
|
|
if ownerID := reserved[key][plan.DetailID]; ownerID > 0 {
|
|
|
|
|
|
return fmt.Errorf("档口入库码记录 %d 的顺运宝商品已被记录 %d 占用,请刷新后重新匹配", plan.ID, ownerID)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func innerCodeStatusHoldsDetail(status model.InnerCodeStatus) bool {
|
|
|
|
|
|
switch status {
|
2026-08-15 17:08:17 +08:00
|
|
|
|
case model.InnerCodeReady, model.InnerCodeQueued, model.InnerCodeApplying, model.InnerCodeUpdated,
|
2026-08-15 09:41:45 +08:00
|
|
|
|
model.InnerCodeAlreadyFilled, model.InnerCodeNeedsCheck:
|
|
|
|
|
|
return true
|
|
|
|
|
|
default:
|
|
|
|
|
|
return false
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 08:59:40 +08:00
|
|
|
|
const innerCodeListColumns = `id,business_date,source_row,COALESCE(print_sequence,0),order_number,
|
2026-08-18 15:25:58 +08:00
|
|
|
|
COALESCE(shop_name,''),stall,COALESCE(source_sku_raw,''),spec_raw,spec_key,inner_code,source_duplicate_count,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
COALESCE(apply_batch_id,''),COALESCE(apply_queued_at,''),
|
2026-08-15 08:59:40 +08:00
|
|
|
|
COALESCE(stock_id,0),COALESCE(detail_id,0),COALESCE(syb_spec,''),COALESCE(syb_sku,''),
|
|
|
|
|
|
COALESCE(syb_variation_sku,''),COALESCE(purchase_platform,''),COALESCE(purchase_code,''),
|
2026-08-18 15:18:35 +08:00
|
|
|
|
COALESCE(remote_inner_code,''),COALESCE(remote_items_json,''),status,COALESCE(result_message,''),created_by_user_id,
|
2026-08-15 08:59:40 +08:00
|
|
|
|
COALESCE(applied_by_user_id,''),COALESCE(planned_at,''),COALESCE(apply_started_at,''),
|
2026-08-15 10:49:28 +08:00
|
|
|
|
COALESCE(applied_at,''),COALESCE(deleted_at,''),COALESCE(deleted_by_user_id,''),created_at,updated_at`
|
2026-08-15 08:59:40 +08:00
|
|
|
|
|
|
|
|
|
|
func innerCodeFilterClause(filter InnerCodeListFilter) (string, []any) {
|
2026-08-15 10:49:28 +08:00
|
|
|
|
clauses := []string{"business_date=?", "deleted_at IS NULL"}
|
2026-08-15 08:59:40 +08:00
|
|
|
|
args := []any{filter.BusinessDate}
|
|
|
|
|
|
if filter.Status != "" {
|
|
|
|
|
|
clauses = append(clauses, "status=?")
|
|
|
|
|
|
args = append(args, filter.Status)
|
|
|
|
|
|
}
|
|
|
|
|
|
if filter.Keyword != "" {
|
|
|
|
|
|
like := "%" + escapeLike(filter.Keyword) + "%"
|
|
|
|
|
|
clauses = append(clauses, `(order_number LIKE ? ESCAPE '!' OR stall LIKE ? ESCAPE '!'
|
|
|
|
|
|
OR spec_raw LIKE ? ESCAPE '!' OR inner_code LIKE ? ESCAPE '!')`)
|
|
|
|
|
|
args = append(args, like, like, like, like)
|
|
|
|
|
|
}
|
|
|
|
|
|
return " WHERE " + strings.Join(clauses, " AND "), args
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// ListInnerCodeRecords 分页读取独立页面记录。
|
|
|
|
|
|
func ListInnerCodeRecords(q Execer, filter InnerCodeListFilter, limit, offset int) ([]model.InnerCodeRecord, error) {
|
|
|
|
|
|
where, args := innerCodeFilterClause(filter)
|
|
|
|
|
|
query := `SELECT ` + innerCodeListColumns + ` FROM syb_inner_code_records` + where +
|
|
|
|
|
|
` ORDER BY source_row,id LIMIT ? OFFSET ?`
|
|
|
|
|
|
args = append(args, limit, offset)
|
|
|
|
|
|
rows, err := q.Query(query, args...)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("查询档口入库码列表失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
|
result := make([]model.InnerCodeRecord, 0)
|
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
|
row, err := scanInnerCodeRecord(rows)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, err
|
|
|
|
|
|
}
|
|
|
|
|
|
result = append(result, row)
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("遍历档口入库码列表失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return result, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// CountInnerCodeRecords 统计与列表完全相同的筛选结果。
|
|
|
|
|
|
func CountInnerCodeRecords(q Execer, filter InnerCodeListFilter) (int, error) {
|
|
|
|
|
|
where, args := innerCodeFilterClause(filter)
|
|
|
|
|
|
var count int
|
|
|
|
|
|
if err := q.QueryRow(`SELECT COUNT(*) FROM syb_inner_code_records`+where, args...).Scan(&count); err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("统计档口入库码列表失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return count, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 17:08:17 +08:00
|
|
|
|
// CountInnerCodeStatuses 统计当前业务日期总量和后台回写关键状态。
|
2026-08-15 08:59:40 +08:00
|
|
|
|
func CountInnerCodeStatuses(q Execer, businessDate string) (InnerCodeStatusCounts, error) {
|
|
|
|
|
|
var result InnerCodeStatusCounts
|
|
|
|
|
|
err := q.QueryRow(`SELECT COUNT(*),
|
2026-08-15 17:08:17 +08:00
|
|
|
|
COALESCE(SUM(status='ready'),0),COALESCE(SUM(status='queued'),0),
|
|
|
|
|
|
COALESCE(SUM(status='applying'),0),COALESCE(SUM(status='needs_check'),0)
|
2026-08-15 10:49:28 +08:00
|
|
|
|
FROM syb_inner_code_records WHERE business_date=? AND deleted_at IS NULL`, businessDate).
|
2026-08-15 17:08:17 +08:00
|
|
|
|
Scan(&result.Total, &result.Ready, &result.Queued, &result.Applying, &result.NeedsCheck)
|
2026-08-15 08:59:40 +08:00
|
|
|
|
if err != nil {
|
|
|
|
|
|
return result, fmt.Errorf("统计档口入库码状态失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return result, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
type innerCodeRowScanner interface {
|
|
|
|
|
|
Scan(...any) error
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func scanInnerCodeRecord(scanner innerCodeRowScanner) (model.InnerCodeRecord, error) {
|
|
|
|
|
|
var row model.InnerCodeRecord
|
|
|
|
|
|
err := scanner.Scan(&row.ID, &row.BusinessDate, &row.SourceRow, &row.PrintSequence,
|
2026-08-18 15:25:58 +08:00
|
|
|
|
&row.OrderNumber, &row.ShopName, &row.Stall, &row.SourceSKURaw, &row.SpecRaw, &row.SpecKey,
|
2026-08-15 17:08:17 +08:00
|
|
|
|
&row.InnerCode, &row.SourceDuplicateCount, &row.ApplyBatchID, &row.ApplyQueuedAt,
|
|
|
|
|
|
&row.StockID, &row.DetailID,
|
2026-08-15 08:59:40 +08:00
|
|
|
|
&row.SybSpec, &row.SybSKU, &row.SybVariationSKU, &row.PurchasePlatform,
|
2026-08-18 15:18:35 +08:00
|
|
|
|
&row.PurchaseCode, &row.RemoteInnerCode, &row.RemoteItemsJSON, &row.Status, &row.ResultMessage,
|
2026-08-15 08:59:40 +08:00
|
|
|
|
&row.CreatedByUserID, &row.AppliedByUserID, &row.PlannedAt, &row.ApplyStartedAt,
|
2026-08-15 10:49:28 +08:00
|
|
|
|
&row.AppliedAt, &row.DeletedAt, &row.DeletedByUserID, &row.CreatedAt, &row.UpdatedAt)
|
2026-08-15 08:59:40 +08:00
|
|
|
|
if err != nil {
|
|
|
|
|
|
return row, fmt.Errorf("读取档口入库码记录失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return row, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 17:08:17 +08:00
|
|
|
|
// QueueInnerCodeApplyBatch 把所选 ready 记录原子加入同一后台批次。
|
|
|
|
|
|
// 任一记录不存在、已删除或状态变化时整批回滚,避免页面选择与实际队列不一致。
|
|
|
|
|
|
func QueueInnerCodeApplyBatch(db *sql.DB, ids []int64, batchID, actorUserID, queuedAt string) (int, error) {
|
|
|
|
|
|
if len(ids) == 0 {
|
|
|
|
|
|
return 0, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
placeholders := make([]string, len(ids))
|
|
|
|
|
|
args := make([]any, 0, len(ids)+4)
|
|
|
|
|
|
args = append(args, batchID, queuedAt, actorUserID, queuedAt)
|
|
|
|
|
|
for index, id := range ids {
|
|
|
|
|
|
placeholders[index] = "?"
|
|
|
|
|
|
args = append(args, id)
|
|
|
|
|
|
}
|
2026-08-15 09:10:10 +08:00
|
|
|
|
tx, err := db.Begin()
|
|
|
|
|
|
if err != nil {
|
2026-08-15 17:08:17 +08:00
|
|
|
|
return 0, fmt.Errorf("开始档口入库码排队事务失败: %w", err)
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
|
|
|
|
|
defer tx.Rollback()
|
2026-08-15 17:08:17 +08:00
|
|
|
|
result, err := tx.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET status='queued',apply_batch_id=?,apply_queued_at=?,applied_by_user_id=?,
|
|
|
|
|
|
result_message='已进入后台回写队列',apply_started_at=NULL,updated_at=?
|
|
|
|
|
|
WHERE deleted_at IS NULL AND status='ready' AND id IN (`+strings.Join(placeholders, ",")+`)`, args...)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("档口入库码加入后台队列失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
affected, err := result.RowsAffected()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("读取档口入库码排队数量失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected != int64(len(ids)) {
|
|
|
|
|
|
return 0, ErrInnerCodeApplyConflict
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := tx.Commit(); err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("提交档口入库码排队事务失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return int(affected), nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// ListQueuedInnerCodeIDs 返回批次下一组待处理记录。分组只控制数据库读取,不放宽逐条远端门禁。
|
|
|
|
|
|
func ListQueuedInnerCodeIDs(q Execer, batchID string, limit int) ([]int64, error) {
|
|
|
|
|
|
rows, err := q.Query(`SELECT id FROM syb_inner_code_records
|
|
|
|
|
|
WHERE apply_batch_id=? AND status='queued' ORDER BY id LIMIT ?`, batchID, limit)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("读取档口入库码后台队列失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
|
ids := make([]int64, 0, limit)
|
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
|
var id int64
|
|
|
|
|
|
if err := rows.Scan(&id); err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("读取档口入库码后台队列编号失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
ids = append(ids, id)
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("遍历档口入库码后台队列失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return ids, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// ClaimQueuedInnerCodeForApply 原子领取批次中的一条 queued 记录。
|
|
|
|
|
|
func ClaimQueuedInnerCodeForApply(db *sql.DB, id int64, batchID, actorUserID, now string) (*model.InnerCodeRecord, bool, error) {
|
|
|
|
|
|
tx, err := db.Begin()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, false, fmt.Errorf("开始领取档口入库码后台事务失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer tx.Rollback()
|
|
|
|
|
|
result, err := tx.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET status='applying',applied_by_user_id=?,apply_started_at=?,updated_at=?
|
|
|
|
|
|
WHERE id=? AND apply_batch_id=? AND status='queued'`, actorUserID, now, now, id, batchID)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, false, fmt.Errorf("领取档口入库码后台记录失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
affected, err := result.RowsAffected()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, false, fmt.Errorf("读取档口入库码后台领取结果失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected != 1 {
|
2026-08-15 09:10:10 +08:00
|
|
|
|
return nil, false, nil
|
|
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
record, err := scanInnerCodeRecord(tx.QueryRow(`SELECT `+innerCodeListColumns+
|
|
|
|
|
|
` FROM syb_inner_code_records WHERE id=? AND apply_batch_id=?`, id, batchID))
|
2026-08-15 09:10:10 +08:00
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, false, err
|
|
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
if err := tx.Commit(); err != nil {
|
|
|
|
|
|
return nil, false, fmt.Errorf("提交档口入库码后台领取失败: %w", err)
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
return &record, true, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// InterruptInnerCodeApplyBatch 收敛异常批次:未知远端结果转需核对,未开始记录恢复可回写。
|
|
|
|
|
|
func InterruptInnerCodeApplyBatch(db *sql.DB, batchID, interruptedAt string) (int, int, error) {
|
|
|
|
|
|
tx, err := db.Begin()
|
2026-08-15 09:10:10 +08:00
|
|
|
|
if err != nil {
|
2026-08-15 17:08:17 +08:00
|
|
|
|
return 0, 0, fmt.Errorf("开始收敛档口入库码后台批次失败: %w", err)
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
defer tx.Rollback()
|
|
|
|
|
|
applying, err := tx.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET status='needs_check',result_message='后台回写异常中断,远端结果不确定;系统不会自动重写',updated_at=?
|
|
|
|
|
|
WHERE apply_batch_id=? AND status='applying'`, interruptedAt, batchID)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, 0, fmt.Errorf("收敛档口入库码未知回写结果失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
queued, err := tx.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET status='ready',result_message='后台回写中断,本条尚未发送远端请求,已恢复为可回写',updated_at=?
|
|
|
|
|
|
WHERE apply_batch_id=? AND status='queued'`, interruptedAt, batchID)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, 0, fmt.Errorf("释放档口入库码后台队列失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
applyingCount, err := applying.RowsAffected()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, 0, fmt.Errorf("读取档口入库码未知结果数量失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
queuedCount, err := queued.RowsAffected()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, 0, fmt.Errorf("读取档口入库码后台释放数量失败: %w", err)
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
|
|
|
|
|
if err := tx.Commit(); err != nil {
|
2026-08-15 17:08:17 +08:00
|
|
|
|
return 0, 0, fmt.Errorf("提交档口入库码后台收敛失败: %w", err)
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
return int(applyingCount), int(queuedCount), nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// GetInnerCodeApplyBatchProgress 从现有业务表聚合批次进度。
|
|
|
|
|
|
func GetInnerCodeApplyBatchProgress(q Execer, batchID string) (*InnerCodeApplyBatchProgress, error) {
|
|
|
|
|
|
p := &InnerCodeApplyBatchProgress{BatchID: batchID}
|
|
|
|
|
|
err := q.QueryRow(`SELECT COUNT(*),
|
|
|
|
|
|
COALESCE(SUM(status='queued'),0),COALESCE(SUM(status='applying'),0),
|
|
|
|
|
|
COALESCE(SUM(status='updated'),0),COALESCE(SUM(status='already_filled'),0),
|
|
|
|
|
|
COALESCE(SUM(status='skipped'),0),COALESCE(SUM(status='failed'),0),
|
|
|
|
|
|
COALESCE(SUM(status='needs_check'),0),COALESCE(SUM(status='ready'),0)
|
|
|
|
|
|
FROM syb_inner_code_records WHERE apply_batch_id=?`, batchID).
|
|
|
|
|
|
Scan(&p.Total, &p.Queued, &p.Applying, &p.Updated, &p.AlreadyFilled,
|
|
|
|
|
|
&p.Skipped, &p.Failed, &p.NeedsCheck, &p.Ready)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, fmt.Errorf("读取档口入库码后台进度失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if p.Total == 0 {
|
|
|
|
|
|
return nil, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
return p, nil
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// FinishInnerCodeApply 保存一条已领取记录的最终结果。
|
|
|
|
|
|
func FinishInnerCodeApply(q Execer, id int64, status model.InnerCodeStatus, message, remoteCode, finishedAt string) error {
|
|
|
|
|
|
result, err := q.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET status=?,result_message=?,remote_inner_code=?,
|
|
|
|
|
|
applied_at=CASE WHEN ? IN ('updated','already_filled') THEN ? ELSE applied_at END,
|
|
|
|
|
|
updated_at=?
|
|
|
|
|
|
WHERE id=? AND status='applying'`, status, message, nullableString(remoteCode), status, finishedAt, finishedAt, id)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("保存档口入库码回写结果失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
|
|
|
|
|
|
return fmt.Errorf("档口入库码记录 %d 已不在回写中,拒绝覆盖结果", id)
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-18 15:18:35 +08:00
|
|
|
|
// SaveInnerCodeRemoteItems 在每个远端动作前后保存逐件检查点。
|
|
|
|
|
|
// 只有 applying 记录可更新,避免后台执行器覆盖已经被人工处理的结果。
|
|
|
|
|
|
func SaveInnerCodeRemoteItems(q Execer, id int64, remoteItemsJSON, message, updatedAt string) error {
|
|
|
|
|
|
result, err := q.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET remote_items_json=?,result_message=?,updated_at=?
|
|
|
|
|
|
WHERE id=? AND status='applying'`, nullableString(remoteItemsJSON), nullableString(message), updatedAt, id)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("保存档口入库码逐件检查点失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
|
|
|
|
|
|
return fmt.Errorf("档口入库码记录 %d 已不在回写中,拒绝保存逐件检查点", id)
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 09:10:10 +08:00
|
|
|
|
// GetInnerCodeForRecheck 读取一条需核对记录。
|
|
|
|
|
|
func GetInnerCodeForRecheck(q Execer, id int64) (*model.InnerCodeRecord, error) {
|
|
|
|
|
|
record, err := scanInnerCodeRecord(q.QueryRow(`SELECT `+innerCodeListColumns+
|
2026-08-15 10:49:28 +08:00
|
|
|
|
` FROM syb_inner_code_records WHERE id=? AND deleted_at IS NULL`, id))
|
2026-08-15 09:10:10 +08:00
|
|
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
|
|
|
|
return nil, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil, err
|
|
|
|
|
|
}
|
|
|
|
|
|
return &record, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 10:49:28 +08:00
|
|
|
|
// SoftDeleteInnerCodeRecords 隐藏整批记录并保留状态、规划和回写审计。
|
|
|
|
|
|
// 任一记录已经删除或不存在时回滚整批,避免页面提示的数量与实际不一致。
|
|
|
|
|
|
func SoftDeleteInnerCodeRecords(db *sql.DB, ids []int64, actorUserID, deletedAt string) (int, error) {
|
|
|
|
|
|
if len(ids) == 0 {
|
|
|
|
|
|
return 0, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
placeholders := make([]string, len(ids))
|
|
|
|
|
|
args := make([]any, 0, len(ids)+3)
|
|
|
|
|
|
args = append(args, deletedAt, actorUserID, deletedAt)
|
|
|
|
|
|
for index, id := range ids {
|
|
|
|
|
|
placeholders[index] = "?"
|
|
|
|
|
|
args = append(args, id)
|
|
|
|
|
|
}
|
|
|
|
|
|
tx, err := db.Begin()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("开始删除档口入库码事务失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
defer tx.Rollback()
|
|
|
|
|
|
result, err := tx.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET deleted_at=?,deleted_by_user_id=?,updated_at=?
|
|
|
|
|
|
WHERE deleted_at IS NULL AND id IN (`+strings.Join(placeholders, ",")+`)`, args...)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("软删除档口入库码记录失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
affected, err := result.RowsAffected()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("读取档口入库码删除数量失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected != int64(len(ids)) {
|
|
|
|
|
|
return 0, ErrInnerCodeDeleteConflict
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := tx.Commit(); err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("提交档口入库码删除事务失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return int(affected), nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 09:10:10 +08:00
|
|
|
|
// SaveInnerCodeRecheck 保存只读重新核对的远端结果,不执行状态领取或写入。
|
2026-08-18 15:18:35 +08:00
|
|
|
|
func SaveInnerCodeRecheck(q Execer, id int64, status model.InnerCodeStatus, message, remoteCode, remoteItemsJSON, checkedAt string) error {
|
2026-08-15 09:10:10 +08:00
|
|
|
|
result, err := q.Exec(`UPDATE syb_inner_code_records
|
2026-08-18 15:18:35 +08:00
|
|
|
|
SET status=?,result_message=?,remote_inner_code=?,remote_items_json=?,
|
2026-08-15 09:10:10 +08:00
|
|
|
|
applied_at=CASE WHEN ?='updated' THEN ? ELSE applied_at END,updated_at=?
|
2026-08-18 15:18:35 +08:00
|
|
|
|
WHERE id=? AND status='needs_check'`, status, message, nullableString(remoteCode), nullableString(remoteItemsJSON), status, checkedAt, checkedAt, id)
|
2026-08-15 09:10:10 +08:00
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("保存档口入库码核对结果失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
|
|
|
|
|
|
return fmt.Errorf("档口入库码记录 %d 已不需要核对", id)
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 17:08:17 +08:00
|
|
|
|
// InterruptApplyingInnerCodes 在启动时收敛后台状态:未知写入转需核对,未开始队列恢复可回写。
|
2026-08-15 09:10:10 +08:00
|
|
|
|
func InterruptApplyingInnerCodes(q Execer, interruptedAt string) (int, error) {
|
2026-08-15 17:08:17 +08:00
|
|
|
|
applying, err := q.Exec(`UPDATE syb_inner_code_records
|
2026-08-15 09:10:10 +08:00
|
|
|
|
SET status='needs_check',result_message='Admin 在回写完成前退出,请重新核对远端结果;系统不会自动重写',
|
|
|
|
|
|
updated_at=? WHERE status='applying'`, interruptedAt)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("恢复中断的档口入库码回写失败: %w", err)
|
|
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
applyingCount, err := applying.RowsAffected()
|
2026-08-15 09:10:10 +08:00
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("读取中断档口入库码数量失败: %w", err)
|
|
|
|
|
|
}
|
2026-08-15 17:08:17 +08:00
|
|
|
|
queued, err := q.Exec(`UPDATE syb_inner_code_records
|
|
|
|
|
|
SET status='ready',result_message='Admin 在后台回写开始前退出,本条未发送远端请求,已恢复为可回写',
|
|
|
|
|
|
updated_at=? WHERE status='queued'`, interruptedAt)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("恢复未开始的档口入库码队列失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
queuedCount, err := queued.RowsAffected()
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return 0, fmt.Errorf("读取恢复档口入库码队列数量失败: %w", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
return int(applyingCount + queuedCount), nil
|
2026-08-15 09:10:10 +08:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 08:47:22 +08:00
|
|
|
|
func nullablePositiveInt(value int) any {
|
|
|
|
|
|
if value <= 0 {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
return value
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 08:51:38 +08:00
|
|
|
|
func nullablePositiveInt64(value int64) any {
|
|
|
|
|
|
if value <= 0 {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
return value
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-08-15 08:47:22 +08:00
|
|
|
|
func nullableString(value string) any {
|
|
|
|
|
|
if value == "" {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
return value
|
|
|
|
|
|
}
|
2026-08-18 15:25:58 +08:00
|
|
|
|
|
|
|
|
|
|
func nullableStringPreserveSpace(value string) any {
|
|
|
|
|
|
if strings.TrimSpace(value) == "" {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
return value
|
|
|
|
|
|
}
|