Files
mesh-controller/internal/inventory/adoption.go
T
jschoubben b5df244096 The mesh says what filters a converged machine: filters kept per node, shown by node show, named by status, and previewed with their fates (hq ADR 0168)
A host reports every table and chain that refuses traffic with its owner,
and a converged machine's found firewall's state. The controller keeps both
on the node's record (migration 0054), shows them on node show, names every
converged machine something other than the mesh filters in status — text
and JSON, and such a machine is not well — and the converge preview lists
what filters the machine with the fate of each: retired with the front end,
left as the runtime's, left as a ban, or left in force and not the mesh's.
What was invisible for eleven hours (issues 144, 145) is said by name.
2026-10-02 12:00:12 +02:00

373 lines
12 KiB
Go

package inventory
import (
"context"
"encoding/json"
"errors"
"fmt"
"sort"
"time"
"github.com/jackc/pgx/v5"
)
// A node is adopted or converged (novox/hq ADR 0100).
//
// Adopted: what is found on the machine is kept until its module is taken, and the firewall found
// there stays in force. Converged: the machine is what the mesh declares, as every node was before
// adoption existed. The controller is authoritative, and every declaration it sends says which.
// ErrNotAdopted is taking a module on a node that is converged. On a converged node every assigned
// module converges already; there is nothing to take.
var ErrNotAdopted = errors.New("the node is converged, so every module on it is taken already")
// ErrNotAssigned is taking a module that is not on the node. Taking is the cutover of a module
// the node runs; one it does not run has nothing to cut over.
var ErrNotAssigned = errors.New("that module is not assigned to the node")
// SetAdopted makes a node adopted or converged. Becoming adopted stamps when; converging stamps
// when too. Neither touches what was taken: what was taken stays taken when a node returns to
// adopted, and converging takes the rest by its own act.
func (i *Inventory) SetAdopted(ctx context.Context, name string, adopted bool) error {
node, err := i.NodeByName(ctx, name)
if err != nil {
return err
}
if node.Adopted == adopted {
return nil
}
if adopted {
_, err = i.store.Pool().Exec(ctx,
`update node set adopted = true, adopted_since = now() where id = $1`, node.ID)
return err
}
_, err = i.store.Pool().Exec(ctx,
`update node set adopted = false, adopted_since = null, converged_at = now() where id = $1`,
node.ID)
return err
}
// Take records that a module has been taken on an adopted node: its cutover. From then on the
// module's resources converge on that node like any other, replacing what was found.
//
// Refused on a converged node and for a module not assigned there. Taking again is not an error;
// the first time it was taken is kept.
func (i *Inventory) Take(ctx context.Context, nodeName, module string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
if !node.Adopted {
return fmt.Errorf("%w: %s", ErrNotAdopted, nodeName)
}
return i.take(ctx, node, module)
}
// take is Take without the adopted check, for converging, which takes every assigned module in
// the same act that makes the node converged.
func (i *Inventory) take(ctx context.Context, node Node, module string) error {
var assigned bool
if err := i.store.Pool().QueryRow(ctx,
`select exists (select 1 from assignment where node = $1 and module = $2)`,
node.ID, module).Scan(&assigned); err != nil {
return err
}
if !assigned {
return fmt.Errorf("%w: %s is not on %s; assign it first", ErrNotAssigned, module, node.Name)
}
_, err := i.store.Pool().Exec(ctx,
`insert into taken (node, module) values ($1, $2) on conflict do nothing`, node.ID, module)
return err
}
// Converge makes an adopted node converged in one act: every module assigned there is taken, and
// the node is recorded converged. Returned is what this act took, in name order.
func (i *Inventory) Converge(ctx context.Context, nodeName string) ([]string, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
if !node.Adopted {
return nil, fmt.Errorf("%s is converged already", nodeName)
}
tx, err := i.store.Pool().Begin(ctx)
if err != nil {
return nil, err
}
defer func() { _ = tx.Rollback(context.WithoutCancel(ctx)) }()
rows, err := tx.Query(ctx,
`insert into taken (node, module)
select node, module from assignment where node = $1
on conflict do nothing
returning module`, node.ID)
if err != nil {
return nil, err
}
var took []string
for rows.Next() {
var m string
if err := rows.Scan(&m); err != nil {
rows.Close()
return nil, err
}
took = append(took, m)
}
rows.Close()
if err := rows.Err(); err != nil {
return nil, err
}
if _, err := tx.Exec(ctx,
`update node set adopted = false, adopted_since = null, converged_at = now(),
held = null, reachable = null, firewall = null, adoption_reported = null
where id = $1`,
node.ID); err != nil {
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
sort.Strings(took)
return took, nil
}
// Taken is every module taken on a node, in name order — including one no longer assigned there:
// unassigning does not un-take.
func (i *Inventory) Taken(ctx context.Context, nodeName string) ([]string, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select module from taken where node = $1 order by module`, node.ID)
if err != nil {
return nil, err
}
defer rows.Close()
var out []string
for rows.Next() {
var m string
if err := rows.Scan(&m); err != nil {
return nil, err
}
out = append(out, m)
}
return out, rows.Err()
}
// Held is one file or container an adopted node found and keeps as it was until its module is
// taken. The node's own account, kept as it said it.
type Held struct {
ID string `json:"id"`
Module string `json:"module"`
Kind string `json:"kind"`
Target string `json:"target"`
Since time.Time `json:"since"`
Changed string `json:"changed,omitempty"`
Kept string `json:"kept,omitempty"`
// Facts is what a take compares (novox/hq ADR 0163), as the host reported it: for a found
// container its image and the image's date, the networks and their other members, mounts and
// ports, beside the declared image, ports and volumes, and whether the declared image is the
// older; for a found file whether the declared content differs and how.
Facts map[string]any `json:"facts,omitempty"`
}
// A Stray is a container a machine runs that the mesh neither wrote nor holds (ADR 0163).
type Stray struct {
Kind string `json:"kind"`
Name string `json:"name"`
Detail string `json:"detail,omitempty"`
}
// A Filter is one place on a machine that refuses traffic, with its owner (novox/hq ADR 0168).
type Filter struct {
Where string `json:"where"`
Owner string `json:"owner"`
Refuses string `json:"refuses"`
}
// Owners of a filter, as the host names them (ADR 0168).
const (
FilterMesh = "mesh"
FilterFoundFirewall = "found-firewall"
FilterRuntime = "runtime"
FilterBan = "ban"
FilterOther = "other"
)
// FoundFirewall is the state of a converged machine's found firewall (ADR 0168): in force now or
// not, and how it came to be inactive.
type FoundFirewall struct {
Kind string `json:"kind"`
Active bool `json:"active"`
RetiredBy string `json:"retired_by,omitempty"`
}
// Filtering is what a machine last said filters it (ADR 0168).
type Filtering struct {
Filters []Filter
FoundFirewall *FoundFirewall
}
// Alone is whether the machine is filtered by the mesh alone: nothing in its list but the mesh's
// own, the runtime's plumbing and bans, and no found firewall in force.
func (f Filtering) Alone() bool {
for _, x := range f.Filters {
if x.Owner == FilterOther || x.Owner == FilterFoundFirewall {
return false
}
}
return f.FoundFirewall == nil || !f.FoundFirewall.Active
}
// Others is every filter that is neither the mesh's, the runtime's nor a ban.
func (f Filtering) Others() []Filter {
var out []Filter
for _, x := range f.Filters {
if x.Owner == FilterOther || x.Owner == FilterFoundFirewall {
out = append(out, x)
}
}
return out
}
// RecordFiltering keeps what a machine last said filters it, replacing what was there (ADR 0168).
func (i *Inventory) RecordFiltering(ctx context.Context, nodeID string, filters []Filter, found *FoundFirewall) error {
raw, err := json.Marshal(nonNil(filters))
if err != nil {
return err
}
var foundRaw any
if found != nil {
b, err := json.Marshal(found)
if err != nil {
return err
}
foundRaw = string(b)
}
_, err = i.store.Pool().Exec(ctx,
`update node set filters = $2, found_firewall = $3 where id = $1`, nodeID, raw, foundRaw)
return err
}
// FilteringOf is what a machine last said filters it; empty for a machine that never said.
func (i *Inventory) FilteringOf(ctx context.Context, name string) (Filtering, error) {
var filtersRaw, foundRaw []byte
err := i.store.Pool().QueryRow(ctx,
`select filters, found_firewall from node where name = $1`, name).Scan(&filtersRaw, &foundRaw)
if errors.Is(err, pgx.ErrNoRows) {
return Filtering{}, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
if err != nil {
return Filtering{}, err
}
var out Filtering
if len(filtersRaw) > 0 {
if err := json.Unmarshal(filtersRaw, &out.Filters); err != nil {
return Filtering{}, err
}
}
if len(foundRaw) > 0 {
out.FoundFirewall = &FoundFirewall{}
if err := json.Unmarshal(foundRaw, out.FoundFirewall); err != nil {
return Filtering{}, err
}
}
return out, nil
}
// Reach is one thing reachable on an adopted node: a listening socket or a published port.
type Reach struct {
Protocol string `json:"protocol"`
Address string `json:"address"`
Port int `json:"port"`
By string `json:"by,omitempty"`
Published bool `json:"published,omitempty"`
ContainerPort int `json:"container-port,omitempty"`
}
// Adoption is what an adopted node last said about adoption, and when.
type Adoption struct {
Held []Held
Firewall string
Reachable []Reach
// Strays is what the machine runs that nobody asked for, as last reported (ADR 0163).
Strays []Stray
// At is when it said so; zero when it never has.
At time.Time
}
// RecordAdoption keeps what a node last reported about adoption, replacing what was there: the
// question is the machine as it is now.
func (i *Inventory) RecordAdoption(ctx context.Context, node string, held []Held, firewall string,
reachable []Reach) error {
return i.RecordAdoptionWithStrays(ctx, node, held, firewall, reachable, nil)
}
// RecordAdoptionWithStrays is RecordAdoption with what the machine says strays on it (ADR 0163).
func (i *Inventory) RecordAdoptionWithStrays(ctx context.Context, node string, held []Held, firewall string,
reachable []Reach, strays []Stray) error {
straysRaw, err := json.Marshal(nonNil(strays))
if err != nil {
return err
}
heldRaw, err := json.Marshal(nonNil(held))
if err != nil {
return err
}
reachRaw, err := json.Marshal(nonNil(reachable))
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`update node set held = $2, firewall = nullif($3, ''), reachable = $4, strays = $5,
adoption_reported = now(), last_seen = now()
where id = $1`, node, heldRaw, firewall, reachRaw, straysRaw)
return err
}
func nonNil[T any](s []T) []T {
if s == nil {
return []T{}
}
return s
}
// AdoptionOf is what a node last reported about adoption.
func (i *Inventory) AdoptionOf(ctx context.Context, name string) (Adoption, error) {
var heldRaw, reachRaw, straysRaw []byte
var firewall *string
var at *time.Time
err := i.store.Pool().QueryRow(ctx,
`select held, firewall, reachable, adoption_reported, strays from node where name = $1`, name).
Scan(&heldRaw, &firewall, &reachRaw, &at, &straysRaw)
if errors.Is(err, pgx.ErrNoRows) {
return Adoption{}, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
if err != nil {
return Adoption{}, err
}
var out Adoption
if firewall != nil {
out.Firewall = *firewall
}
if at != nil {
out.At = *at
}
if len(heldRaw) > 0 {
if err := json.Unmarshal(heldRaw, &out.Held); err != nil {
return Adoption{}, err
}
}
if len(straysRaw) > 0 {
if err := json.Unmarshal(straysRaw, &out.Strays); err != nil {
return Adoption{}, err
}
}
if len(reachRaw) > 0 {
if err := json.Unmarshal(reachRaw, &out.Reachable); err != nil {
return Adoption{}, err
}
}
return out, nil
}