visionA/visionA-backend/internal/agent/postgres_repository.go
jim800121chen 59c57fa481 feat(device): WP-B repository 接 agents 模型 + exchange 重塑(A' 走向第二階段 Go 層)
- 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>
2026-07-16 11:40:17 +08:00

195 lines
6.9 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// Package agent 的 Postgres 持久層實作migration 0005WP-B B2
//
// PostgresRepository 實作與 InMemoryRepository 相同的 Repository interface。
//
// 對齊 migrations/0005_create_agents.up.sql 的 agents 表;語意對齊 in-memoryagent.go
// - GetByOwnerTx / GetOrCreateAgentTx 略過 deleted_at IS NOT NULL 的紀錄。
// - GetOrCreateAgentTx 復用路徑「tx 內局部 UPDATE」只改 last_paired_at + 非空上報欄),
// 不全欄覆寫,避免 lost-updateWP-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。新建路徑INSERTid 由 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 lockpg_advisory_xact_lock以 (ownerUserID) 為鍵
// 序列化同 owner 的 get-or-create同一時間只有一個 tx 能進入「查→建」臨界區,交易結束自動
// 釋放。此法不需動 schema守 WP-B「不碰 migration」邊界
// ⚠️ advisory_xact_lock 需在交易內才會自動釋放——exchange 一律以 db.WithTx 包住
// (見 pairing_exchange.go。若 q 為 pool非 txlock 在該單一語句結束即釋放、
// 無法涵蓋整個「查→建」,故呼叫端務必在 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-createtx-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 agentadvisory 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 TEXTplatform / agent_versionNULL → 空字串(對齊 in-memory zero value
// nullable TIMESTAMPTZlast_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
}