- Device struct 加 4 欄(agent_id/agent_local_device_id/registered_at/
is_representative)+ deviceColumns 13→17 + scanDevice/SaveTx 讀寫新欄
- 新增 internal/agent package(domain + interface + in-memory + PG repo):
GetOrCreateAgentTx/GetByOwnerTx,advisory lock 序列化同 owner get-or-create
- exchange 重塑:建/復用 agent → representative device(綁 session_tokens、
serial=NULL)→ loop 建 N 顆真 USB device(R1 完整 N 顆非只第一顆)
- List filter is_representative=false + DeviceListItem 回傳 agent_id/registered_at
- 併入 WP-0/0005 follow-up Minor:Mi#2 lost-update 收斂(tx 內查詢+局部更新)
/ Mi#3 過時註解 / Mi#4 空 serial 回 ErrNotFound / Mi#5 serial 白名單
^0x[0-9A-Fa-f]{8}$ + 去重 / S-1 firmware forward-compat / S-2 device Name 衍生
守 ADR-018 A'(一 owner N agents、session_tokens FK 物理不動、不加 owner
unique 為多機器留路)。Reviewer 通過(0C/0M)。5 套件 dbtest 130 全綠
(db 19/device 37/agent 13/api 172/cmd 60)、build/vet/test 綠。
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
195 lines
6.9 KiB
Go
195 lines
6.9 KiB
Go
// Package agent 的 Postgres 持久層實作(migration 0005,WP-B B2)。
|
||
//
|
||
// PostgresRepository 實作與 InMemoryRepository 相同的 Repository interface。
|
||
//
|
||
// 對齊 migrations/0005_create_agents.up.sql 的 agents 表;語意對齊 in-memory(agent.go):
|
||
// - GetByOwnerTx / GetOrCreateAgentTx 略過 deleted_at IS NOT NULL 的紀錄。
|
||
// - GetOrCreateAgentTx 復用路徑「tx 內局部 UPDATE」(只改 last_paired_at + 非空上報欄),
|
||
// 不全欄覆寫,避免 lost-update(WP-0 Mi#2 同精神)。
|
||
package agent
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"time"
|
||
|
||
"github.com/jackc/pgx/v5"
|
||
"github.com/jackc/pgx/v5/pgxpool"
|
||
|
||
"visiona-backend/internal/db"
|
||
)
|
||
|
||
// PostgresRepository 是 Agent 的 PostgreSQL 持久層實作。
|
||
type PostgresRepository struct {
|
||
pool *pgxpool.Pool
|
||
}
|
||
|
||
// NewPostgresRepository 建立一個以 pgxpool 為後端的 Repository。
|
||
func NewPostgresRepository(pool *pgxpool.Pool) *PostgresRepository {
|
||
return &PostgresRepository{pool: pool}
|
||
}
|
||
|
||
// 編譯時檢查:確保 PostgresRepository 實作 Repository。
|
||
var _ Repository = (*PostgresRepository)(nil)
|
||
|
||
// agentColumns 是 SELECT 共用欄位清單(順序須與 scanAgent 對齊)。
|
||
const agentColumns = `id, owner_user_id, name, platform, agent_version,
|
||
last_paired_at, created_at, updated_at, deleted_at`
|
||
|
||
// GetByOwnerTx 取得該 owner 的第一個未刪除 agent;不存在回 ErrNotFound。
|
||
//
|
||
// 一 owner 一 agent 語意下最多一筆 active;仍加 ORDER BY created_at + LIMIT 1 確保
|
||
// 多筆殘留時取最早那筆為決定性結果(避免非決定性 row)。
|
||
func (r *PostgresRepository) GetByOwnerTx(ctx context.Context, q db.Querier, ownerUserID string) (*Agent, error) {
|
||
const sql = `SELECT ` + agentColumns + `
|
||
FROM agents
|
||
WHERE owner_user_id = $1 AND deleted_at IS NULL
|
||
ORDER BY created_at ASC
|
||
LIMIT 1`
|
||
|
||
row := q.QueryRow(ctx, sql, ownerUserID)
|
||
a, err := scanAgent(row)
|
||
if errors.Is(err, pgx.ErrNoRows) {
|
||
return nil, ErrNotFound
|
||
}
|
||
if err != nil {
|
||
return nil, fmt.Errorf("agent: pg GetByOwner: %w", err)
|
||
}
|
||
return a, nil
|
||
}
|
||
|
||
// GetOrCreateAgentTx 復用或新建該 owner 的 agent(在傳入 Querier / tx 上執行)。
|
||
//
|
||
// 復用路徑:局部 UPDATE last_paired_at + 非空上報欄(COALESCE 保留既有值),不全欄覆寫
|
||
// ——避免 lost-update。新建路徑:INSERT,id 由 DB gen_random_uuid() 產、created_at/updated_at
|
||
// 用 DEFAULT。
|
||
//
|
||
// 併發正確性(無 owner unique 約束下防 insert race):
|
||
//
|
||
// agents 表沒有 owner_user_id 的 unique 約束(migration 0005 只有 id PK),故「空表併發
|
||
// get-or-create」若只靠 SELECT FOR UPDATE 會失效——FOR UPDATE 只能鎖「已存在的列」,空表
|
||
// 時多個並發 tx 都讀到 no rows、全走 INSERT,建出多筆 agent。
|
||
// 因此改用 transaction-scoped advisory lock(pg_advisory_xact_lock)以 (ownerUserID) 為鍵
|
||
// 序列化同 owner 的 get-or-create:同一時間只有一個 tx 能進入「查→建」臨界區,交易結束自動
|
||
// 釋放。此法不需動 schema(守 WP-B「不碰 migration」邊界)。
|
||
// ⚠️ advisory_xact_lock 需在交易內才會自動釋放——exchange 一律以 db.WithTx 包住
|
||
// (見 pairing_exchange.go)。若 q 為 pool(非 tx),lock 在該單一語句結束即釋放、
|
||
// 無法涵蓋整個「查→建」,故呼叫端務必在 tx 內使用(測試亦如是)。
|
||
func (r *PostgresRepository) GetOrCreateAgentTx(
|
||
ctx context.Context, q db.Querier, ownerUserID, name, platform, agentVersion string, pairedAt time.Time,
|
||
) (*Agent, error) {
|
||
if ownerUserID == "" {
|
||
return nil, errors.New("agent: GetOrCreateAgentTx requires ownerUserID")
|
||
}
|
||
|
||
// 0) advisory lock 序列化同 owner 的 get-or-create(tx-scoped,交易結束自動釋放)。
|
||
// 以 owner UUID 文字的 hashtext 當鎖鍵(碰撞僅造成不同 owner 偶爾序列化,不影響正確性)。
|
||
if _, lErr := q.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtext($1))`, ownerUserID); lErr != nil {
|
||
return nil, fmt.Errorf("agent: pg GetOrCreate advisory lock: %w", lErr)
|
||
}
|
||
|
||
// 1) 嘗試取既有 active agent(advisory lock 已序列化,此處 FOR UPDATE 進一步鎖既有列)。
|
||
const selForUpdate = `SELECT ` + agentColumns + `
|
||
FROM agents
|
||
WHERE owner_user_id = $1 AND deleted_at IS NULL
|
||
ORDER BY created_at ASC
|
||
LIMIT 1
|
||
FOR UPDATE`
|
||
|
||
existing, err := scanAgent(q.QueryRow(ctx, selForUpdate, ownerUserID))
|
||
switch {
|
||
case err == nil:
|
||
// 2a) 復用:tx 內局部更新(只改 last_paired_at + 非空上報欄;空值保留既有)。
|
||
const upd = `UPDATE agents SET
|
||
last_paired_at = $2,
|
||
name = COALESCE(NULLIF($3, ''), name),
|
||
platform = COALESCE(NULLIF($4, ''), platform),
|
||
agent_version = COALESCE(NULLIF($5, ''), agent_version),
|
||
updated_at = now()
|
||
WHERE id = $1
|
||
RETURNING ` + agentColumns
|
||
|
||
updated, uErr := scanAgent(q.QueryRow(ctx, upd,
|
||
existing.ID, pairedAt.UTC(), name, platform, agentVersion))
|
||
if uErr != nil {
|
||
return nil, fmt.Errorf("agent: pg GetOrCreate update: %w", uErr)
|
||
}
|
||
return updated, nil
|
||
|
||
case errors.Is(err, pgx.ErrNoRows):
|
||
// 2b) 新建:id / created_at / updated_at 由 DB 產;name 空時走 agents.name DEFAULT。
|
||
const ins = `INSERT INTO agents
|
||
(owner_user_id, name, platform, agent_version, last_paired_at)
|
||
VALUES ($1, COALESCE(NULLIF($2, ''), 'local-agent'), NULLIF($3, ''), NULLIF($4, ''), $5)
|
||
RETURNING ` + agentColumns
|
||
|
||
created, iErr := scanAgent(q.QueryRow(ctx, ins,
|
||
ownerUserID, name, platform, agentVersion, pairedAt.UTC()))
|
||
if iErr != nil {
|
||
return nil, fmt.Errorf("agent: pg GetOrCreate insert: %w", iErr)
|
||
}
|
||
return created, nil
|
||
|
||
default:
|
||
return nil, fmt.Errorf("agent: pg GetOrCreate select: %w", err)
|
||
}
|
||
}
|
||
|
||
// ==========================================================================
|
||
// scan helper
|
||
// ==========================================================================
|
||
|
||
type rowScanner interface {
|
||
Scan(dest ...any) error
|
||
}
|
||
|
||
// scanAgent 從一列掃出 *Agent。欄位順序須與 agentColumns 對齊。
|
||
//
|
||
// nullable TEXT(platform / agent_version)NULL → 空字串(對齊 in-memory zero value);
|
||
// nullable TIMESTAMPTZ(last_paired_at / deleted_at)以 *time.Time 接,NULL → nil。
|
||
func scanAgent(row rowScanner) (*Agent, error) {
|
||
var (
|
||
a Agent
|
||
platform *string
|
||
agentVersion *string
|
||
)
|
||
err := row.Scan(
|
||
&a.ID,
|
||
&a.OwnerUserID,
|
||
&a.Name,
|
||
&platform,
|
||
&agentVersion,
|
||
&a.LastPairedAt,
|
||
&a.CreatedAt,
|
||
&a.UpdatedAt,
|
||
&a.DeletedAt,
|
||
)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
a.Platform = derefString(platform)
|
||
a.AgentVersion = derefString(agentVersion)
|
||
|
||
a.CreatedAt = a.CreatedAt.UTC()
|
||
a.UpdatedAt = a.UpdatedAt.UTC()
|
||
if a.LastPairedAt != nil {
|
||
t := a.LastPairedAt.UTC()
|
||
a.LastPairedAt = &t
|
||
}
|
||
if a.DeletedAt != nil {
|
||
t := a.DeletedAt.UTC()
|
||
a.DeletedAt = &t
|
||
}
|
||
return &a, nil
|
||
}
|
||
|
||
// derefString 解指標字串,nil 視為空字串(對齊 in-memory zero value)。
|
||
func derefString(s *string) string {
|
||
if s == nil {
|
||
return ""
|
||
}
|
||
return *s
|
||
}
|