feat: 增加顺运宝批量 AI 规格匹配 (#202)

This commit is contained in:
chengma
2026-08-14 10:30:08 +08:00
parent 0201dab98c
commit d64579e9b3
26 changed files with 1289 additions and 26 deletions
+230
View File
@@ -0,0 +1,230 @@
package repository
import (
"database/sql"
"errors"
"fmt"
"cmautobuy/admin/model"
)
func InsertAIMatchBatch(q Execer, batch model.AIMatchBatch, items []model.AIMatchBatchItem) error {
_, err := q.Exec(`INSERT INTO ai_match_batches(batch_id,status,created_by_user_id,total_count,
provider_id,provider_name,provider_base_url,model,timeout_seconds,max_concurrency,
confidence_threshold_bps,config_fingerprint,rules_version,prompt_version,created_at)
VALUES(?,?,?, ?,?,?,?,?,?,?,?,?,?,?,?)`, batch.BatchID, batch.Status, batch.CreatedByUserID,
batch.TotalCount, batch.ProviderID, batch.ProviderName, batch.ProviderBaseURL, batch.Model,
batch.TimeoutSeconds, batch.MaxConcurrency, batch.ConfidenceThresholdBPS,
batch.ConfigFingerprint, batch.RulesVersion, batch.PromptVersion, batch.CreatedAt)
if err != nil {
return fmt.Errorf("创建 AI 匹配批次失败: %w", err)
}
for _, item := range items {
if _, err := q.Exec(`INSERT INTO ai_match_batch_items(item_id,batch_id,syb_id,identity_hash,
leader_syb_id,context_version,position,status) VALUES(?,?,?,?,?,?,?,?)`, item.ItemID,
item.BatchID, item.SybID, item.IdentityHash, item.LeaderSybID, item.ContextVersion,
item.Position, item.Status); err != nil {
return fmt.Errorf("创建 AI 匹配批次明细失败: %w", err)
}
}
return nil
}
func MarkAIMatchBatchRunning(q Execer, batchID, startedAt string) error {
result, err := q.Exec(`UPDATE ai_match_batches SET status='running',started_at=?
WHERE batch_id=? AND status='queued'`, startedAt, batchID)
if err != nil {
return fmt.Errorf("启动 AI 匹配批次失败: %w", err)
}
affected, err := result.RowsAffected()
if err != nil {
return err
}
if affected != 1 {
return fmt.Errorf("AI 匹配批次不是待运行状态")
}
return nil
}
func ListAIMatchLeaderItems(q Execer, batchID string) ([]model.AIMatchBatchItem, error) {
rows, err := q.Query(`SELECT item_id,batch_id,syb_id,identity_hash,leader_syb_id,
context_version,position,status,outcome,message,mapping_source,option_key,confidence_bps,
model_called,started_at,finished_at FROM ai_match_batch_items
WHERE batch_id=? AND syb_id=leader_syb_id AND status='queued' ORDER BY position`, batchID)
if err != nil {
return nil, fmt.Errorf("读取 AI 匹配批次工作项失败: %w", err)
}
defer rows.Close()
return scanAIMatchItems(rows)
}
func ListAIMatchIdentityItems(q Execer, batchID, identityHash string) ([]model.AIMatchBatchItem, error) {
rows, err := q.Query(`SELECT item_id,batch_id,syb_id,identity_hash,leader_syb_id,
context_version,position,status,outcome,message,mapping_source,option_key,confidence_bps,
model_called,started_at,finished_at FROM ai_match_batch_items
WHERE batch_id=? AND identity_hash=? ORDER BY position`, batchID, identityHash)
if err != nil {
return nil, fmt.Errorf("读取 AI 匹配批次同规格明细失败: %w", err)
}
defer rows.Close()
return scanAIMatchItems(rows)
}
func ListAIMatchBatchItems(q Execer, batchID string) ([]model.AIMatchBatchItem, error) {
rows, err := q.Query(`SELECT item_id,batch_id,syb_id,identity_hash,leader_syb_id,
context_version,position,status,outcome,message,mapping_source,option_key,confidence_bps,
model_called,started_at,finished_at FROM ai_match_batch_items WHERE batch_id=? ORDER BY position`, batchID)
if err != nil {
return nil, fmt.Errorf("读取 AI 匹配批次明细失败: %w", err)
}
defer rows.Close()
return scanAIMatchItems(rows)
}
func scanAIMatchItems(rows *sql.Rows) ([]model.AIMatchBatchItem, error) {
var result []model.AIMatchBatchItem
for rows.Next() {
var item model.AIMatchBatchItem
var outcome, message, source, optionKey, startedAt, finishedAt sql.NullString
var confidence sql.NullInt64
var called int
if err := rows.Scan(&item.ItemID, &item.BatchID, &item.SybID, &item.IdentityHash,
&item.LeaderSybID, &item.ContextVersion, &item.Position, &item.Status, &outcome,
&message, &source, &optionKey, &confidence, &called, &startedAt, &finishedAt); err != nil {
return nil, err
}
item.Outcome, item.Message, item.MappingSource, item.OptionKey = outcome.String, message.String, source.String, optionKey.String
item.ConfidenceBPS, item.ConfidenceSet = int(confidence.Int64), confidence.Valid
item.ModelCalled, item.StartedAt, item.FinishedAt = called == 1, startedAt.String, finishedAt.String
result = append(result, item)
}
return result, rows.Err()
}
func MarkAIMatchItemRunning(q Execer, itemID, startedAt string) error {
_, err := q.Exec(`UPDATE ai_match_batch_items SET status='running',started_at=? WHERE item_id=? AND status='queued'`, startedAt, itemID)
if err != nil {
return fmt.Errorf("标记 AI 匹配明细运行失败: %w", err)
}
return nil
}
func FinishAIMatchItem(q Execer, item model.AIMatchBatchItem) error {
_, err := q.Exec(`UPDATE ai_match_batch_items SET status=?,outcome=?,message=?,mapping_source=?,
option_key=?,confidence_bps=?,model_called=?,finished_at=? WHERE item_id=? AND status IN ('queued','running')`,
item.Status, nullableText(item.Outcome), nullableText(item.Message), nullableText(item.MappingSource),
nullableText(item.OptionKey), nullableInt(item.ConfidenceBPS, item.ConfidenceSet), item.ModelCalled,
item.FinishedAt, item.ItemID)
if err != nil {
return fmt.Errorf("完成 AI 匹配明细失败: %w", err)
}
return nil
}
type AIMatchBatchCounts struct{ Processed, Success, Reused, Manual, Failed int }
func CountAIMatchBatchItems(q Execer, batchID string) (AIMatchBatchCounts, error) {
var result AIMatchBatchCounts
err := q.QueryRow(`SELECT
SUM(status NOT IN ('queued','running')),
SUM(status='succeeded'),SUM(status='reused'),SUM(status IN ('manual','stale')),
SUM(status IN ('failed','interrupted')) FROM ai_match_batch_items WHERE batch_id=?`, batchID).
Scan(&result.Processed, &result.Success, &result.Reused, &result.Manual, &result.Failed)
if err != nil {
return result, fmt.Errorf("统计 AI 匹配批次失败: %w", err)
}
return result, nil
}
func FinishAIMatchBatch(q Execer, batchID, status, errorMessage, finishedAt string, counts AIMatchBatchCounts) error {
_, err := q.Exec(`UPDATE ai_match_batches SET status=?,processed_count=?,success_count=?,reused_count=?,
manual_count=?,failed_count=?,error_message=?,finished_at=? WHERE batch_id=?`, status,
counts.Processed, counts.Success, counts.Reused, counts.Manual, counts.Failed,
nullableText(errorMessage), finishedAt, batchID)
if err != nil {
return fmt.Errorf("完成 AI 匹配批次失败: %w", err)
}
return nil
}
func UpdateAIMatchBatchCounts(q Execer, batchID string, counts AIMatchBatchCounts) error {
_, err := q.Exec(`UPDATE ai_match_batches SET processed_count=?,success_count=?,reused_count=?,
manual_count=?,failed_count=? WHERE batch_id=? AND status='running'`, counts.Processed,
counts.Success, counts.Reused, counts.Manual, counts.Failed, batchID)
if err != nil {
return fmt.Errorf("更新 AI 匹配批次进度失败: %w", err)
}
return nil
}
func GetAIMatchBatch(q Execer, batchID string) (*model.AIMatchBatch, error) {
var batch model.AIMatchBatch
var errorMessage, startedAt, finishedAt sql.NullString
err := q.QueryRow(`SELECT batch_id,status,created_by_user_id,total_count,processed_count,
success_count,reused_count,manual_count,failed_count,provider_id,provider_name,
provider_base_url,model,timeout_seconds,max_concurrency,confidence_threshold_bps,
config_fingerprint,rules_version,prompt_version,error_message,created_at,started_at,finished_at
FROM ai_match_batches WHERE batch_id=?`, batchID).Scan(&batch.BatchID, &batch.Status,
&batch.CreatedByUserID, &batch.TotalCount, &batch.ProcessedCount, &batch.SuccessCount,
&batch.ReusedCount, &batch.ManualCount, &batch.FailedCount, &batch.ProviderID,
&batch.ProviderName, &batch.ProviderBaseURL, &batch.Model, &batch.TimeoutSeconds,
&batch.MaxConcurrency, &batch.ConfidenceThresholdBPS, &batch.ConfigFingerprint,
&batch.RulesVersion, &batch.PromptVersion, &errorMessage, &batch.CreatedAt, &startedAt, &finishedAt)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, fmt.Errorf("读取 AI 匹配批次失败: %w", err)
}
batch.ErrorMessage, batch.StartedAt, batch.FinishedAt = errorMessage.String, startedAt.String, finishedAt.String
return &batch, nil
}
func GetLatestAIMatchBatchID(q Execer, userID string, admin bool) (string, error) {
query := `SELECT batch_id FROM ai_match_batches`
var args []any
if !admin {
query += ` WHERE created_by_user_id=?`
args = append(args, userID)
}
query += ` ORDER BY created_at DESC,batch_id DESC LIMIT 1`
var batchID string
err := q.QueryRow(query, args...).Scan(&batchID)
if errors.Is(err, sql.ErrNoRows) {
return "", nil
}
if err != nil {
return "", fmt.Errorf("读取最近 AI 匹配批次失败: %w", err)
}
return batchID, nil
}
func ListUnfinishedAIMatchBatchIDs(q Execer) ([]string, error) {
rows, err := q.Query(`SELECT batch_id FROM ai_match_batches WHERE status IN ('queued','running') ORDER BY created_at,batch_id`)
if err != nil {
return nil, err
}
defer rows.Close()
var result []string
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
return nil, err
}
result = append(result, id)
}
return result, rows.Err()
}
func InterruptAIMatchBatch(q Execer, batchID, finishedAt string) error {
if _, err := q.Exec(`UPDATE ai_match_batch_items SET status='interrupted',outcome='failed',
message='Admin 在批次完成前退出,可重新勾选未成功条目安全重试',finished_at=?
WHERE batch_id=? AND status IN ('queued','running')`, finishedAt, batchID); err != nil {
return err
}
counts, err := CountAIMatchBatchItems(q, batchID)
if err != nil {
return err
}
return FinishAIMatchBatch(q, batchID, "interrupted", "Admin 在批次完成前退出", finishedAt, counts)
}
+98 -2
View File
@@ -20,7 +20,7 @@ import (
"cmautobuy/admin/spec"
)
const mysqlSchemaVersion = 21
const mysqlSchemaVersion = 22
// OpenMySQL 打开生产 MySQL 8 数据库。错误信息绝不包含完整 DSN 或密码。
func OpenMySQL(cfg config.DatabaseConfig) (*sql.DB, error) {
@@ -676,10 +676,75 @@ func MigrateMySQL(db *sql.DB) error {
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 21, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
return fmt.Errorf("记录 MySQL schema v21 失败: %w", err)
}
current = 21
}
if current < 22 {
if err := migrateMySQLV22(db); err != nil {
return fmt.Errorf("执行 MySQL schema v22 失败: %w", err)
}
if err := checkMySQLV22Shape(db); err != nil {
return fmt.Errorf("MySQL schema v22 自检失败,未记录版本: %w", err)
}
if _, err := db.Exec(`INSERT INTO schema_migrations (version, applied_at) VALUES (?, ?)`, 22, time.Now().UTC().Format(time.RFC3339Nano)); err != nil {
return fmt.Errorf("记录 MySQL schema v22 失败: %w", err)
}
}
return CheckMySQLSchema(db)
}
func migrateMySQLV22(db *sql.DB) error {
statements := []string{
`CREATE TABLE IF NOT EXISTS ai_match_batches (
batch_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
status VARCHAR(16) NOT NULL,
created_by_user_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
total_count INT NOT NULL,processed_count INT NOT NULL DEFAULT 0,
success_count INT NOT NULL DEFAULT 0,reused_count INT NOT NULL DEFAULT 0,
manual_count INT NOT NULL DEFAULT 0,failed_count INT NOT NULL DEFAULT 0,
provider_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
provider_name VARCHAR(191) NOT NULL,provider_base_url VARCHAR(2048) NOT NULL,
model VARCHAR(191) NOT NULL,timeout_seconds INT NOT NULL,max_concurrency INT NOT NULL,
confidence_threshold_bps INT NOT NULL,config_fingerprint CHAR(64) COLLATE utf8mb4_bin NOT NULL,
rules_version VARCHAR(32) NOT NULL,prompt_version VARCHAR(32) NOT NULL,
error_message VARCHAR(500),created_at VARCHAR(35) NOT NULL,
started_at VARCHAR(35),finished_at VARCHAR(35),
KEY idx_ai_match_batch_user (created_by_user_id,created_at DESC,batch_id),
KEY idx_ai_match_batch_status (status,created_at,batch_id),
CONSTRAINT fk_ai_match_batch_user FOREIGN KEY (created_by_user_id) REFERENCES users(user_id),
CONSTRAINT chk_ai_match_batch_status CHECK (status IN ('queued','running','partial','succeeded','failed','interrupted')),
CONSTRAINT chk_ai_match_batch_total CHECK (total_count BETWEEN 1 AND 100),
CONSTRAINT chk_ai_match_batch_counts CHECK (processed_count BETWEEN 0 AND total_count AND success_count>=0 AND reused_count>=0 AND manual_count>=0 AND failed_count>=0),
CONSTRAINT chk_ai_match_batch_runtime CHECK (timeout_seconds BETWEEN 1 AND 120 AND max_concurrency BETWEEN 1 AND 16 AND confidence_threshold_bps BETWEEN 0 AND 10000)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`,
`CREATE TABLE IF NOT EXISTS ai_match_batch_items (
item_id VARCHAR(191) COLLATE utf8mb4_bin PRIMARY KEY,
batch_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
syb_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
identity_hash CHAR(64) COLLATE utf8mb4_bin NOT NULL,
leader_syb_id VARCHAR(191) COLLATE utf8mb4_bin NOT NULL,
context_version CHAR(64) COLLATE utf8mb4_bin NOT NULL,
position INT NOT NULL,status VARCHAR(16) NOT NULL,outcome VARCHAR(32),
message VARCHAR(500),mapping_source VARCHAR(16),option_key VARCHAR(191) COLLATE utf8mb4_bin,
confidence_bps INT,model_called TINYINT NOT NULL DEFAULT 0,
started_at VARCHAR(35),finished_at VARCHAR(35),
UNIQUE KEY uq_ai_match_batch_syb (batch_id,syb_id),
KEY idx_ai_match_batch_work (batch_id,status,leader_syb_id,position),
KEY idx_ai_match_batch_identity (batch_id,identity_hash,position),
CONSTRAINT fk_ai_match_item_batch FOREIGN KEY (batch_id) REFERENCES ai_match_batches(batch_id) ON DELETE CASCADE,
CONSTRAINT chk_ai_match_item_status CHECK (status IN ('queued','running','succeeded','reused','manual','failed','stale','interrupted')),
CONSTRAINT chk_ai_match_item_source CHECK (mapping_source IS NULL OR mapping_source IN ('manual','rule','ai')),
CONSTRAINT chk_ai_match_item_confidence CHECK (confidence_bps IS NULL OR confidence_bps BETWEEN 0 AND 10000),
CONSTRAINT chk_ai_match_item_called CHECK (model_called IN (0,1))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci`,
}
for _, statement := range statements {
if _, err := db.Exec(statement); err != nil {
return err
}
}
return nil
}
// migrateMySQLV21 给当前有效规格映射增加来源,并建立只追加的 AI 决策审计。
func migrateMySQLV21(db *sql.DB) error {
columns := []struct{ name, ddl string }{
@@ -1944,6 +2009,7 @@ func CheckMySQLSchema(db *sql.DB) error {
"syb_allowed_shops",
"ai_provider_configs", "ai_provider_audits",
"ai_spec_match_decisions",
"ai_match_batches", "ai_match_batch_items",
}
if err := checkMySQLSchema(db, mysqlRequiredTables); err != nil {
return err
@@ -2002,7 +2068,37 @@ func CheckMySQLSchema(db *sql.DB) error {
if err := checkMySQLV20Shape(db); err != nil {
return err
}
return checkMySQLV21Shape(db)
if err := checkMySQLV21Shape(db); err != nil {
return err
}
return checkMySQLV22Shape(db)
}
func checkMySQLV22Shape(db *sql.DB) error {
if err := checkMySQLSchema(db, []string{"ai_match_batches", "ai_match_batch_items"}); err != nil {
return err
}
for _, item := range []struct{ kind, table, name string }{
{"constraint", "ai_match_batches", "chk_ai_match_batch_status"},
{"constraint", "ai_match_batches", "chk_ai_match_batch_counts"},
{"constraint", "ai_match_batches", "fk_ai_match_batch_user"},
{"constraint", "ai_match_batch_items", "chk_ai_match_item_status"},
{"constraint", "ai_match_batch_items", "fk_ai_match_item_batch"},
{"index", "ai_match_batch_items", "uq_ai_match_batch_syb"},
{"index", "ai_match_batch_items", "idx_ai_match_batch_identity"},
} {
var exists bool
var err error
if item.kind == "index" {
exists, err = mysqlIndexExists(db, item.table, item.name)
} else {
exists, err = mysqlConstraintExists(db, item.table, item.name)
}
if err != nil || !exists {
return fmt.Errorf("AI 匹配批次%s %s.%s 缺失: %v", item.kind, item.table, item.name, err)
}
}
return nil
}
func checkMySQLV21Shape(db *sql.DB) error {
+2
View File
@@ -590,6 +590,7 @@ type SybOrderContext struct {
MappingProviderID string
MappingModel string
MappingConfidenceBPS int
MappingConfidenceSet bool
MappingReason string
MappingSourceVersion string
MappingContextVersion string
@@ -639,6 +640,7 @@ func scanSybOrderContext(s rowScanner) (SybOrderContext, error) {
c.MappingProviderID = mappingProviderID.String
c.MappingModel = mappingModel.String
c.MappingConfidenceBPS = int(mappingConfidence.Int64)
c.MappingConfidenceSet = mappingConfidence.Valid
c.MappingReason = mappingReason.String
c.MappingSourceVersion = mappingSourceVersion.String
c.MappingContextVersion = mappingContextVersion.String