package inventory import ( "context" "crypto/rand" "crypto/sha256" "encoding/base64" "encoding/hex" "encoding/json" "errors" "fmt" "strings" "time" "github.com/jackc/pgx/v5" "github.com/novox/mesh-control/internal/store" ) // Inventory is this context, holding the store it exclusively owns. type Inventory struct{ store *store.Store } // Open connects to the inventory store. func Open(ctx context.Context) (*Inventory, error) { s, err := store.Open(ctx, Name) if err != nil { return nil, err } return &Inventory{store: s}, nil } func (i *Inventory) Close() { i.store.Close() } // Ready waits for the database to answer. func (i *Inventory) Ready(ctx context.Context, within time.Duration) error { return i.store.Ready(ctx, within) } // Node is a machine the mesh knows about. type Node struct { ID string Name string Created time.Time // LastSeen is when this node was last heard from, or zero if it never has been. // // Zero and long-ago are different answers and are kept different. A node that has never // spoken has not joined properly; a node that spoke last month is a node running last // month's assignments, and until this existed those looked the same as a node that is // current (novox/hq 09-the-node-lifecycle). LastSeen time.Time } // Silent is how long since this node was last heard from, and whether it ever was. func (n Node) Silent() (time.Duration, bool) { if n.LastSeen.IsZero() { return 0, false } return time.Since(n.LastSeen), true } // ErrNoSuchNode is returned when a name matches no record. var ErrNoSuchNode = errors.New("no node of that name") // ErrNameTaken is returned when a node of that name already exists. // // Its own error rather than the driver's, because "that name is taken" is an ordinary answer a // person can act on, and a unique-violation from PostgreSQL is not. var ErrNameTaken = errors.New("a node of that name already exists") // AddNode creates a node record. // // The record comes first and the machine second: a token is issued *for* a node record // (novox/hq 09-the-node-lifecycle), so the record is what a token binds to and must exist before // there is anything to join. func (i *Inventory) AddNode(ctx context.Context, name string) (Node, error) { name = strings.TrimSpace(name) if name == "" { return Node{}, errors.New("a node needs a name: it is how a token is issued for it") } var n Node err := i.store.Pool().QueryRow(ctx, `insert into node (name) values ($1) returning id, name, created`, name).Scan(&n.ID, &n.Name, &n.Created) if err != nil { if strings.Contains(err.Error(), "node_name_key") { return Node{}, fmt.Errorf("%w: %s", ErrNameTaken, name) } return Node{}, err } return n, nil } // Nodes are every node record, oldest first. func (i *Inventory) Nodes(ctx context.Context) ([]Node, error) { rows, err := i.store.Pool().Query(ctx, `select id, name, created, last_seen from node order by created, name`) if err != nil { return nil, err } defer rows.Close() var nodes []Node for rows.Next() { var n Node var seen *time.Time if err := rows.Scan(&n.ID, &n.Name, &n.Created, &seen); err != nil { return nil, err } if seen != nil { n.LastSeen = *seen } nodes = append(nodes, n) } return nodes, rows.Err() } // NodeByName finds one node record. func (i *Inventory) NodeByName(ctx context.Context, name string) (Node, error) { var n Node err := i.store.Pool().QueryRow(ctx, `select id, name, created from node where name = $1`, name).Scan(&n.ID, &n.Name, &n.Created) if errors.Is(err, pgx.ErrNoRows) { return Node{}, fmt.Errorf("%w: %s", ErrNoSuchNode, name) } return n, err } // Issued is a token that has just been made. The secret is in it exactly once. type Issued struct { Node Node Secret string Expires time.Time } // hashSecret is what gets stored in place of the secret. // // SHA-256 rather than a password hash, and that is deliberate rather than a shortcut. bcrypt and // its relatives are slow on purpose because a password is low-entropy and guessable; this secret // is 256 bits from the system's random source, so there is nothing to guess and the slowness would // buy nothing while making every redemption expensive. func hashSecret(secret string) string { sum := sha256.Sum256([]byte(secret)) return hex.EncodeToString(sum[:]) } // IssueToken mints a one-time right to join, for a node record. // // The secret is returned once and never again. What is stored is its hash, so a copy of this // database is not a set of working credentials. // // Any outstanding token for the same node is expired first. Two live tokens for one node record // are two machines able to join as the same node, and nothing downstream could tell which was // meant — novox/hq ADR 0004's stolen-laptop case arriving before enrolment rather than after. func (i *Inventory) IssueToken(ctx context.Context, nodeName string, validFor time.Duration) (Issued, error) { if validFor <= 0 { return Issued{}, errors.New("a token needs a lifetime: one that never expires is a " + "permanent credential, which is the thing this is designed not to be") } node, err := i.NodeByName(ctx, nodeName) if err != nil { return Issued{}, err } raw := make([]byte, 32) if _, err := rand.Read(raw); err != nil { return Issued{}, fmt.Errorf("cannot generate a token secret: %w", err) } secret := base64.RawURLEncoding.EncodeToString(raw) expires := time.Now().Add(validFor) tx, err := i.store.Pool().Begin(ctx) if err != nil { return Issued{}, err } defer func() { _ = tx.Rollback(context.WithoutCancel(ctx)) }() // Expired rather than deleted: what was issued and then withdrawn is worth being able to see. if _, err := tx.Exec(ctx, `update enrolment_token set expires = now() where node = $1 and redeemed is null and expires > now()`, node.ID); err != nil { return Issued{}, err } if _, err := tx.Exec(ctx, `insert into enrolment_token (node, secret, expires) values ($1, $2, $3)`, node.ID, hashSecret(secret), expires); err != nil { return Issued{}, err } if err := tx.Commit(ctx); err != nil { return Issued{}, err } return Issued{Node: node, Secret: secret, Expires: expires}, nil } // ErrTokenRefused is what redemption returns for anything that is not a live token. // // One error for every reason — unknown, already used, expired — and deliberately so. Whoever is // presenting a token that does not work is either a machine whose operator can be told out of // band, or somebody guessing, and the second must not learn which of their guesses was a real // token that had expired. var ErrTokenRefused = errors.New("that token cannot be used") // Redeem spends a token and reports which node it was for. // // It does not issue an identity. What a node presents afterwards to prove it is that node is not // decided anywhere (novox/hq ADR 0004 names the property, not the mechanism), and guessing at it // in a migration is the most expensive guess available here. // // The update is the check: one statement that both finds a live token and marks it used, so two // simultaneous redemptions of one secret cannot both succeed. Reading first and writing second // would leave exactly that gap. func (i *Inventory) Redeem(ctx context.Context, secret string) (Node, error) { var id string err := i.store.Pool().QueryRow(ctx, `update enrolment_token set redeemed = now() where secret = $1 and redeemed is null and expires > now() returning node`, hashSecret(secret)).Scan(&id) if errors.Is(err, pgx.ErrNoRows) { return Node{}, ErrTokenRefused } if err != nil { return Node{}, err } var n Node err = i.store.Pool().QueryRow(ctx, `select id, name, created from node where id = $1`, id).Scan(&n.ID, &n.Name, &n.Created) return n, err } // RecordProfile keeps the last thing a node said about what it can do. // // The last one, not a history: the control plane needs to know what this machine can run *now* in // order to decide what it should run, and an old profile is worse than none — it describes a // machine that may have been rebuilt since. func (i *Inventory) RecordProfile(ctx context.Context, node string, profile map[string]any) error { raw, err := json.Marshal(profile) if err != nil { return err } _, err = i.store.Pool().Exec(ctx, `update node set profile = $2, last_seen = now() where id = $1`, node, raw) return err } // Seen records that a node was heard from. // // Separate from the profile because it happens far more often: a node reports it is alive // constantly and describes itself rarely. func (i *Inventory) Seen(ctx context.Context, node string) error { _, err := i.store.Pool().Exec(ctx, `update node set last_seen = now() where id = $1`, node) return err } // RecordOwned keeps the last account a node gave of what it holds. // // A copy for recovery and never a source (novox/hq 09-the-node-lifecycle). Nothing here decides // anything from it; it is handed back to a node that has lost its own store, and if that node // then disagrees, the node wins — it is the one that can see the machine. // // Replaced rather than appended. A history of what a node used to own answers a question nobody // asks, and the one question this does answer — what is on that machine now — is only answered by // the latest. func (i *Inventory) RecordOwned(ctx context.Context, node string, owned []string) error { raw, err := json.Marshal(owned) if err != nil { return err } _, err = i.store.Pool().Exec(ctx, `update node set owned = $2, owned_reported = now(), last_seen = now() where id = $1`, node, raw) return err } // Owned is what a node last said it holds, and when it said so. // // The age is returned with it rather than left to the caller to look up, because an answer about // a machine is worth much less without one — and this repository has already been bitten by a // cache with no age on it. func (i *Inventory) Owned(ctx context.Context, node string) ([]string, time.Time, error) { var raw []byte var reported *time.Time err := i.store.Pool().QueryRow(ctx, `select owned, owned_reported from node where id = $1`, node).Scan(&raw, &reported) if errors.Is(err, pgx.ErrNoRows) { return nil, time.Time{}, fmt.Errorf("%w: %s", ErrNoSuchNode, node) } if err != nil { return nil, time.Time{}, err } if len(raw) == 0 || reported == nil { // Never reported is not the same as reported nothing. A node that has applied nothing // holds nothing; a node that has never spoken is unknown, and handing back an empty list // as though it were a report would tell a rebuilding node it owns nothing and have it // remove whatever it found. return nil, time.Time{}, nil } var owned []string if err := json.Unmarshal(raw, &owned); err != nil { return nil, time.Time{}, err } return owned, *reported, nil } // Overlay is what the mesh knows about one node's place on the private network. type Overlay struct { Node string Name string Key string Endpoint string Site string Hub bool Address string } // Reachable reports whether other nodes can dial this one. // // From the endpoint alone, which is declared. Never from the shape of an address: that inference // is wrong for carrier-grade NAT, wrong for IPv6, and wrong for a routable address behind a // closed firewall (novox/hq ADR 0007). func (o Overlay) Reachable() bool { return strings.TrimSpace(o.Endpoint) != "" } // RecordSealingKey keeps the public half of the key this node's secrets are sealed to. // // Replacing whatever was there. A node that rejoins has generated a new one, and everything // sealed to the old key is unreadable to it -- which is why this does not merge and why what it // invalidates is reported rather than repaired silently. func (i *Inventory) RecordSealingKey(ctx context.Context, node, key string) error { if key == "" { return nil } _, err := i.store.Pool().Exec(ctx, `update node set sealing_key = $2 where id = $1`, node, key) return err } // SealingKeyOf is the key to seal something to for a node, empty if it has none. func (i *Inventory) SealingKeyOf(ctx context.Context, name string) (string, error) { var key *string err := i.store.Pool().QueryRow(ctx, `select sealing_key from node where name = $1`, name).Scan(&key) if err != nil { return "", err } if key == nil { return "", nil } return *key, nil } // RecordOverlayKey keeps the public half a node generated. func (i *Inventory) RecordOverlayKey(ctx context.Context, node, key string) error { if strings.TrimSpace(key) == "" { return errors.New("a node reported an empty overlay key") } _, err := i.store.Pool().Exec(ctx, `update node set overlay_key = $2 where id = $1`, node, key) return err } // Overlays is every node's place on the private network, which is what computing the graph needs. // // Every node at once, deliberately: a peer list is derived from all of them, and that is the // whole reason this is the control plane's work rather than a node's. func (i *Inventory) Overlays(ctx context.Context) ([]Overlay, error) { rows, err := i.store.Pool().Query(ctx, `select id, name, coalesce(overlay_key,''), coalesce(endpoint,''), coalesce(site,''), is_hub, coalesce(host(overlay_address),'') from node order by name`) if err != nil { return nil, err } defer rows.Close() var out []Overlay for rows.Next() { var o Overlay if err := rows.Scan(&o.Node, &o.Name, &o.Key, &o.Endpoint, &o.Site, &o.Hub, &o.Address); err != nil { return nil, err } out = append(out, o) } return out, rows.Err() } // ErrNotOneHub is what the mesh says when the graph cannot be computed. // // Its own error because it is not a fault in any node: it means nobody has said which node is the // hub, and a mesh with no hub has no path between sites at all. The old arrangement inferred this // from an address prefix and failed silently when nobody knew the convention. var ErrNotOneHub = errors.New("this mesh has no hub, so there is no path between sites") // SetPlace declares where a node is and how it is reached. func (i *Inventory) SetPlace(ctx context.Context, name, endpoint, site string, hub bool, address string) error { node, err := i.NodeByName(ctx, name) if err != nil { return err } var addr any if strings.TrimSpace(address) != "" { addr = address } _, err = i.store.Pool().Exec(ctx, `update node set endpoint = nullif($2,''), site = nullif($3,''), is_hub = $4, overlay_address = $5::inet where id = $1`, node.ID, endpoint, site, hub, addr) return err } // Outcomes a node's last report can have. const ( // OutcomeApplied is everything the declaration asked for. OutcomeApplied = "applied" // OutcomeFailed is some of it. The machine is in a state nobody declared. OutcomeFailed = "failed" // OutcomeRefused is none of it: the host would not accept the declaration at all, so the // machine is exactly as it was. A different situation from failing, with a different remedy — // one is fixed on the machine and the other in what was sent. OutcomeRefused = "refused" ) // Doing is what one machine did with what it was last sent. type Doing struct { Node string Outcome string Refused string Failed []FailedResource Applied int At time.Time } // FailedResource is one thing a node could not do. type FailedResource struct { ID string `json:"id"` Error string `json:"error"` } // Wrong reports whether this machine needs somebody to look at it. func (d Doing) Wrong() bool { return d.Outcome != OutcomeApplied } // RecordDoing keeps what a node said it did. // // One row per node, replaced. The question is the machine's current state — "this failed an hour // ago and then succeeded" is not something anybody needs to look at, and a table of every report // would bury the ones that matter. func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error { failed, err := json.Marshal(d.Failed) if err != nil { return err } _, err = i.store.Pool().Exec(ctx, `insert into node_report (node, outcome, refused, failed, applied, at) values ($1, $2, $3, $4, $5, now()) on conflict (node) do update set outcome = excluded.outcome, refused = excluded.refused, failed = excluded.failed, applied = excluded.applied, at = excluded.at`, node, d.Outcome, d.Refused, failed, d.Applied) return err } // NotDoingWhatTheyWereTold is every machine whose last report was not a clean apply. // // The list somebody wants when they ask what is wrong. A machine that has never reported is // absent rather than listed: it may be new, or switched off, and "never said anything" is a // different situation from "said it could not" — which is what `node list` reports as last heard // from. func (i *Inventory) NotDoingWhatTheyWereTold(ctx context.Context) ([]Doing, error) { rows, err := i.store.Pool().Query(ctx, `select n.name, r.outcome, r.refused, r.failed, r.applied, r.at from node_report r join node n on n.id = r.node where r.outcome <> $1 order by r.at desc`, OutcomeApplied) if err != nil { return nil, err } defer rows.Close() var out []Doing for rows.Next() { var d Doing var failed []byte if err := rows.Scan(&d.Node, &d.Outcome, &d.Refused, &failed, &d.Applied, &d.At); err != nil { return nil, err } if err := json.Unmarshal(failed, &d.Failed); err != nil { return nil, err } out = append(out, d) } return out, rows.Err() } // DoingOf is what one machine last did, and whether it has said anything at all. func (i *Inventory) DoingOf(ctx context.Context, name string) (Doing, bool, error) { node, err := i.NodeByName(ctx, name) if err != nil { return Doing{}, false, err } var d Doing var failed []byte err = i.store.Pool().QueryRow(ctx, `select outcome, refused, failed, applied, at from node_report where node = $1`, node.ID). Scan(&d.Outcome, &d.Refused, &failed, &d.Applied, &d.At) if errors.Is(err, pgx.ErrNoRows) { return Doing{}, false, nil } if err != nil { return Doing{}, false, err } d.Node = name if err := json.Unmarshal(failed, &d.Failed); err != nil { return Doing{}, false, err } return d, true, nil }