package postgres import ( "context" "encoding/json" "fmt" "time" "armature/store" ) // ClusterStore manages node_registry — Postgres-only UNLOGGED table. type ClusterStore struct{} func NewClusterStore() *ClusterStore { return &ClusterStore{} } func (s *ClusterStore) Register(ctx context.Context, nodeID, endpoint string) error { _, err := DB.ExecContext(ctx, ` INSERT INTO node_registry (node_id, endpoint) VALUES ($1, $2) ON CONFLICT (node_id) DO UPDATE SET endpoint = EXCLUDED.endpoint, registered_at = now(), heartbeat = now(), stats = '{}' `, nodeID, endpoint) return err } func (s *ClusterStore) Heartbeat(ctx context.Context, nodeID string, stats json.RawMessage) (int64, error) { result, err := DB.ExecContext(ctx, ` UPDATE node_registry SET heartbeat = now(), stats = $2 WHERE node_id = $1 `, nodeID, stats) if err != nil { return 0, err } return result.RowsAffected() } func (s *ClusterStore) SweepStale(ctx context.Context, threshold time.Duration) (int64, error) { result, err := DB.ExecContext(ctx, ` DELETE FROM node_registry WHERE heartbeat < now() - $1 * interval '1 second' `, threshold.Seconds()) if err != nil { return 0, err } return result.RowsAffected() } func (s *ClusterStore) ListNodes(ctx context.Context) ([]store.ClusterNode, error) { rows, err := DB.QueryContext(ctx, ` SELECT node_id, endpoint, seq, registered_at, heartbeat, stats FROM node_registry ORDER BY seq `) if err != nil { return nil, err } defer rows.Close() var nodes []store.ClusterNode for rows.Next() { var n store.ClusterNode if err := rows.Scan(&n.NodeID, &n.Endpoint, &n.Seq, &n.RegisteredAt, &n.Heartbeat, &n.Stats); err != nil { return nil, fmt.Errorf("scan cluster node: %w", err) } nodes = append(nodes, n) } return nodes, rows.Err() } func (s *ClusterStore) Deregister(ctx context.Context, nodeID string) error { _, err := DB.ExecContext(ctx, `DELETE FROM node_registry WHERE node_id = $1`, nodeID) return err }