// 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 }