Files
mesh-controller/internal/inventory/nodes.go
T
jschoubben e6e1e3bc89
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
A gate judges its own send and the build it sent, and never puts the controller back behind its store (hq issue 352)
On 2026-10-09 a release's gate on the control node read the machine's
report against a newer send another plan had just made there, failed
three builds the machine had reported healthy, and put them back on
every machine to a controller older than the store's schema; that
controller then passed the newer plan's gate from its own health.

- A gate keeps what its send carried (digest, sequence) and reads the
  report against it; a report on the last send is on it too.
- A gate judges only the build the machine was last sent: another build
  there supersedes the judging — no verdict, nothing put back.
- A controller is told its build (MESH_CONTROLLER_VERSION, ${version}
  in a process's env) and records how far it reads the store's schema;
  a put-back to a build that reaches less, or never said, is refused
  and the current build kept, said as urgent.
- A release's open gate holds other sends of its modules there, and a
  plan's own first send waits on it.
2026-10-09 17:15:40 +02:00

1360 lines
54 KiB
Go

package inventory
import (
"context"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"regexp"
"sort"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/store"
)
// Inventory is this context, holding the store it exclusively owns.
type Inventory struct {
store *store.Store
// acting is the gate a write that acts passes (ActsUnder, novox/hq to-be 45 §6); nil passes all.
acting func(ctx context.Context) (uint64, error)
}
// 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/<account> otherwise), so the common case needs no
// entry. What decides who a file under a home is owned by, and which account `ssh <node>` uses.
Account string
AccountHome string
// AgentAccount is the login agents run as on this machine when it is not the operator's — `agent`
// on the control node (novox/hq ADR 0266): an account of their own, without sudo, so no agent there
// can become root without a person. Empty means agents run as the operator account. Stated at the
// controller's terminal only, never by a verb or a setting. AgentAccountHome is its home when not
// /home/<account>; empty means derive it.
AgentAccount string
AgentAccountHome string
// HostVersion is the version of the host this machine reported running (novox/hq 04-ISSUES/087).
// Empty when it has not said since the mesh began keeping it — which is not the same as running
// no host, so nothing derives "behind" from an empty one.
HostVersion string
}
// Home is the account's home directory, derived when not stored: /root for root, /home/<account>
// 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
}
}
// AgentHome is the agent account's home, derived when not stored; empty when no agent account is named.
func (n Node) AgentHome() string {
if n.AgentAccount == "" {
return ""
}
if n.AgentAccountHome != "" {
return n.AgentAccountHome
}
return "/home/" + n.AgentAccount
}
// 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,
agent_account, agent_account_home, host_version`
func scanNode(row pgx.Row) (Node, error) {
var n Node
var seen, since *time.Time
var host *string
if err := row.Scan(&n.ID, &n.Name, &n.Created, &seen, &n.Adopted, &since,
&n.Account, &n.AccountHome, &n.AgentAccount, &n.AgentAccountHome, &host); err != nil {
return Node{}, err
}
if host != nil {
n.HostVersion = *host
}
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.
//
// **Never the node's agent account** (novox/hq ADR 0266): the operator account may become root, and the
// agent account exists so agents cannot; naming the one as the other gives agents root. Refused here as
// SetAgentAccount refuses the other direction.
func (i *Inventory) SetAccount(ctx context.Context, node, account, home string) error {
account = strings.TrimSpace(account)
if account != "" {
n, err := i.NodeByName(ctx, node)
if err != nil {
return err
}
if n.AgentAccount != "" && n.AgentAccount == account {
return fmt.Errorf("%s is %s's agent account: the operator account may become root, and agents run as "+
"%s so that they cannot (novox/hq ADR 0266); clear the agent account first "+
"(node agent-account %s --clear) if the operator is to log in as it", account, node, account, node)
}
}
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
}
// loginName is what a login may be called: what useradd accepts by default, lower case, a letter or
// an underscore first.
var loginName = regexp.MustCompile(`^[a-z_][a-z0-9_-]{0,31}$`)
// SetAgentAccount records the account agents run as on a node, and optionally its home (novox/hq ADR
// 0266). An empty account clears it: agents run as the operator account again.
//
// **Refused, rather than recorded and judged later:** root, which is the very thing the account exists
// to keep agents from; the node's operator account, which may become root without a password and is
// what agents ran as before — naming it here would say the agents have an account of their own while
// they do not; and a name no machine would accept as a login.
func (i *Inventory) SetAgentAccount(ctx context.Context, node, account, home string) error {
account, home = strings.TrimSpace(account), strings.TrimSpace(home)
if account == "" && home != "" {
return errors.New("a home without an agent account says nothing; name the account too")
}
if account != "" {
if err := AgentAccountRefusal(account, home); err != nil {
return err
}
n, err := i.NodeByName(ctx, node)
if err != nil {
return err
}
if n.Account != "" && n.Account == account {
return fmt.Errorf("%s is %s's operator account: agents would run as the operator, who may become "+
"root; clear the agent account instead (node agent-account %s --clear)", account, node, node)
}
}
tag, err := i.store.Pool().Exec(ctx,
`update node set agent_account = $1, agent_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
}
// serviceAccounts are the system and service accounts a machine of the mesh has, or a module of the
// catalogue declares (showcase, and the accounts ADR 0259 gives the router and the channels). The controller
// cannot read a machine's user database, so this list is the controller's half; the node-engine's half is
// refusing to take an existing account below the first login uid as one that must never become root.
var serviceAccounts = map[string]bool{
"root": true, "bin": true, "daemon": true, "sys": true, "adm": true, "nobody": true, "mail": true,
"ftp": true, "http": true, "www-data": true, "git": true, "sshd": true, "dbus": true, "polkitd": true,
"postgres": true, "docker": true, "nats": true, "redis": true, "uuidd": true, "dnsmasq": true,
"avahi": true, "rtkit": true, "colord": true, "geoclue": true, "tss": true, "alpm": true, "usbmux": true,
"showcase": true, "messenger": true, "telegram": true,
}
// serviceAccount says whether a name is a system or service account: one of the list, or a name of
// systemd's own (systemd-…).
func serviceAccount(name string) bool {
return serviceAccounts[name] || strings.HasPrefix(name, "systemd-")
}
// AgentAccountRefusal is why an agent account cannot be named, or nil: root, a malformed login, or a
// home that is not an absolute path.
func AgentAccountRefusal(account, home string) error {
if account == "root" {
return errors.New("agents may not run as root: the agent account exists to keep them from it " +
"(novox/hq ADR 0266)")
}
if !loginName.MatchString(account) {
return fmt.Errorf("%q is not a login name: lower case letters, digits, _ and -, a letter or _ first, "+
"at most 32", account)
}
if serviceAccount(account) {
return fmt.Errorf("%q is a system or service account a machine or a module already has: the agent "+
"account is one the mesh creates for agents alone, which nothing else runs as or owns files as "+
"(novox/hq ADR 0266); name a new one, such as agent", account)
}
if home != "" && !strings.HasPrefix(home, "/") {
return fmt.Errorf("the agent account's home %q is not an absolute path", home)
}
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
}
// **And the account that secret is the password of** (novox/hq 04-ISSUES/146). The composed
// user list names an enrolment user for every node with a live token, and nothing minted a
// credential for it — so the composer left it out as a user with no password and every
// enrolment was refused by the server before the mesh heard of it.
//
// Recorded rather than minted: the token's secret IS the password, which is what lets a
// machine's first connection be authenticated by the thing it is enrolling with. It cannot be
// chosen here, because it has already been handed to whoever will present it.
//
// Outside the transaction on purpose. The token is what the mesh promised; a credential that
// the next composition rewrites anyway is not worth failing an issue over, and a token with no
// account is recoverable by issuing another, while an account with no token is a user nobody
// can be.
if err := i.RecordBusPassword(ctx, BusUser{
Username: broker.Principal{Kind: broker.KindEnrolment, Node: node.Name}.Username(),
Kind: BusEnrolment,
Node: node.Name,
}, secret); err != nil {
return Issued{}, fmt.Errorf(
"the token for %s was issued and the bus account it is the password of was not "+
"recorded, so this token cannot connect: %w", node.Name, 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))
}
// BindTokenToKey records the tunnel key a node's live token was issued for (novox/hq ADR 0169), so
// enrolment takes that key and no other. Refused when the node has no live token to bind: a key
// recorded against nothing would be a promise nothing keeps.
func (i *Inventory) BindTokenToKey(ctx context.Context, node, key string) error {
tag, err := i.store.Pool().Exec(ctx,
`update enrolment_token set overlay_key = $2
where node = $1 and redeemed is null and expires > now()`, node, key)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("no live token to issue for the tunnel key")
}
return nil
}
// TokenKey is the tunnel key a token was issued for, or empty when it was issued without one.
func (i *Inventory) TokenKey(ctx context.Context, secret string) (string, error) {
var key *string
err := i.store.Pool().QueryRow(ctx,
`select overlay_key from enrolment_token where secret = $1`, hashSecret(secret)).Scan(&key)
if err != nil || key == nil {
return "", err
}
return *key, nil
}
// 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
}
// RecordOutwardLinks keeps the links a machine reported as facing outside it.
//
// A reported fact, not a setting (novox/hq ADR 0140). It replaces the networks a machine used to be
// told to say it routes: the filter blocked everything passing through and then allowed the machine's
// own containers back by naming their address ranges, and every way of keeping that list correct
// failed — a constant describes one machine, and a recorded range goes stale in silence. The filter
// now constrains what arrives from outside and says nothing about what did not, and the one thing it
// needs is which links "outside" arrives on. The machine reads that from its own routing table on
// every apply, so it cannot go stale and nobody types it.
//
// An empty list clears it, which is what a machine with no route off itself reports. The mesh then
// composes no filter for that machine at all.
func (i *Inventory) RecordOutwardLinks(ctx context.Context, id string, links []string) error {
var kept []string
for _, name := range links {
if name = strings.TrimSpace(name); name != "" {
kept = append(kept, name)
}
}
if len(kept) == 0 {
_, err := i.store.Pool().Exec(ctx,
`update node set outward_links = null where id = $1`, id)
return err
}
body, err := json.Marshal(kept)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`update node set outward_links = $2 where id = $1`, id, string(body))
return err
}
// OutwardLinksOf is the links a machine reported as facing outside it, empty when it has reported
// none — which is a machine the mesh composes no filter for.
func (i *Inventory) OutwardLinksOf(ctx context.Context, name string) ([]string, error) {
var body []byte
err := i.store.Pool().QueryRow(ctx,
`select outward_links from node where name = $1`, name).Scan(&body)
if errors.Is(err, pgx.ErrNoRows) {
return nil, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
if err != nil {
return nil, err
}
if len(body) == 0 {
return nil, nil
}
var links []string
if err := json.Unmarshal(body, &links); err != nil {
return nil, fmt.Errorf("the outward links recorded for %s are not a list: %w", name, err)
}
return links, 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.
// **And whether this report was news**, which is what makes a fact about it worth stating (novox/hq
// ADR 0134). A machine reconciles continuously and reports each time; the same outcome about the same
// declaration is the same state said again, and a fact per report would be a fact per minute per
// machine that tells nobody anything. Read here because the previous row is read here anyway.
//
// **An account that claims no order clears the order kept** (novox/hq to-be 45 §6): the account kept is
// then one the next ordered report has nothing to compare against, and it is taken as it was before
// reports carried an order. RecordOrderedDoing keeps an ordered one.
func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) (news bool, err error) {
news, err = recordDoing(ctx, i.store.Pool(), node, d)
if err != nil {
return false, err
}
_, err = i.store.Pool().Exec(ctx,
`update node_report set reported_sequence = null, reported_epoch = null, report_sequence = null
where node = $1`, node)
return news, err
}
// queries is what recordDoing writes through: the pool, or a transaction an ordered account holds.
type queries interface {
QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
}
func recordDoing(ctx context.Context, q queries, node string, d Doing) (news bool, err error) {
failed, err := json.Marshal(d.Failed)
if err != nil {
return false, err
}
var before Doing
var beforeFailed []byte
found := q.QueryRow(ctx,
`select outcome, refused, failed, failing_since, failures, coalesce(declared,'')
from node_report where node = $1`,
node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times,
&before.Declared)
switch {
case errors.Is(found, pgx.ErrNoRows):
news = true
case found != nil:
return false, found
default:
if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil {
return false, err
}
news = before.Outcome != d.Outcome || before.Declared != d.Declared || !sameFailure(before, d)
}
var since *time.Time
times := 0
if d.Outcome != OutcomeApplied {
now := time.Now()
since, times = &now, 1
if found == nil && sameFailure(before, d) && before.Since != nil {
since, times = before.Since, before.Times+1
}
}
_, err = q.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)
if err != nil {
return false, err
}
return news, nil
}
// 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, coalesce(declared, '')
from node_report where node = $1`, node.ID).
Scan(&d.Outcome, &d.Refused, &failed, &d.Applied, &d.At, &d.Since, &d.Times, &d.Declared)
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
// Declared is the digest of the declaration the last report was about, and ReportedSequence that
// declaration's sequence as the report claimed it (zero from an engine that claims none): what a
// gate holds against the send it made, rather than against the send made last (novox/hq issue 352).
Declared string
ReportedSequence int64
}
// SentDeclaration is what one send carried to a machine, as a gate keeps it: the declaration's digest and
// its sequence (novox/hq issue 352). A report about this declaration, or about one sequenced after it, is a
// report on what the gate sent — whatever the machine was sent since.
type SentDeclaration struct {
Digest string `json:"digest"`
Sequence int64 `json:"sequence,omitempty"`
}
// ReportsOn says a report is about this send: the declaration itself; one the same machine was sequenced
// after it; or the declaration the machine was sent last (Current), which is this send or a later one —
// sends to a machine are made one after another. A send kept without a sequence is matched by its digest
// and by the last send alone.
func (s SentDeclaration) ReportsOn(r Reported) bool {
if r.Current || (s.Digest != "" && r.Declared == s.Digest) {
return true
}
return s.Sequence > 0 && r.ReportedSequence >= s.Sequence
}
// SentTo is the declaration a machine was last sent, by name: its digest and sequence, and false when it
// was never sent one.
func (i *Inventory) SentTo(ctx context.Context, name string) (SentDeclaration, bool, error) {
var digest *string
var seq *int64
err := i.store.Pool().QueryRow(ctx, `select sent, sequence from node where name = $1`, name).Scan(&digest, &seq)
if errors.Is(err, pgx.ErrNoRows) {
return SentDeclaration{}, false, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
if err != nil || digest == nil || *digest == "" {
return SentDeclaration{}, false, err
}
s := SentDeclaration{Digest: *digest}
if seq != nil {
s.Sequence = *seq
}
return s, true, nil
}
// 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,
coalesce(r.declared, ''), coalesce(r.reported_sequence, 0)
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, &r.Declared, &r.ReportedSequence); 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, and the build of each module
// it carried.
//
// **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.
//
// **And which build of each module** (novox/hq issue 259, ADR 0221): module name to the commit its
// build was made from. A digest cannot say whether a machine differs because a module moved to a build
// its upgrade policy holds back, or because of something a push made — a grant — and only the second
// is a push's to send to a machine it did not name. Nil records that it is not known, as for a
// declaration sent by hand.
func (i *Inventory) RecordSent(ctx context.Context, node, digest string, builds map[string]string) error {
return i.RecordSentUnder(ctx, node, digest, builds, 0)
}
// RecordSentUnder is RecordSent for a declaration that carried an epoch (novox/hq to-be 45 §6): kept
// beside the digest, so what the mesh would send is composed with it. Zero is one that carried none.
func (i *Inventory) RecordSentUnder(ctx context.Context, node, digest string, builds map[string]string,
epoch uint64) error {
var sentEpoch *int64
if epoch > 0 {
e := int64(epoch)
sentEpoch = &e
}
var carried *string
if builds != nil {
raw, err := json.Marshal(builds)
if err != nil {
return err
}
text := string(raw)
carried = &text
}
_, err := i.store.Pool().Exec(ctx,
`update node set sent = $2, sent_at = now(), sent_builds = $3::jsonb, sent_epoch = $4 where id = $1`,
node, digest, carried, sentEpoch)
return err
}
// AwaitingSince is since when a machine has had something waiting for its next send (novox/hq issue
// 275): the oldest assignment on it made after it was last sent, or — when nothing was assigned since —
// when it was last sent, or when it joined if it never was. The moment a bound on "not pushed yet" is
// read from: every change a push carries happened after the last push, and an assignment is the one
// that says when.
func (i *Inventory) AwaitingSince(ctx context.Context, name string) (time.Time, error) {
var since time.Time
err := i.store.Pool().QueryRow(ctx,
`select coalesce(
(select min(a.assigned) from assignment a
where a.node = n.id and (n.sent_at is null or a.assigned > n.sent_at)),
n.sent_at, n.created)
from node n where n.name = $1`, name).Scan(&since)
if errors.Is(err, pgx.ErrNoRows) {
return time.Time{}, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
return since, err
}
// SentBuilds is the build of each module a machine was last sent, by its name: module to the commit
// its build was made from (novox/hq issue 259). Known is false when that was not kept — a machine
// last sent before it was, sent a declaration by hand, or one the mesh does not know.
func (i *Inventory) SentBuilds(ctx context.Context, name string) (builds map[string]string, known bool, err error) {
var raw []byte
err = i.store.Pool().QueryRow(ctx, `select sent_builds from node where name = $1`, name).Scan(&raw)
if errors.Is(err, pgx.ErrNoRows) {
return nil, false, nil
}
if err != nil || raw == nil {
return nil, false, err
}
builds = map[string]string{}
if err := json.Unmarshal(raw, &builds); err != nil {
return nil, false, err
}
return builds, true, nil
}
// RecordSentBusUsers keeps a digest of the bus's user list a machine was just sent, by its name
// (novox/hq issue 249): whether the machine holding the bus must go first is whether this differs
// from the list composed now.
func (i *Inventory) RecordSentBusUsers(ctx context.Context, name, digest string) error {
_, err := i.store.Pool().Exec(ctx,
`update node set sent_bus_users = $2 where name = $1`, name, digest)
return err
}
// SentBusUsers is the digest of the bus's user list a machine was last sent, empty for none.
func (i *Inventory) SentBusUsers(ctx context.Context, name string) (string, error) {
var sent string
err := i.store.Pool().QueryRow(ctx,
`select sent_bus_users from node where name = $1`, name).Scan(&sent)
if errors.Is(err, pgx.ErrNoRows) {
return "", nil
}
return sent, err
}
// Outstanding is the digest of the declaration a machine was last sent, by its name, and empty
// for one that has never been sent anything.
//
// **By name rather than by id**, because the caller is the serving loop and what a node puts in a
// report is its name. Asked of one machine rather than read from Waiting's sweep, because it is
// asked per message: a report names the declaration it is about, and a report about one the mesh
// has already moved past is not acted on (design 25 §3).
func (i *Inventory) Outstanding(ctx context.Context, name string) (string, error) {
var sent string
err := i.store.Pool().QueryRow(ctx,
`select coalesce(sent, '') from node where name = $1`, name).Scan(&sent)
if errors.Is(err, pgx.ErrNoRows) {
// Not an error worth carrying up: a report from a machine the mesh has no record of has
// nothing to be stale against, and whatever is wrong with it is the listener's to say.
return "", nil
}
return sent, 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
}
// RecordHostVersion keeps the version of the host a machine reported running (novox/hq 04-ISSUES/087).
//
// Never cleared by a report that carries none: a bare word that the node is there says nothing about
// its host, and a machine whose host predates ADR 0141 reports none at all. So an empty version means
// the mesh has not been told, and the caller does not write it.
func (i *Inventory) RecordHostVersion(ctx context.Context, id, version string) error {
version = strings.TrimSpace(version)
if version == "" {
return nil
}
_, err := i.store.Pool().Exec(ctx,
`update node set host_version = $2, last_seen = now() where id = $1`, id, version)
return err
}
// NextSequence takes the next number for a declaration to this node, one higher than the last it
// was sent (novox/hq 04-ISSUES/107).
//
// **One statement, so two composers cannot take the same number.** The caller holds the node while it
// composes and sends, so in practice there is one; the increment is atomic anyway, because a rule
// that is true only while a lock is held somewhere else is a rule nobody can see from here.
func (i *Inventory) NextSequence(ctx context.Context, id string) (int64, error) {
var n int64
err := i.store.Pool().QueryRow(ctx,
`update node set sequence = coalesce(sequence, 0) + 1 where id = $1 returning sequence`, id).Scan(&n)
if err != nil {
return 0, fmt.Errorf("taking the next sequence for %s: %w", id, err)
}
return n, nil
}
// Sequence is the number of the last declaration this node was sent, and zero for one sent nothing
// since sends were numbered. Read, not taken: what the mesh WOULD send is composed with this, so it
// is byte for byte what it DID send when nothing else changed — a comparison that took a fresh number
// would read every machine as behind for ever (novox/hq 04-ISSUES/107).
func (i *Inventory) Sequence(ctx context.Context, id string) (int64, error) {
var n *int64
if err := i.store.Pool().QueryRow(ctx, `select sequence from node where id = $1`, id).Scan(&n); err != nil {
return 0, fmt.Errorf("reading the sequence of %s: %w", id, err)
}
if n == nil {
return 0, nil
}
return *n, nil
}
// RootSearchPending keeps when the controller first saw this node's agent account waiting for the node-engine's
// setuid search (novox/hq ADR 0266) and answers it: the first time it is seen since the last complete verdict,
// at, kept; seen again, what was kept. Kept across the engine's restarts and the controller's own.
func (i *Inventory) RootSearchPending(ctx context.Context, node string, at time.Time) (time.Time, error) {
var since time.Time
err := i.store.Pool().QueryRow(ctx,
`update node set agent_root_pending_since = coalesce(agent_root_pending_since, $2)
where name = $1 returning agent_root_pending_since`, node, at).Scan(&since)
if errors.Is(err, pgx.ErrNoRows) {
return time.Time{}, fmt.Errorf("%w: %s", ErrNoSuchNode, node)
}
return since, err
}
// RootSearchJudged forgets when a search began pending: a complete verdict came.
func (i *Inventory) RootSearchJudged(ctx context.Context, node string) error {
_, err := i.store.Pool().Exec(ctx,
`update node set agent_root_pending_since = null where name = $1 and agent_root_pending_since is not null`, node)
return err
}