feat: 档口入库码后台批量回写 (#246)
This commit is contained in:
+194
-44
@@ -21,6 +21,9 @@ var ErrInnerCodeRestoreConflict = errors.New("已回写或需核对的删除记
|
||||
// ErrInnerCodeDeleteConflict 表示批量删除时记录已经不可见或不存在,整批不会部分删除。
|
||||
var ErrInnerCodeDeleteConflict = errors.New("部分档口入库码记录已删除或不存在")
|
||||
|
||||
// ErrInnerCodeApplyConflict 表示所选记录有一条已不再可回写,整批不会部分入队。
|
||||
var ErrInnerCodeApplyConflict = errors.New("部分档口入库码记录已不再可回写")
|
||||
|
||||
// InnerCodeListFilter 是独立页面可组合的查询条件。
|
||||
type InnerCodeListFilter struct {
|
||||
BusinessDate string
|
||||
@@ -32,9 +35,30 @@ type InnerCodeListFilter struct {
|
||||
type InnerCodeStatusCounts struct {
|
||||
Total int
|
||||
Ready int
|
||||
Queued int
|
||||
Applying int
|
||||
NeedsCheck int
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// InnerCodeImportOutcome 说明幂等导入是新增、更新还是恢复软删除记录。
|
||||
type InnerCodeImportOutcome string
|
||||
|
||||
@@ -133,6 +157,7 @@ func UpsertInnerCodeImportRow(tx *sql.Tx, row model.InnerCodeImportRow, now stri
|
||||
status='pending',stock_id=NULL,detail_id=NULL,syb_spec=NULL,syb_sku=NULL,
|
||||
syb_variation_sku=NULL,purchase_platform=NULL,purchase_code=NULL,
|
||||
remote_inner_code=NULL,result_message=NULL,planned_at=NULL,
|
||||
apply_batch_id=NULL,apply_queued_at=NULL,apply_started_at=NULL,
|
||||
deleted_at=NULL,deleted_by_user_id=NULL,updated_at=?
|
||||
WHERE id=?`, row.SourceRow, nullablePositiveInt(row.PrintSequence), nullableString(row.ShopName),
|
||||
row.SpecRaw, row.InnerCode, row.SourceDuplicateCount, now, id)
|
||||
@@ -147,18 +172,21 @@ func UpsertInnerCodeImportRow(tx *sql.Tx, row model.InnerCodeImportRow, now stri
|
||||
UPDATE syb_inner_code_records
|
||||
SET source_row=?,print_sequence=?,shop_name=?,spec_raw=?,inner_code=?,
|
||||
source_duplicate_count=?,
|
||||
status=CASE WHEN status IN ('updated','already_filled','applying','needs_check')
|
||||
status=CASE WHEN status IN ('queued','updated','already_filled','applying','needs_check')
|
||||
THEN status ELSE 'pending' END,
|
||||
stock_id=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN stock_id ELSE NULL END,
|
||||
detail_id=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN detail_id ELSE NULL END,
|
||||
syb_spec=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN syb_spec ELSE NULL END,
|
||||
syb_sku=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN syb_sku ELSE NULL END,
|
||||
syb_variation_sku=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN syb_variation_sku ELSE NULL END,
|
||||
purchase_platform=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN purchase_platform ELSE NULL END,
|
||||
purchase_code=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN purchase_code ELSE NULL END,
|
||||
remote_inner_code=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN remote_inner_code ELSE NULL END,
|
||||
result_message=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN result_message ELSE NULL END,
|
||||
planned_at=CASE WHEN status IN ('updated','already_filled','applying','needs_check') THEN planned_at ELSE NULL END,
|
||||
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,
|
||||
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,
|
||||
updated_at=?
|
||||
WHERE id=?`,
|
||||
row.SourceRow, nullablePositiveInt(row.PrintSequence), nullableString(row.ShopName), row.SpecRaw,
|
||||
@@ -171,7 +199,8 @@ func UpsertInnerCodeImportRow(tx *sql.Tx, row model.InnerCodeImportRow, now stri
|
||||
|
||||
func innerCodeStatusPreservesImportResult(status model.InnerCodeStatus) bool {
|
||||
switch status {
|
||||
case model.InnerCodeApplying, model.InnerCodeUpdated, model.InnerCodeAlreadyFilled, model.InnerCodeNeedsCheck:
|
||||
case model.InnerCodeQueued, model.InnerCodeApplying, model.InnerCodeUpdated,
|
||||
model.InnerCodeAlreadyFilled, model.InnerCodeNeedsCheck:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
@@ -239,7 +268,7 @@ func ListInnerCodePlanningContext(q Execer, businessDate string, orderNumbers []
|
||||
rows, err := q.Query(`SELECT `+innerCodeListColumns+`
|
||||
FROM syb_inner_code_records
|
||||
WHERE business_date=? AND order_number IN (`+strings.Join(placeholders, ",")+`)
|
||||
AND (deleted_at IS NULL OR status IN ('applying','updated','already_filled','needs_check'))
|
||||
AND (deleted_at IS NULL OR status IN ('queued','applying','updated','already_filled','needs_check'))
|
||||
ORDER BY source_row,id`, args...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("查询档口入库码匹配上下文失败: %w", err)
|
||||
@@ -274,7 +303,8 @@ func SaveInnerCodePlans(db *sql.DB, plans []model.InnerCodeRecord, plannedAt str
|
||||
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=?,
|
||||
result_message=?,planned_at=?,updated_at=?
|
||||
result_message=?,planned_at=?,apply_batch_id=NULL,apply_queued_at=NULL,
|
||||
apply_started_at=NULL,updated_at=?
|
||||
WHERE id=? AND deleted_at IS NULL AND status IN ('pending','ready','skipped','failed')`,
|
||||
nullablePositiveInt64(plan.StockID), nullablePositiveInt64(plan.DetailID),
|
||||
nullableString(plan.SybSpec), nullableString(plan.SybSKU), nullableString(plan.SybVariationSKU),
|
||||
@@ -322,7 +352,7 @@ func lockAndValidateInnerCodePlanClaims(tx *sql.Tx, plans []model.InnerCodeRecor
|
||||
rows, err := tx.Query(`SELECT id,COALESCE(detail_id,0),status
|
||||
FROM syb_inner_code_records
|
||||
WHERE business_date=? AND order_number=?
|
||||
AND (deleted_at IS NULL OR status IN ('applying','updated','already_filled','needs_check'))
|
||||
AND (deleted_at IS NULL OR status IN ('queued','applying','updated','already_filled','needs_check'))
|
||||
ORDER BY id FOR UPDATE`, key.BusinessDate, key.OrderNumber)
|
||||
if err != nil {
|
||||
return fmt.Errorf("锁定档口入库码订单匹配上下文失败: %w", err)
|
||||
@@ -362,7 +392,7 @@ func lockAndValidateInnerCodePlanClaims(tx *sql.Tx, plans []model.InnerCodeRecor
|
||||
|
||||
func innerCodeStatusHoldsDetail(status model.InnerCodeStatus) bool {
|
||||
switch status {
|
||||
case model.InnerCodeReady, model.InnerCodeApplying, model.InnerCodeUpdated,
|
||||
case model.InnerCodeReady, model.InnerCodeQueued, model.InnerCodeApplying, model.InnerCodeUpdated,
|
||||
model.InnerCodeAlreadyFilled, model.InnerCodeNeedsCheck:
|
||||
return true
|
||||
default:
|
||||
@@ -372,6 +402,7 @@ func innerCodeStatusHoldsDetail(status model.InnerCodeStatus) bool {
|
||||
|
||||
const innerCodeListColumns = `id,business_date,source_row,COALESCE(print_sequence,0),order_number,
|
||||
COALESCE(shop_name,''),stall,spec_raw,spec_key,inner_code,source_duplicate_count,
|
||||
COALESCE(apply_batch_id,''),COALESCE(apply_queued_at,''),
|
||||
COALESCE(stock_id,0),COALESCE(detail_id,0),COALESCE(syb_spec,''),COALESCE(syb_sku,''),
|
||||
COALESCE(syb_variation_sku,''),COALESCE(purchase_platform,''),COALESCE(purchase_code,''),
|
||||
COALESCE(remote_inner_code,''),status,COALESCE(result_message,''),created_by_user_id,
|
||||
@@ -429,13 +460,14 @@ func CountInnerCodeRecords(q Execer, filter InnerCodeListFilter) (int, error) {
|
||||
return count, nil
|
||||
}
|
||||
|
||||
// CountInnerCodeStatuses 统计当前业务日期总量、可回写和需核对数量。
|
||||
// CountInnerCodeStatuses 统计当前业务日期总量和后台回写关键状态。
|
||||
func CountInnerCodeStatuses(q Execer, businessDate string) (InnerCodeStatusCounts, error) {
|
||||
var result InnerCodeStatusCounts
|
||||
err := q.QueryRow(`SELECT COUNT(*),
|
||||
COALESCE(SUM(status='ready'),0),COALESCE(SUM(status='needs_check'),0)
|
||||
COALESCE(SUM(status='ready'),0),COALESCE(SUM(status='queued'),0),
|
||||
COALESCE(SUM(status='applying'),0),COALESCE(SUM(status='needs_check'),0)
|
||||
FROM syb_inner_code_records WHERE business_date=? AND deleted_at IS NULL`, businessDate).
|
||||
Scan(&result.Total, &result.Ready, &result.NeedsCheck)
|
||||
Scan(&result.Total, &result.Ready, &result.Queued, &result.Applying, &result.NeedsCheck)
|
||||
if err != nil {
|
||||
return result, fmt.Errorf("统计档口入库码状态失败: %w", err)
|
||||
}
|
||||
@@ -450,7 +482,8 @@ func scanInnerCodeRecord(scanner innerCodeRowScanner) (model.InnerCodeRecord, er
|
||||
var row model.InnerCodeRecord
|
||||
err := scanner.Scan(&row.ID, &row.BusinessDate, &row.SourceRow, &row.PrintSequence,
|
||||
&row.OrderNumber, &row.ShopName, &row.Stall, &row.SpecRaw, &row.SpecKey,
|
||||
&row.InnerCode, &row.SourceDuplicateCount, &row.StockID, &row.DetailID,
|
||||
&row.InnerCode, &row.SourceDuplicateCount, &row.ApplyBatchID, &row.ApplyQueuedAt,
|
||||
&row.StockID, &row.DetailID,
|
||||
&row.SybSpec, &row.SybSKU, &row.SybVariationSKU, &row.PurchasePlatform,
|
||||
&row.PurchaseCode, &row.RemoteInnerCode, &row.Status, &row.ResultMessage,
|
||||
&row.CreatedByUserID, &row.AppliedByUserID, &row.PlannedAt, &row.ApplyStartedAt,
|
||||
@@ -461,41 +494,148 @@ func scanInnerCodeRecord(scanner innerCodeRowScanner) (model.InnerCodeRecord, er
|
||||
return row, nil
|
||||
}
|
||||
|
||||
// ClaimInnerCodeForApply 原子领取一条 ready 记录。返回 claimed=false 表示状态已变化。
|
||||
func ClaimInnerCodeForApply(db *sql.DB, id int64, actorUserID, now string) (*model.InnerCodeRecord, bool, error) {
|
||||
// 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)
|
||||
}
|
||||
tx, err := db.Begin()
|
||||
if err != nil {
|
||||
return nil, false, fmt.Errorf("开始领取档口入库码事务失败: %w", err)
|
||||
return 0, fmt.Errorf("开始档口入库码排队事务失败: %w", err)
|
||||
}
|
||||
defer tx.Rollback()
|
||||
record, err := scanInnerCodeRecord(tx.QueryRow(`SELECT `+innerCodeListColumns+
|
||||
` FROM syb_inner_code_records WHERE id=? AND deleted_at IS NULL FOR UPDATE`, id))
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
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 {
|
||||
return nil, false, nil
|
||||
}
|
||||
record, err := scanInnerCodeRecord(tx.QueryRow(`SELECT `+innerCodeListColumns+
|
||||
` FROM syb_inner_code_records WHERE id=? AND apply_batch_id=?`, id, batchID))
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
if record.Status != model.InnerCodeReady {
|
||||
return &record, false, nil
|
||||
if err := tx.Commit(); err != nil {
|
||||
return nil, false, fmt.Errorf("提交档口入库码后台领取失败: %w", err)
|
||||
}
|
||||
result, err := tx.Exec(`UPDATE syb_inner_code_records
|
||||
SET status='applying',applied_by_user_id=?,apply_started_at=?,updated_at=?
|
||||
WHERE id=? AND deleted_at IS NULL AND status='ready'`, actorUserID, now, now, id)
|
||||
return &record, true, nil
|
||||
}
|
||||
|
||||
// InterruptInnerCodeApplyBatch 收敛异常批次:未知远端结果转需核对,未开始记录恢复可回写。
|
||||
func InterruptInnerCodeApplyBatch(db *sql.DB, batchID, interruptedAt string) (int, int, error) {
|
||||
tx, err := db.Begin()
|
||||
if err != nil {
|
||||
return nil, false, fmt.Errorf("领取档口入库码记录失败: %w", err)
|
||||
return 0, 0, fmt.Errorf("开始收敛档口入库码后台批次失败: %w", err)
|
||||
}
|
||||
if affected, err := result.RowsAffected(); err != nil || affected != 1 {
|
||||
return &record, false, nil
|
||||
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)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return nil, false, fmt.Errorf("提交档口入库码领取失败: %w", err)
|
||||
return 0, 0, fmt.Errorf("提交档口入库码后台收敛失败: %w", err)
|
||||
}
|
||||
record.Status = model.InnerCodeApplying
|
||||
record.AppliedByUserID = actorUserID
|
||||
record.ApplyStartedAt = now
|
||||
record.UpdatedAt = now
|
||||
return &record, true, nil
|
||||
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
|
||||
}
|
||||
|
||||
// FinishInnerCodeApply 保存一条已领取记录的最终结果。
|
||||
@@ -579,19 +719,29 @@ func SaveInnerCodeRecheck(q Execer, id int64, status model.InnerCodeStatus, mess
|
||||
return nil
|
||||
}
|
||||
|
||||
// InterruptApplyingInnerCodes 在启动时把未知结果的 applying 收敛为 needs_check。
|
||||
// InterruptApplyingInnerCodes 在启动时收敛后台状态:未知写入转需核对,未开始队列恢复可回写。
|
||||
func InterruptApplyingInnerCodes(q Execer, interruptedAt string) (int, error) {
|
||||
result, err := q.Exec(`UPDATE syb_inner_code_records
|
||||
applying, err := q.Exec(`UPDATE syb_inner_code_records
|
||||
SET status='needs_check',result_message='Admin 在回写完成前退出,请重新核对远端结果;系统不会自动重写',
|
||||
updated_at=? WHERE status='applying'`, interruptedAt)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("恢复中断的档口入库码回写失败: %w", err)
|
||||
}
|
||||
affected, err := result.RowsAffected()
|
||||
applyingCount, err := applying.RowsAffected()
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("读取中断档口入库码数量失败: %w", err)
|
||||
}
|
||||
return int(affected), nil
|
||||
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
|
||||
}
|
||||
|
||||
func nullablePositiveInt(value int) any {
|
||||
|
||||
@@ -24,11 +24,11 @@ func TestInterruptApplyingInnerCodes_只收敛回写中记录(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := db.Exec(`INSERT INTO syb_inner_code_records(id,status,updated_at) VALUES
|
||||
(1,'applying','old'),(2,'ready','old')`); err != nil {
|
||||
(1,'applying','old'),(2,'ready','old'),(3,'queued','old')`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
count, err := InterruptApplyingInnerCodes(db, "2026-08-15T01:00:00Z")
|
||||
if err != nil || count != 1 {
|
||||
if err != nil || count != 2 {
|
||||
t.Fatalf("count=%d err=%v", count, err)
|
||||
}
|
||||
var status, message, updatedAt string
|
||||
@@ -42,6 +42,9 @@ func TestInterruptApplyingInnerCodes_只收敛回写中记录(t *testing.T) {
|
||||
if err := db.QueryRow(`SELECT status FROM syb_inner_code_records WHERE id=2`).Scan(&status); err != nil || status != "ready" {
|
||||
t.Fatalf("ready 记录不应变化 status=%s err=%v", status, err)
|
||||
}
|
||||
if err := db.QueryRow(`SELECT status,result_message FROM syb_inner_code_records WHERE id=3`).Scan(&status, &message); err != nil || status != "ready" || message == "" {
|
||||
t.Fatalf("queued 记录应恢复可回写 status=%s message=%q err=%v", status, message, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestInnerCodeImportWriteError_唯一冲突不泄漏索引细节(t *testing.T) {
|
||||
@@ -65,7 +68,7 @@ func TestSoftDeleteInnerCodeRecords_所有状态只写删除审计(t *testing.T)
|
||||
)`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
statuses := []string{"pending", "ready", "applying", "updated", "already_filled", "skipped", "failed", "needs_check"}
|
||||
statuses := []string{"pending", "ready", "queued", "applying", "updated", "already_filled", "skipped", "failed", "needs_check"}
|
||||
for index, status := range statuses {
|
||||
if _, err := db.Exec(`INSERT INTO syb_inner_code_records(id,status,updated_at) VALUES(?,?,?)`, index+1, status, "old"); err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -120,6 +123,95 @@ func TestSoftDeleteInnerCodeRecords_有失效ID时整批回滚(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueueInnerCodeApplyBatch_状态冲突时整批回滚(t *testing.T) {
|
||||
db, err := sql.Open("sqlite", "file:inner_code_queue_conflict?mode=memory&cache=shared")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
if _, err := db.Exec(`CREATE TABLE syb_inner_code_records (
|
||||
id INTEGER PRIMARY KEY,status TEXT NOT NULL,apply_batch_id TEXT,apply_queued_at TEXT,
|
||||
applied_by_user_id TEXT,result_message TEXT,apply_started_at TEXT,updated_at TEXT,
|
||||
deleted_at TEXT
|
||||
); INSERT INTO syb_inner_code_records(id,status,updated_at) VALUES
|
||||
(1,'ready','old'),(2,'pending','old')`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := QueueInnerCodeApplyBatch(db, []int64{1, 2}, "ICB-1", "user-1", "2026-08-15T05:00:00Z"); !errors.Is(err, ErrInnerCodeApplyConflict) {
|
||||
t.Fatalf("期望整批排队冲突,实际 %v", err)
|
||||
}
|
||||
var status string
|
||||
var batchID sql.NullString
|
||||
if err := db.QueryRow(`SELECT status,apply_batch_id FROM syb_inner_code_records WHERE id=1`).Scan(&status, &batchID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if status != "ready" || batchID.Valid {
|
||||
t.Fatalf("冲突后不应部分入队 status=%s batch=%+v", status, batchID)
|
||||
}
|
||||
}
|
||||
|
||||
func TestInnerCodeApplyBatchProgress_聚合终态与恢复记录(t *testing.T) {
|
||||
db, err := sql.Open("sqlite", "file:inner_code_batch_progress?mode=memory&cache=shared")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
if _, err := db.Exec(`CREATE TABLE syb_inner_code_records (
|
||||
id INTEGER PRIMARY KEY,status TEXT NOT NULL,apply_batch_id TEXT
|
||||
); INSERT INTO syb_inner_code_records(id,status,apply_batch_id) VALUES
|
||||
(1,'queued','ICB-1'),(2,'applying','ICB-1'),(3,'updated','ICB-1'),
|
||||
(4,'failed','ICB-1'),(5,'needs_check','ICB-1'),(6,'ready','ICB-1')`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
progress, err := GetInnerCodeApplyBatchProgress(db, "ICB-1")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if progress == nil || progress.Total != 6 || progress.Queued != 1 || progress.Applying != 1 ||
|
||||
progress.Updated != 1 || progress.Failed != 1 || progress.NeedsCheck != 1 ||
|
||||
progress.Ready != 1 || progress.Processed() != 3 {
|
||||
t.Fatalf("progress=%+v", progress)
|
||||
}
|
||||
}
|
||||
|
||||
func TestInterruptInnerCodeApplyBatch_区分未知结果和未开始记录(t *testing.T) {
|
||||
db, err := sql.Open("sqlite", "file:inner_code_batch_interrupt?mode=memory&cache=shared")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
if _, err := db.Exec(`CREATE TABLE syb_inner_code_records (
|
||||
id INTEGER PRIMARY KEY,status TEXT NOT NULL,apply_batch_id TEXT,result_message TEXT,updated_at TEXT
|
||||
); INSERT INTO syb_inner_code_records(id,status,apply_batch_id,updated_at) VALUES
|
||||
(1,'applying','ICB-1','old'),(2,'queued','ICB-1','old'),(3,'queued','ICB-2','old')`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
needsCheck, released, err := InterruptInnerCodeApplyBatch(db, "ICB-1", "2026-08-15T06:00:00Z")
|
||||
if err != nil || needsCheck != 1 || released != 1 {
|
||||
t.Fatalf("needs_check=%d released=%d err=%v", needsCheck, released, err)
|
||||
}
|
||||
rows, err := db.Query(`SELECT id,status,result_message,updated_at FROM syb_inner_code_records ORDER BY id`)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer rows.Close()
|
||||
wants := []string{"needs_check", "ready", "queued"}
|
||||
index := 0
|
||||
for rows.Next() {
|
||||
var id int
|
||||
var status, message, updatedAt string
|
||||
var nullableMessage sql.NullString
|
||||
if err := rows.Scan(&id, &status, &nullableMessage, &updatedAt); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
message = nullableMessage.String
|
||||
if status != wants[index] || (id < 3 && (message == "" || updatedAt != "2026-08-15T06:00:00Z")) {
|
||||
t.Fatalf("id=%d status=%s message=%q updated=%s", id, status, message, updatedAt)
|
||||
}
|
||||
index++
|
||||
}
|
||||
}
|
||||
|
||||
func TestFinishInnerCodeApply_软删除后仍保存后台结果(t *testing.T) {
|
||||
db, err := sql.Open("sqlite", "file:inner_code_finish_deleted?mode=memory&cache=shared")
|
||||
if err != nil {
|
||||
@@ -153,7 +245,7 @@ func TestInnerCodeStatusPreservesImportResult(t *testing.T) {
|
||||
t.Errorf("%s 不应保留旧规划", status)
|
||||
}
|
||||
}
|
||||
for _, status := range []string{"applying", "updated", "already_filled", "needs_check"} {
|
||||
for _, status := range []string{"queued", "applying", "updated", "already_filled", "needs_check"} {
|
||||
if !innerCodeStatusPreservesImportResult(model.InnerCodeStatus(status)) {
|
||||
t.Errorf("%s 应保留远端结果和审计", status)
|
||||
}
|
||||
|
||||
@@ -20,7 +20,7 @@ import (
|
||||
"cmautobuy/admin/spec"
|
||||
)
|
||||
|
||||
const mysqlSchemaVersion = 24
|
||||
const mysqlSchemaVersion = 25
|
||||
|
||||
// OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。
|
||||
func OpenMySQL(cfg config.DatabaseConfig) (*sql.DB, error) {
|
||||
@@ -712,10 +712,64 @@ func MigrateMySQL(db *sql.DB) error {
|
||||
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 24, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
|
||||
return fmt.Errorf("记录 MySQL schema v24 失败: %w", err)
|
||||
}
|
||||
current = 24
|
||||
}
|
||||
if current < 25 {
|
||||
if err := migrateMySQLV25(db); err != nil {
|
||||
return fmt.Errorf("执行 MySQL schema v25 失败: %w", err)
|
||||
}
|
||||
if err := checkMySQLV25Shape(db); err != nil {
|
||||
return fmt.Errorf("MySQL schema v25 自检失败,未记录版本: %w", err)
|
||||
}
|
||||
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 25, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
|
||||
return fmt.Errorf("记录 MySQL schema v25 失败: %w", err)
|
||||
}
|
||||
}
|
||||
return CheckMySQLSchema(db)
|
||||
}
|
||||
|
||||
// migrateMySQLV25 在现有档口入库码业务表上增加后台批次状态,不另建批次表。
|
||||
func migrateMySQLV25(db *sql.DB) error {
|
||||
for _, column := range []struct{ name, ddl string }{
|
||||
{"apply_batch_id", `ALTER TABLE syb_inner_code_records ADD COLUMN apply_batch_id VARCHAR(191) COLLATE utf8mb4_bin NULL AFTER source_duplicate_count`},
|
||||
{"apply_queued_at", `ALTER TABLE syb_inner_code_records ADD COLUMN apply_queued_at VARCHAR(35) NULL AFTER apply_batch_id`},
|
||||
} {
|
||||
exists, err := mysqlColumnExists(db, "syb_inner_code_records", column.name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !exists {
|
||||
if _, err := db.Exec(column.ddl); err != nil {
|
||||
return fmt.Errorf("增加 syb_inner_code_records.%s 失败: %w", column.name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
clause, exists, err := mysqlCheckConstraintClause(db, "syb_inner_code_records", "chk_inner_code_status")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if exists && !strings.Contains(strings.ToLower(clause), "queued") {
|
||||
if _, err := db.Exec(`ALTER TABLE syb_inner_code_records DROP CHECK chk_inner_code_status`); err != nil {
|
||||
return fmt.Errorf("更新档口入库码状态约束前删除旧约束失败: %w", err)
|
||||
}
|
||||
exists = false
|
||||
}
|
||||
if !exists {
|
||||
if _, err := db.Exec(`ALTER TABLE syb_inner_code_records ADD CONSTRAINT chk_inner_code_status
|
||||
CHECK (status IN ('pending','ready','queued','applying','updated','already_filled','skipped','failed','needs_check'))`); err != nil {
|
||||
return fmt.Errorf("增加档口入库码后台状态约束失败: %w", err)
|
||||
}
|
||||
}
|
||||
if exists, err := mysqlIndexExists(db, "syb_inner_code_records", "idx_inner_code_apply_batch"); err != nil {
|
||||
return err
|
||||
} else if !exists {
|
||||
if _, err := db.Exec(`ALTER TABLE syb_inner_code_records ADD INDEX idx_inner_code_apply_batch (apply_batch_id,status,id)`); err != nil {
|
||||
return fmt.Errorf("增加档口入库码后台批次索引失败: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// migrateMySQLV24 为档口入库码增加可恢复软删除审计。
|
||||
func migrateMySQLV24(db *sql.DB) error {
|
||||
for _, column := range []struct{ name, ddl string }{
|
||||
@@ -1939,6 +1993,23 @@ func mysqlConstraintExists(db *sql.DB, table, constraint string) (bool, error) {
|
||||
return count == 1, nil
|
||||
}
|
||||
|
||||
func mysqlCheckConstraintClause(db *sql.DB, table, constraint string) (string, bool, error) {
|
||||
var clause string
|
||||
err := db.QueryRow(`SELECT cc.check_clause
|
||||
FROM information_schema.table_constraints tc
|
||||
JOIN information_schema.check_constraints cc
|
||||
ON cc.constraint_schema=tc.constraint_schema AND cc.constraint_name=tc.constraint_name
|
||||
WHERE tc.constraint_schema=DATABASE() AND tc.table_name=? AND tc.constraint_name=?
|
||||
AND tc.constraint_type='CHECK'`, table, constraint).Scan(&clause)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return "", false, nil
|
||||
}
|
||||
if err != nil {
|
||||
return "", false, fmt.Errorf("读取 MySQL 约束 %s.%s 表达式失败: %w", table, constraint, err)
|
||||
}
|
||||
return clause, true, nil
|
||||
}
|
||||
|
||||
func mysqlIndexExists(db *sql.DB, table, index string) (bool, error) {
|
||||
var count int
|
||||
if err := db.QueryRow(`SELECT COUNT(DISTINCT index_name) FROM information_schema.statistics
|
||||
@@ -2181,7 +2252,27 @@ func CheckMySQLSchema(db *sql.DB) error {
|
||||
if err := checkMySQLV23Shape(db); err != nil {
|
||||
return err
|
||||
}
|
||||
return checkMySQLV24Shape(db)
|
||||
if err := checkMySQLV24Shape(db); err != nil {
|
||||
return err
|
||||
}
|
||||
return checkMySQLV25Shape(db)
|
||||
}
|
||||
|
||||
func checkMySQLV25Shape(db *sql.DB) error {
|
||||
for _, name := range []string{"apply_batch_id", "apply_queued_at"} {
|
||||
exists, err := mysqlColumnExists(db, "syb_inner_code_records", name)
|
||||
if err != nil || !exists {
|
||||
return fmt.Errorf("档口入库码后台回写字段 %s 缺失: %v", name, err)
|
||||
}
|
||||
}
|
||||
if exists, err := mysqlIndexExists(db, "syb_inner_code_records", "idx_inner_code_apply_batch"); err != nil || !exists {
|
||||
return fmt.Errorf("档口入库码后台回写索引 idx_inner_code_apply_batch 缺失: %v", err)
|
||||
}
|
||||
clause, exists, err := mysqlCheckConstraintClause(db, "syb_inner_code_records", "chk_inner_code_status")
|
||||
if err != nil || !exists || !strings.Contains(strings.ToLower(clause), "queued") {
|
||||
return fmt.Errorf("档口入库码状态约束未包含 queued: %v", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func checkMySQLV24Shape(db *sql.DB) error {
|
||||
|
||||
@@ -1094,7 +1094,7 @@ func TestMySQLMigrate_V22升级V23且唯一约束生效(t *testing.T) {
|
||||
}
|
||||
mustExec(t, db, `SET FOREIGN_KEY_CHECKS=0`)
|
||||
mustExec(t, db, `DROP TABLE syb_inner_code_records`)
|
||||
mustExec(t, db, `DELETE FROM schema_migrations WHERE version=23`)
|
||||
mustExec(t, db, `DELETE FROM schema_migrations WHERE version>=23`)
|
||||
mustExec(t, db, `SET FOREIGN_KEY_CHECKS=1`)
|
||||
if err := MigrateMySQL(db); err != nil {
|
||||
t.Fatalf("v22 升级 v23 失败: %v", err)
|
||||
@@ -1129,7 +1129,7 @@ func TestMySQLMigrate_V23升级V24且软删除字段有效(t *testing.T) {
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records DROP FOREIGN KEY fk_inner_code_deleted_by`)
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records DROP INDEX idx_inner_code_deleted`)
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records DROP COLUMN deleted_by_user_id,DROP COLUMN deleted_at`)
|
||||
mustExec(t, db, `DELETE FROM schema_migrations WHERE version=24`)
|
||||
mustExec(t, db, `DELETE FROM schema_migrations WHERE version>=24`)
|
||||
if err := MigrateMySQL(db); err != nil {
|
||||
t.Fatalf("v23 升级 v24 失败: %v", err)
|
||||
}
|
||||
@@ -1141,6 +1141,31 @@ func TestMySQLMigrate_V23升级V24且软删除字段有效(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMySQLMigrate_V24升级V25且后台队列状态有效(t *testing.T) {
|
||||
db := openMySQLMigrationTestDB(t)
|
||||
defer db.Close()
|
||||
cleanMySQLTestSchema(t, db)
|
||||
defer cleanMySQLTestSchema(t, db)
|
||||
if err := MigrateMySQL(db); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records DROP INDEX idx_inner_code_apply_batch`)
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records DROP COLUMN apply_queued_at,DROP COLUMN apply_batch_id`)
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records DROP CHECK chk_inner_code_status`)
|
||||
mustExec(t, db, `ALTER TABLE syb_inner_code_records ADD CONSTRAINT chk_inner_code_status
|
||||
CHECK (status IN ('pending','ready','applying','updated','already_filled','skipped','failed','needs_check'))`)
|
||||
mustExec(t, db, `DELETE FROM schema_migrations WHERE version>=25`)
|
||||
if err := MigrateMySQL(db); err != nil {
|
||||
t.Fatalf("v24 升级 v25 失败: %v", err)
|
||||
}
|
||||
if err := MigrateMySQL(db); err != nil {
|
||||
t.Fatalf("v25 重复迁移失败: %v", err)
|
||||
}
|
||||
if err := checkMySQLV25Shape(db); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpsertInnerCodeImportRow_软删除记录按状态安全恢复(t *testing.T) {
|
||||
db := openMySQLMigrationTestDB(t)
|
||||
defer db.Close()
|
||||
|
||||
Reference in New Issue
Block a user