A public route used to carry its whole hostname as a literal in the module manifest, so running the same catalogue against a different domain meant overriding that literal on every routed module, per node. The mesh was, in effect, holding a map of names to services: the one thing it should never hold, because the subdomain is the operator's choice and the domain is the node's. Compose instead. A route contribution carries a `label` (the subdomain); a node carries its `public_domain` as node-level configuration; the mesh joins `<label>.<public-domain>` and grants exactly that, interpreting neither half. Held as a node property beside the node's other node-level facts (endpoint, site, overlay address), not in a module's settings — the ADR calls it node-level, and the settings table is keyed per module. Additive, so an unmigrated catalogue keeps working: a contribution that still carries a full `name` and no `label` passes through unchanged, and the catalogue can migrate module by module. A labelled contribution on a node with no public domain composes nothing, reading downstream as a route that named no host. And propagate: each granted route name is published into internal resolution mesh-wide, mapped to the node that serves it, alongside the `<node>.internal` names every container already gets. So a container — and an internal ACME validator, which cannot complete a challenge for a name it cannot reach — resolves a routed name to the proxy that serves it. Name-agnostic throughout: the mesh propagates whatever names it was told to serve and knows nothing about what they mean. novox/hq 02-DECISIONS/0056 Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
659 lines
23 KiB
Go
659 lines
23 KiB
Go
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
|
|
}
|
|
|
|
// SetPublicDomain records the domain a node's routed names are composed under.
|
|
//
|
|
// A node-level fact (novox/hq ADR 0056), kept beside the node's other node-level configuration
|
|
// rather than in a module's settings: the subdomain is a module's to choose and the domain is the
|
|
// node's, and the mesh joins the two without interpreting either. An empty domain clears it — a
|
|
// node that stops facing the outside composes no names — which is why this writes null rather than
|
|
// refusing.
|
|
func (i *Inventory) SetPublicDomain(ctx context.Context, name, domain string) error {
|
|
node, err := i.NodeByName(ctx, name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`update node set public_domain = nullif($2, '') where id = $1`, node.ID, strings.TrimSpace(domain))
|
|
return err
|
|
}
|
|
|
|
// PublicDomainOf is the domain a node composes its routed names under, empty if it has none.
|
|
func (i *Inventory) PublicDomainOf(ctx context.Context, name string) (string, error) {
|
|
var domain *string
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select public_domain from node where name = $1`, name).Scan(&domain)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return "", fmt.Errorf("%w: %s", ErrNoSuchNode, name)
|
|
}
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if domain == nil {
|
|
return "", nil
|
|
}
|
|
return *domain, 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
|
|
// Declared is the digest of the declaration the report was about; empty when the machine
|
|
// did not say.
|
|
Declared string
|
|
}
|
|
|
|
// 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, declared)
|
|
values ($1, $2, $3, $4, $5, now(), $6)
|
|
on conflict (node) do update set outcome = excluded.outcome, refused = excluded.refused,
|
|
failed = excluded.failed, applied = excluded.applied, at = excluded.at,
|
|
declared = excluded.declared`,
|
|
node, d.Outcome, d.Refused, failed, d.Applied, d.Declared)
|
|
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
|
|
}
|
|
|
|
// Reported is one machine's last word set beside what was last asked of it.
|
|
//
|
|
// The pair is what makes "has it caught up" answerable: a machine that reported *after* it was
|
|
// sent the current declaration has acted on it; one that reported before is still working, or
|
|
// has not started — and "not waiting" alone cannot tell those apart, because the sent digest is
|
|
// recorded at send, not at apply. The lab asserted on a machine the moment its declaration was
|
|
// current and found containers that did not exist yet.
|
|
type Reported struct {
|
|
Node string
|
|
Outcome string
|
|
// At is when it last reported; nil if it never has.
|
|
At *time.Time
|
|
// Sent is when the current declaration went to it; nil if nothing ever did.
|
|
Sent *time.Time
|
|
// Current is whether the last report names the declaration last sent — the machine has
|
|
// acted on the current words, not merely spoken after they were written. False also covers
|
|
// a machine that has not said which, which is every host from before reports carried it.
|
|
Current bool
|
|
}
|
|
|
|
// LastReports is every machine's last report beside when it was last sent a declaration.
|
|
func (i *Inventory) LastReports(ctx context.Context) ([]Reported, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select n.name, coalesce(r.outcome, ''), r.at, n.sent_at,
|
|
r.declared is not null and r.declared <> '' and r.declared = n.sent
|
|
from node n left join node_report r on r.node = n.id
|
|
order by n.name`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []Reported
|
|
for rows.Next() {
|
|
var r Reported
|
|
if err := rows.Scan(&r.Node, &r.Outcome, &r.At, &r.Sent, &r.Current); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, r)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// RecordSent keeps a digest of the declaration a machine was last sent.
|
|
//
|
|
// **A digest rather than the declaration.** The mesh can compute what a machine should be at any
|
|
// moment; keeping a copy would be a second account of it, able to disagree with the first. What
|
|
// cannot be recomputed is what was *actually sent*, and that is the whole difference between a
|
|
// machine that is out of date and one that has never been told.
|
|
func (i *Inventory) RecordSent(ctx context.Context, node, digest string) error {
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`update node set sent = $2, sent_at = now() where id = $1`, node, digest)
|
|
return err
|
|
}
|
|
|
|
// Waiting is every machine whose declaration has changed since it was last sent one.
|
|
//
|
|
// The caller works out what each machine should be now, because only it can — resolution is the
|
|
// control plane's and this context holds records. What is answered here is the comparison.
|
|
//
|
|
// **A machine that has never been sent anything is waiting**, and says so differently: it is not
|
|
// out of date, it has never been told, and the remedy is the same push while the situation is not
|
|
// the same at all.
|
|
func (i *Inventory) Waiting(ctx context.Context, would map[string]string) ([]Machine, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select name, coalesce(sent, ''), sent_at from node order by name`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []Machine
|
|
for rows.Next() {
|
|
var m Machine
|
|
var at *time.Time
|
|
if err := rows.Scan(&m.Node, &m.Sent, &at); err != nil {
|
|
return nil, err
|
|
}
|
|
m.SentAt = at
|
|
wanted, known := would[m.Node]
|
|
if !known {
|
|
// Nothing was computed for it — it resolves to nothing, or the caller did not ask.
|
|
// Silence rather than a guess: saying "waiting" about a machine nobody worked out
|
|
// would be inventing a comparison.
|
|
continue
|
|
}
|
|
if m.Sent == wanted {
|
|
continue
|
|
}
|
|
m.Never = m.Sent == ""
|
|
out = append(out, m)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// Machine is one machine that has not been sent what it should be.
|
|
type Machine struct {
|
|
Node string
|
|
// Sent is the digest it last received, empty if it has never received one.
|
|
Sent string
|
|
SentAt *time.Time
|
|
// Never is true when it has never been sent anything, which is a different situation from
|
|
// being out of date and reads differently to whoever is looking.
|
|
Never bool
|
|
}
|