package inventory import ( "context" "crypto/rand" "crypto/sha256" "encoding/base64" "encoding/hex" "encoding/json" "errors" "fmt" "sort" "strings" "time" "github.com/jackc/pgx/v5" "github.com/novox/mesh-controller/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 // Adopted is whether this node is adopted rather than converged (novox/hq ADR 0100): what is // found on it is kept until its module is taken, and the firewall found on it stays in force. // AdoptedSince is when it last became so; zero for a converged node. Adopted bool AdoptedSince time.Time // Account is the operator's login on this machine — `jochens` on novox, `ace` on ace (novox/hq // to-be 29). Empty when none is known yet. AccountHome is where that account's home is; empty // means derive it (/root for root, /home/ otherwise), so the common case needs no // entry. What decides who a file under a home is owned by, and which account `ssh ` uses. Account string AccountHome string } // Home is the account's home directory, derived when not stored: /root for root, /home/ // otherwise. Empty only when there is no account at all. func (n Node) Home() string { if n.AccountHome != "" { return n.AccountHome } switch n.Account { case "": return "" case "root": return "/root" default: return "/home/" + n.Account } } // 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) { return i.AddNodeAs(ctx, name, false) } // AddNodeAs creates a node record, adopted or converged (novox/hq ADR 0100). The operator says // which; a node added without saying is converged, as every node was before adoption existed. func (i *Inventory) AddNodeAs(ctx context.Context, name string, adopted bool) (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 var since *time.Time err := i.store.Pool().QueryRow(ctx, `insert into node (name, adopted, adopted_since) values ($1, $2, case when $2 then now() end) returning id, name, created, adopted, adopted_since`, name, adopted).Scan(&n.ID, &n.Name, &n.Created, &n.Adopted, &since) if err != nil { if strings.Contains(err.Error(), "node_name_key") { return Node{}, fmt.Errorf("%w: %s", ErrNameTaken, name) } return Node{}, err } if since != nil { n.AdoptedSince = *since } return n, nil } // nodeColumns and scanNode are the one reading of a node row, so every way of finding a node // says whether it is adopted. const nodeColumns = `id, name, created, last_seen, adopted, adopted_since, account, account_home` func scanNode(row pgx.Row) (Node, error) { var n Node var seen, since *time.Time if err := row.Scan(&n.ID, &n.Name, &n.Created, &seen, &n.Adopted, &since, &n.Account, &n.AccountHome); err != nil { return Node{}, err } if seen != nil { n.LastSeen = *seen } if since != nil { n.AdoptedSince = *since } return n, nil } // SetAccount records the operator account on a node — its human login — and optionally where that // account's home is (novox/hq to-be 29). An empty home means the mesh derives it. Clearing the // account (empty name) is allowed: a machine may stop having a known operator. func (i *Inventory) SetAccount(ctx context.Context, node, account, home string) error { tag, err := i.store.Pool().Exec(ctx, `update node set account = $1, account_home = $2 where name = $3`, account, home, node) if err != nil { return err } if tag.RowsAffected() == 0 { return fmt.Errorf("%w: %s", ErrNoSuchNode, node) } return 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 `+nodeColumns+` from node order by created, name`) if err != nil { return nil, err } defer rows.Close() var nodes []Node for rows.Next() { n, err := scanNode(rows) if err != nil { return nil, err } nodes = append(nodes, n) } return nodes, rows.Err() } // NodeByName finds one node record. func (i *Inventory) NodeByName(ctx context.Context, name string) (Node, error) { n, err := scanNode(i.store.Pool().QueryRow(ctx, `select `+nodeColumns+` from node where name = $1`, name)) 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") // ClaimLease is how long a claimed token is held for the one presenter that claimed it. Long // enough for an enrolment to be tried again through a store restart; short enough that a host // which gave up and was started over, with keys of its own, is not kept waiting long. const ClaimLease = 2 * time.Minute // ErrTokenInUse is a token another presenter holds a claim on right now. Not a refusal: the claim // lapses, and asking again after it is the answer. var ErrTokenInUse = errors.New("the token is being used by another enrolment") // Claim takes a token for one presenter — `by`, which names the key presenting it — for the length // of a lease, and says which node it enrols. The same presenter may claim it again, as may anyone // once the lease has lapsed; nothing is spent until Spend (novox/hq 04-ISSUES/083). // // A token this same presenter already spent may be claimed again when `again` says so — the // caller has proof the presenter holds the key's private half — and only while its claim's lease // is live: its spend reached the store and the answer did not reach the node, which asked again. // Refusing it then would lock out a machine the mesh holds as enrolled. Without the proof a spent // token stays spent to everyone, as ADR 0004 says. func (i *Inventory) Claim(ctx context.Context, secret, by string, again bool) (Node, error) { var id string err := i.store.Pool().QueryRow(ctx, `update enrolment_token set claimed_by = $2, claimed_until = case when redeemed is null then now() + $3::interval else claimed_until end where secret = $1 and expires > now() and ((redeemed is null and (claimed_by is null or claimed_by = $2 or claimed_until < now())) or (redeemed is not null and $4 and claimed_by = $2 and claimed_until > now())) returning node`, hashSecret(secret), by, ClaimLease.String(), again).Scan(&id) if errors.Is(err, pgx.ErrNoRows) { // Unusable, or held by someone else — told apart, because the second passes. var held bool probe := i.store.Pool().QueryRow(ctx, `select true from enrolment_token where secret = $1 and redeemed is null and expires > now()`, hashSecret(secret)).Scan(&held) switch { case probe == nil && held: return Node{}, ErrTokenInUse case probe != nil && !errors.Is(probe, pgx.ErrNoRows): // The store went away between the two questions: "not now", not a refusal. return Node{}, probe } return Node{}, ErrTokenRefused } if err != nil { return Node{}, err } return scanNode(i.store.Pool().QueryRow(ctx, `select `+nodeColumns+` from node where id = $1`, id)) } // Spend makes a claimed token used, only for the presenter holding the claim. The last write to the // store in an enrolment, so a token is spent exactly when the node it enrolled is complete. Spent // again by the same presenter is not an error: an answer lost after the first spend. func (i *Inventory) Spend(ctx context.Context, secret, by string) error { tag, err := i.store.Pool().Exec(ctx, `update enrolment_token set redeemed = coalesce(redeemed, now()) where secret = $1 and claimed_by = $2`, hashSecret(secret), by) if err != nil { return err } if tag.RowsAffected() == 0 { return ErrTokenRefused } return nil } // Redeem spends a token in one step and reports which node it was for. Enrolment claims and then // spends (Claim, Spend); this is the one-step form, kept for what spends a token outright. // // 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 } return scanNode(i.store.Pool().QueryRow(ctx, `select `+nodeColumns+` from node where id = $1`, id)) } // 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 0066), 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 // Since is when the machine first reported THIS failure — the same outcome, the same refusal, // the same failed resources by id — and Times is how many reports in a row have said it. A node // re-applies on a steady interval and reports each time (novox/hq ADR 0010), so a failure // that will never succeed arrives as the same report over and over, indistinguishable from // one that just happened until somebody counts (novox/hq 04-ISSUES/065). Nil and zero for a // machine doing what it was told. Since *time.Time Times int } // 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 } // StuckAfter is how many identical reports in a row make a failure one that will not fix itself. // // Three, because a node reports after every apply and applies on its reconcile interval: one // failure is an event, two may be the same event still under way, three separate applies saying // the same words is a machine looping on something that is not going to change (novox/hq // 04-ISSUES/065). Not a duration: a laptop that was shut for a week has had one attempt. const StuckAfter = 3 // Stuck reports whether this machine has been failing the same way for long enough that waiting // is no longer a plan. The failure is still the host's own words; this only says it is not new. func (d Doing) Stuck() bool { return d.Wrong() && d.Times >= StuckAfter && d.Since != nil } // sameFailure is whether two reports describe one failure: the same outcome, the same refusal, and // the same failed resources BY ID. Not by the host's words: an error that carries a duration, a // counter or a temporary path would read as new on every report, and the resource looping on it — // which is what stuck is for — would never be said to be (novox/hq 04-ISSUES/065). func sameFailure(a, b Doing) bool { if a.Outcome != b.Outcome || a.Refused != b.Refused || len(a.Failed) != len(b.Failed) { return false } ids := func(d Doing) []string { out := make([]string, 0, len(d.Failed)) for _, f := range d.Failed { out = append(out, f.ID) } sort.Strings(out) return out } x, y := ids(a), ids(b) for i := range x { if x[i] != y[i] { return false } } return true } // 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. // // **What the row also keeps is whether this failure is the one before.** The same outcome, the // same refusal, the same failed resources by id: then the failure did not just happen, it is // still happening, and the row keeps when it began and counts one more report. Any difference // starts again — a machine failing on a new resource is a new situation, not a longer one — and // a clean apply clears both (novox/hq 04-ISSUES/065). The previous row is read first and the // comparison made here, so "the same" is a rule this package states rather than a jsonb equality // that would restart the count on a changed word in an error. func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error { failed, err := json.Marshal(d.Failed) if err != nil { return err } var since *time.Time times := 0 if d.Outcome != OutcomeApplied { var before Doing var beforeFailed []byte err := i.store.Pool().QueryRow(ctx, `select outcome, refused, failed, failing_since, failures from node_report where node = $1`, node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times) switch { case errors.Is(err, pgx.ErrNoRows): case err != nil: return err default: if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil { return err } } now := time.Now() since, times = &now, 1 if err == nil && sameFailure(before, d) && before.Since != nil { since, times = before.Since, before.Times+1 } } _, err = i.store.Pool().Exec(ctx, `insert into node_report (node, outcome, refused, failed, applied, at, declared, failing_since, failures) values ($1, $2, $3, $4, $5, now(), $6, $7, $8) on conflict (node) do update set outcome = excluded.outcome, refused = excluded.refused, failed = excluded.failed, applied = excluded.applied, at = excluded.at, declared = excluded.declared, failing_since = excluded.failing_since, failures = excluded.failures`, node, d.Outcome, d.Refused, failed, d.Applied, d.Declared, since, times) 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, r.failing_since, r.failures 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, &d.Since, &d.Times); 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, failing_since, failures from node_report where node = $1`, node.ID). Scan(&d.Outcome, &d.Refused, &failed, &d.Applied, &d.At, &d.Since, &d.Times) 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 }