The report carries the digest of the declaration it applied (mesh-host 8211d8b), and the mesh stores it beside the outcome. `reported` rows in the status JSON now say `current`: whether the machine's last word names the declaration last sent. Not derivable from the timestamps beside it, which is why they were not enough: an apply begun under the previous declaration reports after the next send — newer, and still about the old words. The lab lost exactly that race between one test's closing push and the next test's opening one. Empty digests — every host from before reports carried one — read as not current, which errs toward waiting rather than toward asserting on files that are not there yet.
625 lines
22 KiB
Go
625 lines
22 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
|
|
}
|
|
|
|
// 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
|
|
}
|