package inventory import ( "context" "encoding/json" "errors" "fmt" "time" "github.com/jackc/pgx/v5" ) // What each machine says of its long-running resources (novox/hq ADR 0240, to-be 48 §4): the newest // statement per machine, and per module how many statements in a row said a resource of it was unhealthy. // ResourceHealth is one long-running resource's state as the node-engine said it. type ResourceHealth struct { Module string `json:"module"` Resource string `json:"resource"` Kind string `json:"kind"` Target string `json:"target"` State string `json:"state"` Reason string `json:"reason,omitempty"` Since time.Time `json:"since"` Streak int `json:"streak,omitempty"` Restarts int `json:"restarts,omitempty"` } // NodeHealth is a machine's newest statement, as kept. type NodeHealth struct { Node string Contract int // SaidAt is when the engine looked (the machine's clock); HeardAt when this controller heard it. SaidAt time.Time HeardAt time.Time Resources []ResourceHealth // Streaks is, per module, how many statements in a row said a resource of it was unhealthy. Streaks map[string]int } // HealthOf is a machine's newest statement; false when its node-engine has never stated one. func (i *Inventory) HealthOf(ctx context.Context, nodeName string) (NodeHealth, bool, error) { all, err := i.healths(ctx, nodeName) if err != nil { return NodeHealth{}, false, err } h, ok := all[nodeName] return h, ok, nil } // Healths is every machine's newest statement, by machine. A machine absent never stated one: its // node-engine is older than the judging, and its health is not known. func (i *Inventory) Healths(ctx context.Context) (map[string]NodeHealth, error) { return i.healths(ctx, "") } func (i *Inventory) healths(ctx context.Context, only string) (map[string]NodeHealth, error) { rows, err := i.store.Pool().Query(ctx, `select n.name, h.contract, h.said_at, h.heard_at, h.resources, h.streaks from node_health h join node n on n.id = h.node where $1 = '' or n.name = $1`, only) if err != nil { return nil, err } defer rows.Close() out := map[string]NodeHealth{} for rows.Next() { var h NodeHealth var resources, streaks []byte if err := rows.Scan(&h.Node, &h.Contract, &h.SaidAt, &h.HeardAt, &resources, &streaks); err != nil { return nil, err } if err := json.Unmarshal(resources, &h.Resources); err != nil { return nil, fmt.Errorf("%s's health cannot be read: %w", h.Node, err) } if err := json.Unmarshal(streaks, &h.Streaks); err != nil { return nil, fmt.Errorf("%s's health cannot be read: %w", h.Node, err) } out[h.Node] = h } return out, rows.Err() } // RecordHealth keeps a machine's statement in place of the one kept — unless the one kept is newer, by // when the engine looked: then nothing is written, and false says it was refused as older. func (i *Inventory) RecordHealth(ctx context.Context, h NodeHealth) (bool, error) { if h.Resources == nil { h.Resources = []ResourceHealth{} } if h.Streaks == nil { h.Streaks = map[string]int{} } resources, err := json.Marshal(h.Resources) if err != nil { return false, err } streaks, err := json.Marshal(h.Streaks) if err != nil { return false, err } heard := h.HeardAt if heard.IsZero() { heard = time.Now() } var node string err = i.store.Pool().QueryRow(ctx, `insert into node_health (node, contract, said_at, heard_at, resources, streaks) select id, $2, $3, $4, $5, $6 from node where name = $1 on conflict (node) do update set contract = excluded.contract, said_at = excluded.said_at, heard_at = excluded.heard_at, resources = excluded.resources, streaks = excluded.streaks where node_health.said_at <= excluded.said_at returning node`, h.Node, h.Contract, h.SaidAt, heard, resources, streaks).Scan(&node) if errors.Is(err, pgx.ErrNoRows) { if _, nerr := i.NodeByName(ctx, h.Node); nerr != nil { return false, nerr } return false, nil } return err == nil, err }