feat: 增加客户端采购员归属管理 (#54)
This commit is contained in:
@@ -2,12 +2,26 @@ package repository
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"cmautobuy/admin/model"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrClientNotFound = errors.New("客户端不存在")
|
||||
ErrPurchaserNotActive = errors.New("采购员不存在或已禁用")
|
||||
ErrClientNotAssigned = errors.New("客户端当前未绑定采购员")
|
||||
)
|
||||
|
||||
// ClientWithAssignee 是客户端列表查询结果。归属为空表示尚未绑定。
|
||||
type ClientWithAssignee struct {
|
||||
model.Client
|
||||
AssignedUserID string
|
||||
AssignedUsername string
|
||||
}
|
||||
|
||||
// UpsertClient 登记或更新一台客户端。
|
||||
//
|
||||
// # 名称的更新规则(两个接口不一样,这是有意的)
|
||||
@@ -112,6 +126,153 @@ func ListClients(db *sql.DB, keyword string) ([]model.Client, error) {
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// ListClientsForUser 返回带当前负责人的客户端。
|
||||
// visibleUserID 为空表示管理员全量视图,否则只返回当前绑定给该用户的客户端。
|
||||
func ListClientsForUser(db *sql.DB, keyword, visibleUserID string) ([]ClientWithAssignee, error) {
|
||||
query := `SELECT c.client_id, c.name, c.device_address, c.platform, c.pdd_package,
|
||||
c.capabilities, c.last_seen_at, c.created_at, c.updated_at,
|
||||
a.user_id, u.username
|
||||
FROM clients c
|
||||
LEFT JOIN client_user_assignments a
|
||||
ON a.client_id = c.client_id AND a.ended_at IS NULL
|
||||
LEFT JOIN users u ON u.user_id = a.user_id`
|
||||
where := make([]string, 0, 2)
|
||||
args := make([]any, 0, 3)
|
||||
if visibleUserID != "" {
|
||||
where = append(where, `a.user_id = ?`)
|
||||
args = append(args, visibleUserID)
|
||||
}
|
||||
if kw := strings.TrimSpace(keyword); kw != "" {
|
||||
where = append(where, `(c.name LIKE ? OR c.client_id LIKE ?)`)
|
||||
like := "%" + kw + "%"
|
||||
args = append(args, like, like)
|
||||
}
|
||||
if len(where) > 0 {
|
||||
query += ` WHERE ` + strings.Join(where, ` AND `)
|
||||
}
|
||||
query += ` ORDER BY c.last_seen_at DESC, c.client_id`
|
||||
|
||||
rows, err := db.Query(query, args...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("查询客户端归属列表失败: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := make([]ClientWithAssignee, 0)
|
||||
for rows.Next() {
|
||||
var row ClientWithAssignee
|
||||
var name, addr, platform, pkg, caps, userID, username sql.NullString
|
||||
if err := rows.Scan(&row.ClientID, &name, &addr, &platform, &pkg, &caps,
|
||||
&row.LastSeenAt, &row.CreatedAt, &row.UpdatedAt, &userID, &username); err != nil {
|
||||
return nil, fmt.Errorf("读取客户端归属行失败: %w", err)
|
||||
}
|
||||
row.Name, row.DeviceAddress, row.Platform = name.String, addr.String, platform.String
|
||||
row.PddPackage, row.Capabilities = pkg.String, caps.String
|
||||
row.AssignedUserID, row.AssignedUsername = userID.String, username.String
|
||||
out = append(out, row)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// ListActivePurchasers 返回可成为新负责人的启用采购员。
|
||||
func ListActivePurchasers(db *sql.DB) ([]model.User, error) {
|
||||
rows, err := db.Query(`
|
||||
SELECT user_id, username, role, status, password_changed_at, created_at, updated_at
|
||||
FROM users WHERE role = ? AND status = ? ORDER BY username`,
|
||||
model.RolePurchaser, model.UserActive)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("查询可绑定采购员失败: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
users := make([]model.User, 0)
|
||||
for rows.Next() {
|
||||
var user model.User
|
||||
if err := rows.Scan(&user.UserID, &user.Username, &user.Role, &user.Status,
|
||||
&user.PasswordChangedAt, &user.CreatedAt, &user.UpdatedAt); err != nil {
|
||||
return nil, fmt.Errorf("读取可绑定采购员失败: %w", err)
|
||||
}
|
||||
users = append(users, user)
|
||||
}
|
||||
return users, rows.Err()
|
||||
}
|
||||
|
||||
// AssignClient 原子完成首次绑定或转交。返回 changed、transferred。
|
||||
func AssignClient(db *sql.DB, assignment model.ClientUserAssignment) (bool, bool, error) {
|
||||
tx, err := db.Begin()
|
||||
if err != nil {
|
||||
return false, false, fmt.Errorf("开始绑定客户端事务失败: %w", err)
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var exists int
|
||||
if err := tx.QueryRow(`SELECT COUNT(*) FROM clients WHERE client_id = ?`, assignment.ClientID).Scan(&exists); err != nil {
|
||||
return false, false, fmt.Errorf("检查客户端失败: %w", err)
|
||||
}
|
||||
if exists == 0 {
|
||||
return false, false, ErrClientNotFound
|
||||
}
|
||||
if err := tx.QueryRow(`SELECT COUNT(*) FROM users WHERE user_id = ? AND role = ? AND status = ?`,
|
||||
assignment.UserID, model.RolePurchaser, model.UserActive).Scan(&exists); err != nil {
|
||||
return false, false, fmt.Errorf("检查采购员失败: %w", err)
|
||||
}
|
||||
if exists == 0 {
|
||||
return false, false, ErrPurchaserNotActive
|
||||
}
|
||||
|
||||
var currentID, currentUserID string
|
||||
err = tx.QueryRow(`SELECT assignment_id, user_id FROM client_user_assignments
|
||||
WHERE client_id = ? AND ended_at IS NULL`, assignment.ClientID).Scan(¤tID, ¤tUserID)
|
||||
if err != nil && !errors.Is(err, sql.ErrNoRows) {
|
||||
return false, false, fmt.Errorf("读取当前客户端归属失败: %w", err)
|
||||
}
|
||||
if err == nil && currentUserID == assignment.UserID {
|
||||
return false, false, nil
|
||||
}
|
||||
transferred := err == nil
|
||||
if transferred {
|
||||
if _, err := tx.Exec(`UPDATE client_user_assignments
|
||||
SET ended_at = ?, ended_by_user_id = ?, end_reason = 'transfer'
|
||||
WHERE assignment_id = ? AND ended_at IS NULL`,
|
||||
assignment.StartedAt, assignment.AssignedByUserID, currentID); err != nil {
|
||||
return false, false, fmt.Errorf("结束原客户端归属失败: %w", err)
|
||||
}
|
||||
}
|
||||
if _, err := tx.Exec(`INSERT INTO client_user_assignments
|
||||
(assignment_id, client_id, user_id, started_at, assigned_by_user_id)
|
||||
VALUES (?, ?, ?, ?, ?)`, assignment.AssignmentID, assignment.ClientID,
|
||||
assignment.UserID, assignment.StartedAt, assignment.AssignedByUserID); err != nil {
|
||||
return false, false, fmt.Errorf("保存客户端归属失败: %w", err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return false, false, fmt.Errorf("提交客户端绑定事务失败: %w", err)
|
||||
}
|
||||
return true, transferred, nil
|
||||
}
|
||||
|
||||
// UnassignClient 原子结束当前归属,历史记录保留。
|
||||
func UnassignClient(db *sql.DB, clientID, actorUserID, endedAt string) error {
|
||||
tx, err := db.Begin()
|
||||
if err != nil {
|
||||
return fmt.Errorf("开始解绑客户端事务失败: %w", err)
|
||||
}
|
||||
defer tx.Rollback()
|
||||
res, err := tx.Exec(`UPDATE client_user_assignments
|
||||
SET ended_at = ?, ended_by_user_id = ?, end_reason = 'unbind'
|
||||
WHERE client_id = ? AND ended_at IS NULL`, endedAt, actorUserID, clientID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("结束客户端归属失败: %w", err)
|
||||
}
|
||||
n, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return fmt.Errorf("读取解绑结果失败: %w", err)
|
||||
}
|
||||
if n == 0 {
|
||||
return ErrClientNotAssigned
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("提交客户端解绑事务失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteClients 批量删除客户端。返回实际删除的条数。
|
||||
func DeleteClients(db *sql.DB, clientIDs []string) (int64, error) {
|
||||
if len(clientIDs) == 0 {
|
||||
|
||||
+37
-2
@@ -264,7 +264,7 @@ var migrations = [][]string{
|
||||
// 背景见 #20:v1 曾经被原地改写而不是新增版本,导致已经建过库的机器
|
||||
// (user_version 已经越过 v1)永远不会重跑改写后的语句,程序拿着一个
|
||||
// 和代码对不上的库静默启动。
|
||||
const schemaVersion = 6
|
||||
const schemaVersion = 7
|
||||
|
||||
// migrationV4 给 PDD 商品增加店铺名。
|
||||
//
|
||||
@@ -326,6 +326,33 @@ var migrationV6 = []string{
|
||||
`CREATE INDEX idx_web_sessions_expiry ON web_sessions(expires_at);`,
|
||||
}
|
||||
|
||||
// migrationV7 记录客户端当前负责人及完整转交历史,见工单 #54。
|
||||
// client_id 故意不加 clients 外键:客户端记录可删除后由同一稳定编号重新登记,
|
||||
// 归属和审计历史不能因此丢失。
|
||||
var migrationV7 = []string{
|
||||
`CREATE TABLE client_user_assignments (
|
||||
assignment_id TEXT PRIMARY KEY,
|
||||
client_id TEXT NOT NULL,
|
||||
user_id TEXT NOT NULL,
|
||||
started_at TEXT NOT NULL,
|
||||
ended_at TEXT,
|
||||
assigned_by_user_id TEXT NOT NULL,
|
||||
ended_by_user_id TEXT,
|
||||
end_reason TEXT CHECK (end_reason IS NULL OR end_reason IN ('unbind', 'transfer')),
|
||||
FOREIGN KEY (user_id) REFERENCES users(user_id),
|
||||
FOREIGN KEY (assigned_by_user_id) REFERENCES users(user_id),
|
||||
FOREIGN KEY (ended_by_user_id) REFERENCES users(user_id),
|
||||
CHECK ((ended_at IS NULL AND ended_by_user_id IS NULL AND end_reason IS NULL)
|
||||
OR (ended_at IS NOT NULL AND ended_by_user_id IS NOT NULL AND end_reason IS NOT NULL))
|
||||
);`,
|
||||
`CREATE UNIQUE INDEX idx_client_assignment_current
|
||||
ON client_user_assignments(client_id) WHERE ended_at IS NULL;`,
|
||||
`CREATE INDEX idx_client_assignment_user
|
||||
ON client_user_assignments(user_id, ended_at, client_id);`,
|
||||
`CREATE INDEX idx_client_assignment_history
|
||||
ON client_user_assignments(client_id, started_at DESC);`,
|
||||
}
|
||||
|
||||
// Migrate 把数据库升到最新版本。
|
||||
// 已经是最新的就什么都不做,可以重复调用。
|
||||
func Migrate(db *sql.DB) error {
|
||||
@@ -396,6 +423,14 @@ func Migrate(db *sql.DB) error {
|
||||
if err := runSQLMigration(db, 6, migrationV6); err != nil {
|
||||
return err
|
||||
}
|
||||
reached = 6
|
||||
}
|
||||
|
||||
// v7 是纯追加的客户端归属历史表和索引。
|
||||
if reached < 7 {
|
||||
if err := runSQLMigration(db, 7, migrationV7); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
@@ -857,7 +892,7 @@ var requiredTables = []string{
|
||||
"syb_orders", "sku_mappings", "tasks", "clients",
|
||||
"idempotency_keys", "task_claims",
|
||||
"syb_session", "syb_sync_state",
|
||||
"users", "web_sessions",
|
||||
"users", "web_sessions", "client_user_assignments",
|
||||
}
|
||||
|
||||
// requiredColumns 只列出不能靠“表存在”发现的关键追加列。
|
||||
|
||||
@@ -593,7 +593,7 @@ func TestMigrate_v5新增会话表同步状态表和product_spec列(t *testing.T
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrate_v6新增用户和WebSession表(t *testing.T) {
|
||||
func TestMigrate_v6用户表与v7客户端归属表均存在(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
name string
|
||||
db *sql.DB
|
||||
@@ -608,7 +608,7 @@ func TestMigrate_v6新增用户和WebSession表(t *testing.T) {
|
||||
}
|
||||
}
|
||||
tables := existingTableSet(t, c.db)
|
||||
for _, table := range []string{"users", "web_sessions"} {
|
||||
for _, table := range []string{"users", "web_sessions", "client_user_assignments"} {
|
||||
if !tables[table] {
|
||||
t.Errorf("%s:迁移后应该有表 %s", c.name, table)
|
||||
}
|
||||
@@ -617,21 +617,21 @@ func TestMigrate_v6新增用户和WebSession表(t *testing.T) {
|
||||
if err := c.db.QueryRow("PRAGMA user_version").Scan(&version); err != nil {
|
||||
t.Fatalf("%s 读取 user_version 失败: %v", c.name, err)
|
||||
}
|
||||
if version != 6 {
|
||||
t.Errorf("%s user_version = %d,期望 6", c.name, version)
|
||||
if version != schemaVersion {
|
||||
t.Errorf("%s user_version = %d,期望 %d", c.name, version, schemaVersion)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// #52:网页登录上线后显式覆盖每一个已发布版本起点,证明 Client 现实库
|
||||
// 不会因为 v6 的用户表而卡在中间版本。v2 的两种历史结构另由上面的收敛
|
||||
// #52/#54:显式覆盖每一个已发布版本起点,证明 Client 现实库
|
||||
// 不会因为 v6/v7 的新增表而卡在中间版本。v2 的两种历史结构另由上面的收敛
|
||||
// 测试持续覆盖;这里验证顺序发布的 v1-v5 路径。
|
||||
func TestMigrate_v1到v5均可升级到v6(t *testing.T) {
|
||||
func TestMigrate_v1到v5均可升级到最新版本(t *testing.T) {
|
||||
for version := 1; version <= 5; version++ {
|
||||
t.Run(fmt.Sprintf("v%d", version), func(t *testing.T) {
|
||||
db := newPublishedVersionDB(t, version)
|
||||
if err := Migrate(db); err != nil {
|
||||
t.Fatalf("v%d 迁移到 v6 失败: %v", version, err)
|
||||
t.Fatalf("v%d 迁移到最新版本失败: %v", version, err)
|
||||
}
|
||||
if err := CheckSchema(db); err != nil {
|
||||
t.Fatalf("v%d 迁移后 schema 自检失败: %v", version, err)
|
||||
@@ -640,8 +640,8 @@ func TestMigrate_v1到v5均可升级到v6(t *testing.T) {
|
||||
if err := db.QueryRow(`PRAGMA user_version`).Scan(&gotVersion); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if gotVersion != 6 {
|
||||
t.Fatalf("v%d 迁移后 user_version = %d,期望 6", version, gotVersion)
|
||||
if gotVersion != schemaVersion {
|
||||
t.Fatalf("v%d 迁移后 user_version = %d,期望 %d", version, gotVersion, schemaVersion)
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -1120,7 +1120,7 @@ func TestCheckSchema_缺少v4关键列时拒绝(t *testing.T) {
|
||||
if err := migrateV3(db); err != nil {
|
||||
t.Fatalf("准备 v3 数据库失败: %v", err)
|
||||
}
|
||||
// 故意跳过 v4(不加 shop_name),但把 v5/v6 补上——否则 CheckSchema 会先
|
||||
// 故意跳过 v4(不加 shop_name),但把 v5-v7 补上——否则 CheckSchema 会先
|
||||
// 因为缺后续表报错,测不到本测试真正要覆盖的"缺 shop_name"这条路径。
|
||||
if err := runSQLMigration(db, 5, migrationV5); err != nil {
|
||||
t.Fatalf("准备 v5 数据库失败: %v", err)
|
||||
@@ -1128,6 +1128,9 @@ func TestCheckSchema_缺少v4关键列时拒绝(t *testing.T) {
|
||||
if err := runSQLMigration(db, 6, migrationV6); err != nil {
|
||||
t.Fatalf("准备 v6 数据库失败: %v", err)
|
||||
}
|
||||
if err := runSQLMigration(db, 7, migrationV7); err != nil {
|
||||
t.Fatalf("准备 v7 数据库失败: %v", err)
|
||||
}
|
||||
|
||||
err := CheckSchema(db)
|
||||
if err == nil || !strings.Contains(err.Error(), "shop_name") {
|
||||
|
||||
Reference in New Issue
Block a user