1080 lines
37 KiB
Go
1080 lines
37 KiB
Go
package inventory
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
)
|
|
|
|
// ErrNoSuchModule is what the mesh says about a module it has never been told about.
|
|
var ErrNoSuchModule = errors.New("no module of that name")
|
|
|
|
// ErrStillAssigned is why a module cannot be forgotten.
|
|
//
|
|
// Its own error because it is not a fault: it means a machine is running that module now, and
|
|
// removing the record would leave the mesh unable to describe what is on it.
|
|
var ErrStillAssigned = errors.New("that module is still assigned to nodes")
|
|
|
|
// Source is where a module comes from and what has been built from it.
|
|
type Source struct {
|
|
Repository string
|
|
// Path is the module's directory inside that repository (novox/hq ADR 0069). Empty is the
|
|
// repository's root, which is a real answer rather than a missing one.
|
|
Path string
|
|
Ref string
|
|
// BuiltFrom is the commit the manifest the mesh holds was read at.
|
|
BuiltFrom string
|
|
// Head is the newest commit the source is known to have.
|
|
Head string
|
|
}
|
|
|
|
// Current reports whether what the mesh holds is what the source last had.
|
|
//
|
|
// A module with no source is always current: it was handed over directly, and there is nothing
|
|
// it could be behind. Saying "out of date" about it would be inventing a comparison.
|
|
func (s Source) Current() bool {
|
|
if s.Repository == "" || s.Head == "" {
|
|
return true
|
|
}
|
|
return s.BuiltFrom == s.Head
|
|
}
|
|
|
|
// RegisterModule records a module, replacing what was there.
|
|
//
|
|
// Replacing rather than refusing, because a manifest changing is the ordinary case -- a module
|
|
// gains a requirement, a claim, a resource. What matters is that the change is visible the next
|
|
// time a node is resolved, which it is.
|
|
func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, from Source) error {
|
|
raw, err := json.Marshal(m)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// A module registered without provenance keeps whatever it had. Handing over a manifest by
|
|
// hand is a legitimate way to fix something in a hurry, and it should not silently erase the
|
|
// record of where the module normally comes from — which is the only thing that would say,
|
|
// afterwards, that the machine is running something nobody can rebuild.
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`insert into module (name, manifest, version, source, source_path, ref, built_from, source_head)
|
|
values ($1, $2, nullif($3,''), nullif($4,''), $7, nullif($5,''), nullif($6,''), nullif($6,''))
|
|
on conflict (name) do update set
|
|
manifest = excluded.manifest,
|
|
version = excluded.version,
|
|
registered = now(),
|
|
source = coalesce(excluded.source, module.source),
|
|
source_path = case when excluded.source is null then module.source_path
|
|
else excluded.source_path end,
|
|
ref = coalesce(excluded.ref, module.ref),
|
|
built_from = coalesce(excluded.built_from, module.built_from),
|
|
source_head = coalesce(excluded.built_from, module.source_head)`,
|
|
m.Module, raw, m.Version, from.Repository, from.Ref, from.BuiltFrom, from.Path)
|
|
return err
|
|
}
|
|
|
|
// SourceMoved records that a module's source has a newer commit than the mesh has built.
|
|
//
|
|
// This is the whole of noticing. Nothing here builds anything — it writes down that the two
|
|
// halves differ, which is what makes *is this current?* answerable without building, and what
|
|
// makes a module that nobody rebuilt visible rather than silent.
|
|
func (i *Inventory) SourceMoved(ctx context.Context, module, head string) error {
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`update module set source_head = $2, source_seen = now() where name = $1`, module, head)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return fmt.Errorf("%w: %s", ErrNoSuchModule, module)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SourceOf is where a module came from and whether the mesh is behind it.
|
|
func (i *Inventory) SourceOf(ctx context.Context, module string) (Source, error) {
|
|
var s Source
|
|
var repo, ref, built, head *string
|
|
var path string
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select source, source_path, ref, built_from, source_head from module where name = $1`,
|
|
module).Scan(&repo, &path, &ref, &built, &head)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return Source{}, fmt.Errorf("%w: %s", ErrNoSuchModule, module)
|
|
}
|
|
if err != nil {
|
|
return Source{}, err
|
|
}
|
|
for _, pair := range []struct {
|
|
from *string
|
|
to *string
|
|
}{{repo, &s.Repository}, {ref, &s.Ref}, {built, &s.BuiltFrom}, {head, &s.Head}} {
|
|
if pair.from != nil {
|
|
*pair.to = *pair.from
|
|
}
|
|
}
|
|
// Not in the loop above: the path is never null, because "the repository's root" is an answer
|
|
// rather than an absence.
|
|
s.Path = path
|
|
return s, nil
|
|
}
|
|
|
|
// Behind is every module the mesh has not built from what its source now has, with the nodes
|
|
// running the old one.
|
|
//
|
|
// The nodes are the point. "Is this module out of date" is a fact about the catalogue; "which
|
|
// machines are running last week's version" is the question somebody actually has, and it is the
|
|
// one novox/hq ADR 0010 names as the thing that must not be lost.
|
|
func (i *Inventory) Behind(ctx context.Context) (map[string][]string, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select m.name, coalesce(n.name, '')
|
|
from module m
|
|
left join assignment a on a.module = m.name
|
|
left join node n on n.id = a.node
|
|
where m.source is not null
|
|
and m.source_head is not null
|
|
and coalesce(m.built_from, '') is distinct from m.source_head
|
|
order by m.name, n.name`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := map[string][]string{}
|
|
for rows.Next() {
|
|
var module, node string
|
|
if err := rows.Scan(&module, &node); err != nil {
|
|
return nil, err
|
|
}
|
|
if _, seen := out[module]; !seen {
|
|
out[module] = nil
|
|
}
|
|
if node != "" {
|
|
out[module] = append(out[module], node)
|
|
}
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// Catalogue is every module the mesh knows about, which is what resolution needs: the question
|
|
// "how many modules provide this" cannot be asked of a subset.
|
|
func (i *Inventory) Catalogue(ctx context.Context) (map[string]catalogue.Manifest, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select manifest from module order by name`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := map[string]catalogue.Manifest{}
|
|
for rows.Next() {
|
|
var raw []byte
|
|
if err := rows.Scan(&raw); err != nil {
|
|
return nil, err
|
|
}
|
|
var m catalogue.Manifest
|
|
if err := json.Unmarshal(raw, &m); err != nil {
|
|
return nil, err
|
|
}
|
|
out[m.Module] = m
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// Provide records a module the control plane ships with itself.
|
|
//
|
|
// The mesh's own private network is one: what computes its files is in this binary, so the
|
|
// manifest is too. Marked so it can be told apart from a module somebody wrote — not because it
|
|
// behaves differently, but because "where did this come from" must have an answer for everything
|
|
// in the catalogue, and "it came with the control plane" is that answer.
|
|
func (i *Inventory) Provide(ctx context.Context, m catalogue.Manifest) error {
|
|
raw, err := json.Marshal(m)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`insert into module (name, manifest, version, source)
|
|
values ($1, $2, $3, 'the control plane')
|
|
on conflict (name) do update set manifest = excluded.manifest,
|
|
version = excluded.version, source = 'the control plane'`,
|
|
m.Module, raw, m.Version)
|
|
return err
|
|
}
|
|
|
|
// Provided reports whether a module came with the control plane rather than from a repository.
|
|
func (i *Inventory) Provided(ctx context.Context, name string) (bool, error) {
|
|
var source *string
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select source from module where name = $1`, name).Scan(&source)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return source != nil && *source == "the control plane", nil
|
|
}
|
|
|
|
// ErrStillHolds is why a module cannot be forgotten without saying so first.
|
|
//
|
|
// Its own error because it is not a fault either: the operator settings, the module's own secrets
|
|
// and the ports the mesh chose for it are all keyed on the module by name and all cascade when the
|
|
// row goes. Removing the module removes them, silently, and none of them can be recovered — a
|
|
// sealed secret least of all, because the mesh discarded the plaintext when it made it.
|
|
var ErrStillHolds = errors.New("the mesh still holds things for that module")
|
|
|
|
// Holdings is everything keyed on a module that would go with it.
|
|
//
|
|
// **Named one at a time rather than counted.** "3 settings" tells somebody there is something to
|
|
// lose and not whether they can afford to lose it; "the mesh-wide layer, and anchor's" tells them
|
|
// what to write down before they type the command again.
|
|
type Holdings struct {
|
|
// Mesh is true when a mesh-wide settings layer exists for the module.
|
|
Mesh bool
|
|
// Nodes are the machines with a settings layer of their own for it, sorted.
|
|
Nodes []string
|
|
// Secrets are the module's own secrets, as "<name> on <node>", sorted. These are the ones the
|
|
// mesh cannot make again: what is stored is sealed to a machine and the plaintext is gone.
|
|
Secrets []string
|
|
// Ports are the ports the mesh chose for it, as "<wanted> on <node>", sorted. Made once and
|
|
// kept (novox/hq ADR 0038) — removing the module gives that promise up.
|
|
Ports []string
|
|
}
|
|
|
|
// Any reports whether removing the module would discard anything.
|
|
func (h Holdings) Any() bool {
|
|
return h.Mesh || len(h.Nodes) > 0 || len(h.Secrets) > 0 || len(h.Ports) > 0
|
|
}
|
|
|
|
// Lines is what would be lost, one thing per line, for a person about to decide.
|
|
func (h Holdings) Lines() []string {
|
|
var out []string
|
|
if h.Mesh {
|
|
out = append(out, " settings, for the whole mesh")
|
|
}
|
|
for _, n := range h.Nodes {
|
|
out = append(out, " settings, on "+n)
|
|
}
|
|
for _, s := range h.Secrets {
|
|
out = append(out, " its own secret "+s+" — sealed, so the mesh cannot make it again")
|
|
}
|
|
for _, p := range h.Ports {
|
|
out = append(out, " the port "+p)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// HeldFor is everything the mesh keeps that is keyed on one module.
|
|
//
|
|
// Read rather than counted at the moment of removal, because the answer is the whole of what a
|
|
// person needs in order to say yes.
|
|
func (i *Inventory) HeldFor(ctx context.Context, name string) (Holdings, error) {
|
|
var held Holdings
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select coalesce(n.name, '') from settings s left join node n on n.id = s.node
|
|
where s.module = $1 order by n.name nulls first`, name)
|
|
if err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
for rows.Next() {
|
|
var node string
|
|
if err := rows.Scan(&node); err != nil {
|
|
rows.Close()
|
|
return Holdings{}, err
|
|
}
|
|
if node == "" {
|
|
held.Mesh = true
|
|
continue
|
|
}
|
|
held.Nodes = append(held.Nodes, node)
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
|
|
for _, read := range []struct {
|
|
query string
|
|
into *[]string
|
|
}{
|
|
{`select s.name || ' on ' || n.name from module_secret s join node n on n.id = s.node
|
|
where s.module = $1 order by n.name, s.name`, &held.Secrets},
|
|
{`select p.wanted::text || ' on ' || n.name from port_assignment p
|
|
join node n on n.id = p.node where p.module = $1 order by n.name, p.wanted`, &held.Ports},
|
|
} {
|
|
rows, err := i.store.Pool().Query(ctx, read.query, name)
|
|
if err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
for rows.Next() {
|
|
var one string
|
|
if err := rows.Scan(&one); err != nil {
|
|
rows.Close()
|
|
return Holdings{}, err
|
|
}
|
|
*read.into = append(*read.into, one)
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
}
|
|
return held, nil
|
|
}
|
|
|
|
// ForgetModule removes a module, unless a machine is running it, unless the mesh still holds
|
|
// things for it, and never one the control plane provides.
|
|
//
|
|
// **Refused rather than cascaded.** The settings, own-secrets and port assignments all name the
|
|
// module by a foreign key that cascades, so the row going takes them with it and says nothing.
|
|
// That is an action succeeding into a state its own verify would reject (novox/hq 04-ISSUES/017):
|
|
// the command reports "forgotten", the operator re-registers the module a moment later, and what
|
|
// comes back is a module with none of its configuration and none of its secrets — with nothing
|
|
// anywhere naming the moment they were lost.
|
|
func (i *Inventory) ForgetModule(ctx context.Context, name string) error {
|
|
if err := i.mayForget(ctx, name); err != nil {
|
|
return err
|
|
}
|
|
held, err := i.HeldFor(ctx, name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if held.Any() {
|
|
return fmt.Errorf("%w:\n%s\n\nAll of it goes when the module does. Run "+
|
|
"`module forget %s --and-what-it-holds` if that is what you mean",
|
|
ErrStillHolds, strings.Join(held.Lines(), "\n"), name)
|
|
}
|
|
return i.discard(ctx, name)
|
|
}
|
|
|
|
// DiscardModule removes a module and everything the mesh holds for it, having been told to.
|
|
//
|
|
// The same checks as ForgetModule except the one about what is held: a machine running it still
|
|
// refuses, and a module the control plane provides still refuses, because neither of those is
|
|
// something an operator can consent to on the module's behalf.
|
|
func (i *Inventory) DiscardModule(ctx context.Context, name string) (Holdings, error) {
|
|
if err := i.mayForget(ctx, name); err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
held, err := i.HeldFor(ctx, name)
|
|
if err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
if err := i.discard(ctx, name); err != nil {
|
|
return Holdings{}, err
|
|
}
|
|
// Returned so the caller can say what went, rather than "forgotten". A person who has just
|
|
// destroyed a sealed secret should be able to read which one from the output.
|
|
return held, nil
|
|
}
|
|
|
|
// mayForget is the part of forgetting that is not about what is held.
|
|
func (i *Inventory) mayForget(ctx context.Context, name string) error {
|
|
provided, err := i.Provided(ctx, name)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return fmt.Errorf("%w: %s", ErrNoSuchModule, name)
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if provided {
|
|
// Refused rather than removed and re-created on the next migrate, which would look like
|
|
// it worked and quietly come back. Taking it off a machine is what "I do not want this"
|
|
// means, and that is what unassign is for.
|
|
return fmt.Errorf(
|
|
"%s comes with the control plane and cannot be forgotten; unassign it instead", name)
|
|
}
|
|
|
|
var on []string
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select n.name from assignment a join node n on n.id = a.node where a.module = $1
|
|
order by n.name`, name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for rows.Next() {
|
|
var node string
|
|
if err := rows.Scan(&node); err != nil {
|
|
rows.Close()
|
|
return err
|
|
}
|
|
on = append(on, node)
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return err
|
|
}
|
|
if len(on) > 0 {
|
|
return fmt.Errorf("%w: %s. Unassign it first", ErrStillAssigned, strings.Join(on, ", "))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// discard is the removal itself, once it has been decided.
|
|
func (i *Inventory) discard(ctx context.Context, name string) error {
|
|
tag, err := i.store.Pool().Exec(ctx, `delete from module where name = $1`, name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return fmt.Errorf("%w: %s", ErrNoSuchModule, name)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Assign puts a module on a node.
|
|
//
|
|
// Records the intention and checks nothing. Whether the set of assignments can actually become a
|
|
// declaration is resolution's question, asked over the whole set at once — and asking it here,
|
|
// one module at a time, would let an assignment look accepted and then refuse when a second
|
|
// arrives.
|
|
func (i *Inventory) Assign(ctx context.Context, nodeName, module string) error {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := i.runsSomewhere(ctx, module); err != nil {
|
|
return err
|
|
}
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`insert into assignment (node, module) values ($1, $2) on conflict do nothing`,
|
|
node.ID, module)
|
|
if err != nil && strings.Contains(err.Error(), "assignment_module_fkey") {
|
|
return fmt.Errorf("%w: %s", ErrNoSuchModule, module)
|
|
}
|
|
return err
|
|
}
|
|
|
|
// Unassign takes a module off a node.
|
|
func (i *Inventory) Unassign(ctx context.Context, nodeName, module string) error {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`delete from assignment where node = $1 and module = $2`, node.ID, module)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return fmt.Errorf("%s is not assigned to %s", module, nodeName)
|
|
}
|
|
// And its ports go back. Kept-once-chosen is a promise about a module that is still here —
|
|
// held past unassignment, a fixed port stays claimed in the name of something that is gone,
|
|
// and the next module needing it is refused by a ghost. The same argument that lets a machine
|
|
// take a carried port back: a set that only grows keeps a port reserved for nothing.
|
|
return i.ReleasePorts(ctx, nodeName, module)
|
|
}
|
|
|
|
// Assigned is what a person put on this node, which is not the same as what it runs: resolution
|
|
// adds whatever those modules require.
|
|
func (i *Inventory) Assigned(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 assignment 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()
|
|
}
|
|
|
|
// ProfileOf is what a node last said it can do, as resolution needs it: the capabilities that are
|
|
// present, and nothing else.
|
|
func (i *Inventory) ProfileOf(ctx context.Context, nodeName string) (map[string]bool, error) {
|
|
var raw []byte
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select profile from node where name = $1`, nodeName).Scan(&raw)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, fmt.Errorf("%w: %s", ErrNoSuchNode, nodeName)
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
out := map[string]bool{}
|
|
if len(raw) == 0 {
|
|
// A node that has never reported. Not an error, and not an empty machine either — every
|
|
// capability will read as absent, so anything requiring one is refused with "the wrong
|
|
// machine", which is wrong but visible. Better than assuming it can do everything.
|
|
return out, nil
|
|
}
|
|
held, err := profileFrom(raw)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, c := range held {
|
|
if c.Present {
|
|
out[c.Name] = true
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// Capability is one named fact a machine reported about itself.
|
|
//
|
|
// **Detected and never assumed** (novox/hq ADR 0009). The detail is the part that was being
|
|
// thrown away: a capability's presence gates an assignment and *its detail can also carry a
|
|
// value* — `seat: card1-DP-1`, an architecture, an amount of memory. So *can this run here* and
|
|
// *what should it be configured as* are the same fact read two ways, and keeping only the first
|
|
// read makes the second unanswerable.
|
|
//
|
|
// It is also what a refusal should quote. "This machine has no seat" is the right answer and
|
|
// "no graphics devices are present" is the reason, and only the machine knows the reason.
|
|
type Capability struct {
|
|
Name string `json:"name"`
|
|
Present bool `json:"present"`
|
|
// Detail is what the detector observed, in its own words — including when it found nothing,
|
|
// which is the case a person most needs explaining.
|
|
Detail string `json:"detail,omitempty"`
|
|
}
|
|
|
|
// Profile is everything a machine last reported about what it can do.
|
|
//
|
|
// The same read ProfileOf uses, unreduced. Two functions parsing one column would be two things
|
|
// to keep agreeing about what a profile is.
|
|
func (i *Inventory) Profile(ctx context.Context, nodeName string) ([]Capability, error) {
|
|
var raw []byte
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select profile from node where name = $1`, nodeName).Scan(&raw)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, fmt.Errorf("%w: %s", ErrNoSuchNode, nodeName)
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return profileFrom(raw)
|
|
}
|
|
|
|
func profileFrom(raw []byte) ([]Capability, error) {
|
|
if len(raw) == 0 {
|
|
// A machine that has never reported. Not an error and not an empty machine: nil says
|
|
// "nothing is known", where an empty slice would say "it reported having nothing", and
|
|
// those send a person to different places.
|
|
return nil, nil
|
|
}
|
|
var reported struct {
|
|
Capabilities []Capability `json:"capabilities"`
|
|
}
|
|
if err := json.Unmarshal(raw, &reported); err != nil {
|
|
return nil, err
|
|
}
|
|
return reported.Capabilities, nil
|
|
}
|
|
|
|
// SetSettings records what somebody wants a module's configuration to say.
|
|
//
|
|
// An empty node name means the whole mesh. Replacing rather than merging what is already there:
|
|
// this is a statement of the whole layer, so removing a key is done by leaving it out, which is
|
|
// the only way removing one could work at all.
|
|
func (i *Inventory) SetSettings(ctx context.Context, nodeName, module string, values map[string]any) error {
|
|
raw, err := json.Marshal(values)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if nodeName == "" {
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`insert into settings (node, module, values) values (null, $1, $2)
|
|
on conflict (module) where node is null
|
|
do update set values = excluded.values, set_at = now()`, module, raw)
|
|
return wrapModule(err, module)
|
|
}
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
given, err := givenIn(module, raw)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tx, err := i.store.Pool().Begin(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer func() { _ = tx.Rollback(context.WithoutCancel(ctx)) }()
|
|
if len(given) > 0 {
|
|
if err := refuseGivenCollisions(ctx, tx, node.ID, nodeName, module, given); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
_, err = tx.Exec(ctx,
|
|
`insert into settings (node, module, values) values ($1, $2, $3)
|
|
on conflict (node, module) where node is not null
|
|
do update set values = excluded.values, set_at = now()`, node.ID, module, raw)
|
|
if err != nil {
|
|
return wrapModule(err, module)
|
|
}
|
|
// A given port replaces what the mesh assigned for that port: the assignment is given back,
|
|
// so the number is free for the next module rather than held for ever in the name of a port
|
|
// that now lives elsewhere.
|
|
for wanted := range given {
|
|
if _, err := tx.Exec(ctx,
|
|
`delete from port_assignment where node = $1 and module = $2 and wanted = $3`,
|
|
node.ID, module, wanted); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return tx.Commit(ctx)
|
|
}
|
|
|
|
// givenIn is the machine ports a node-level settings layer gives a module, software port →
|
|
// machine port (novox/hq ADR 0100). Nothing when the layer gives none; what is not a port is left
|
|
// for composition to refuse in its own words.
|
|
func givenIn(module string, raw []byte) (map[int]int, error) {
|
|
var layer map[string]any
|
|
if err := json.Unmarshal(raw, &layer); err != nil {
|
|
return nil, err
|
|
}
|
|
entries, ok := layer[catalogue.PortsSetting].(map[string]any)
|
|
if !ok {
|
|
return nil, nil
|
|
}
|
|
out := map[int]int{}
|
|
by := map[int]int{}
|
|
for text, value := range entries {
|
|
wanted, err := strconv.Atoi(text)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
at, ok := value.(float64)
|
|
if !ok {
|
|
continue
|
|
}
|
|
machine := int(at)
|
|
if machine == catalogue.SSHPort {
|
|
return nil, fmt.Errorf("%s cannot be given port %d for its %d: that is ssh's, the one "+
|
|
"port a machine may never lose", module, machine, wanted)
|
|
}
|
|
if other, twice := by[machine]; twice {
|
|
return nil, fmt.Errorf("%s gives machine port %d to both its %d and its %d; a machine "+
|
|
"port has one holder", module, machine, min(other, wanted), max(other, wanted))
|
|
}
|
|
by[machine] = wanted
|
|
out[wanted] = machine
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// refuseGivenCollisions refuses a given machine port another module on the node already has —
|
|
// assigned by the mesh or given by its own setting — or that this module has for another of its
|
|
// ports. A port it was assigned for the same software port is not a collision: the given one
|
|
// replaces it.
|
|
func refuseGivenCollisions(ctx context.Context, tx pgx.Tx, nodeID any, node, module string,
|
|
given map[int]int) error {
|
|
type holder struct {
|
|
module string
|
|
wanted int
|
|
}
|
|
held := map[int]holder{}
|
|
rows, err := tx.Query(ctx,
|
|
`select machine, module, wanted from port_assignment where node = $1`, nodeID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for rows.Next() {
|
|
var h holder
|
|
var machine int
|
|
if err := rows.Scan(&machine, &h.module, &h.wanted); err != nil {
|
|
rows.Close()
|
|
return err
|
|
}
|
|
held[machine] = h
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return err
|
|
}
|
|
rows, err = tx.Query(ctx,
|
|
`select module, values->'ports' from settings
|
|
where node = $1 and module <> $2 and jsonb_typeof(values->'ports') = 'object'`,
|
|
nodeID, module)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for rows.Next() {
|
|
var other string
|
|
var raw []byte
|
|
if err := rows.Scan(&other, &raw); err != nil {
|
|
rows.Close()
|
|
return err
|
|
}
|
|
var theirs map[string]any
|
|
if err := json.Unmarshal(raw, &theirs); err != nil {
|
|
rows.Close()
|
|
return err
|
|
}
|
|
for text, v := range theirs {
|
|
wanted, _ := strconv.Atoi(text)
|
|
if at, ok := v.(float64); ok {
|
|
held[int(at)] = holder{module: other, wanted: wanted}
|
|
}
|
|
}
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return err
|
|
}
|
|
|
|
wanted := make([]int, 0, len(given))
|
|
for w := range given {
|
|
wanted = append(wanted, w)
|
|
}
|
|
sort.Ints(wanted)
|
|
for _, w := range wanted {
|
|
machine := given[w]
|
|
h, taken := held[machine]
|
|
if !taken || (h.module == module && h.wanted == w) {
|
|
continue
|
|
}
|
|
if _, moving := given[h.wanted]; h.module == module && moving {
|
|
// Its own port for another of its software ports, which this same layer moves away.
|
|
continue
|
|
}
|
|
return fmt.Errorf("%w: %s cannot be given %d on %s for its %d — %s already has it for "+
|
|
"its %d", ErrPortTaken, module, machine, node, w, h.module, h.wanted)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func wrapModule(err error, module string) error {
|
|
if err != nil && strings.Contains(err.Error(), "settings_module_fkey") {
|
|
return fmt.Errorf("%w: %s", ErrNoSuchModule, module)
|
|
}
|
|
return err
|
|
}
|
|
|
|
// ClearSettings removes a layer.
|
|
func (i *Inventory) ClearSettings(ctx context.Context, nodeName, module string) error {
|
|
if nodeName == "" {
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`delete from settings where module = $1 and node is null`, module)
|
|
return err
|
|
}
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`delete from settings where module = $1 and node = $2`, module, node.ID)
|
|
return err
|
|
}
|
|
|
|
// SettingsFor is the layers that apply to one module on one node, in the order they are applied.
|
|
//
|
|
// The mesh's first, then the node's, so a node that differs is expressed by differing rather
|
|
// than by restating everything the rest of the mesh already says.
|
|
func (i *Inventory) SettingsFor(ctx context.Context, nodeName, module string) ([]catalogue.Layer, error) {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select node is null, values from settings
|
|
where module = $1 and (node is null or node = $2)
|
|
order by node is null desc`, module, node.ID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var layers []catalogue.Layer
|
|
for rows.Next() {
|
|
var meshWide bool
|
|
var raw []byte
|
|
if err := rows.Scan(&meshWide, &raw); err != nil {
|
|
return nil, err
|
|
}
|
|
values := map[string]any{}
|
|
if err := json.Unmarshal(raw, &values); err != nil {
|
|
return nil, err
|
|
}
|
|
from := nodeName
|
|
if meshWide {
|
|
from = catalogue.MeshWideLayer
|
|
}
|
|
layers = append(layers, catalogue.Layer{From: from, Values: values})
|
|
}
|
|
return layers, rows.Err()
|
|
}
|
|
|
|
// PinProvision records which node a machine gets a provision from.
|
|
//
|
|
// Only needed when more than one node could answer. Recordable before that, because a mesh with
|
|
// one database should not change where an existing machine gets its data the day a second
|
|
// arrives.
|
|
func (i *Inventory) PinProvision(ctx context.Context, nodeName, provision, provider string) error {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
from, err := i.NodeByName(ctx, provider)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = i.store.Pool().Exec(ctx,
|
|
`insert into provision_pin (node, name, provider) values ($1, $2, $3)
|
|
on conflict (node, name) do update set provider = excluded.provider, pinned_at = now()`,
|
|
node.ID, provision, from.ID)
|
|
return err
|
|
}
|
|
|
|
// UnpinProvision removes a choice, putting the question back.
|
|
func (i *Inventory) UnpinProvision(ctx context.Context, nodeName, provision string) error {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`delete from provision_pin where node = $1 and name = $2`, node.ID, provision)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return fmt.Errorf("%s was not told where to get %q from", nodeName, provision)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// PinsFor is what a node was told about where its provisions come from.
|
|
func (i *Inventory) PinsFor(ctx context.Context, nodeName string) (map[string]string, error) {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select p.name, n.name from provision_pin p join node n on n.id = p.provider
|
|
where p.node = $1`, node.ID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := map[string]string{}
|
|
for rows.Next() {
|
|
var name, provider string
|
|
if err := rows.Scan(&name, &provider); err != nil {
|
|
return nil, err
|
|
}
|
|
out[name] = provider
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// pinRows is how many pins a node holds, counted in the table rather than through the join that
|
|
// reads them.
|
|
//
|
|
// For a test that would otherwise pass for the wrong reason: PinsFor joins on the provider, so a
|
|
// pin left behind by a departed node is invisible through it whether it was cleaned up or not.
|
|
func (i *Inventory) pinRows(ctx context.Context, nodeName string) (int, error) {
|
|
node, err := i.NodeByName(ctx, nodeName)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
var n int
|
|
err = i.store.Pool().QueryRow(ctx,
|
|
`select count(*) from provision_pin where node = $1`, node.ID).Scan(&n)
|
|
return n, err
|
|
}
|
|
|
|
// Entry is one module as a catalogue shows it: what it is, where it came from, and who runs it.
|
|
type Entry struct {
|
|
Manifest catalogue.Manifest
|
|
Source Source
|
|
// On is every node this module is assigned to, sorted.
|
|
On []string
|
|
// Provided is true when the module came with the control plane rather than from a repository.
|
|
Provided bool
|
|
}
|
|
|
|
// Catalogued is every module the mesh knows about, with everything a person asks about one.
|
|
//
|
|
// **One query rather than a call per module.** A catalogue that costs a round trip per row is a
|
|
// catalogue nobody lists, and the questions here — what is this, where did it come from, who is
|
|
// running it — are asked together every time.
|
|
func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select m.name, m.manifest,
|
|
coalesce(m.source, ''), m.source_path, coalesce(m.ref, ''),
|
|
coalesce(m.built_from, ''), coalesce(m.source_head, ''),
|
|
coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}')
|
|
from module m
|
|
left join assignment a on a.module = m.name
|
|
left join node n on n.id = a.node
|
|
group by m.name, m.manifest, m.source, m.source_path, m.ref, m.built_from, m.source_head
|
|
order by m.name`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []Entry
|
|
for rows.Next() {
|
|
var raw []byte
|
|
var name string
|
|
var source Source
|
|
var on []string
|
|
if err := rows.Scan(&name, &raw, &source.Repository, &source.Path, &source.Ref,
|
|
&source.BuiltFrom, &source.Head, &on); err != nil {
|
|
return nil, err
|
|
}
|
|
var m catalogue.Manifest
|
|
if err := json.Unmarshal(raw, &m); err != nil {
|
|
return nil, err
|
|
}
|
|
entry := Entry{Manifest: m, Source: source, On: on}
|
|
if source.Repository == providedBy {
|
|
// It came with the control plane. Not a repository, and showing it as one would have
|
|
// somebody go looking for it.
|
|
entry.Provided = true
|
|
entry.Source = Source{}
|
|
}
|
|
out = append(out, entry)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// providedBy is what the source column says for a module the control plane ships.
|
|
const providedBy = "the control plane"
|
|
|
|
// Upgrade is what the mesh decided to do when a module's current version moves.
|
|
type Upgrade struct {
|
|
// RollOut is true when the machines running it should be sent the new version. False means
|
|
// record it and stop — which needs no record of its own, because a machine not running what
|
|
// the mesh would send it is already something the mesh reports.
|
|
RollOut bool
|
|
// Together is true when every machine running it is sent the new version at once, rather than
|
|
// one after another. Only meaningful when RollOut is.
|
|
Together bool
|
|
}
|
|
|
|
// UpgradeOf is what to do when this module moves.
|
|
//
|
|
// A module the mesh does not hold is not an error here: the catalogue may know of modules this
|
|
// mesh has never registered, and being told one of them moved is information, not a fault. The
|
|
// answer is the safe one — record it — because there is nothing to roll out to.
|
|
func (i *Inventory) UpgradeOf(ctx context.Context, module string) (Upgrade, error) {
|
|
var u Upgrade
|
|
var policy string
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select upgrade, upgrade_together from module where name = $1`, module).
|
|
Scan(&policy, &u.Together)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return Upgrade{}, nil
|
|
}
|
|
if err != nil {
|
|
return Upgrade{}, err
|
|
}
|
|
u.RollOut = policy == "roll-out"
|
|
return u, nil
|
|
}
|
|
|
|
// SetUpgradeOf records what to do when this module moves.
|
|
func (i *Inventory) SetUpgradeOf(ctx context.Context, module string, u Upgrade) error {
|
|
policy := "record"
|
|
if u.RollOut {
|
|
policy = "roll-out"
|
|
}
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`update module set upgrade = $2, upgrade_together = $3 where name = $1`,
|
|
module, policy, u.Together)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return fmt.Errorf("this mesh holds no module called %s", module)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Running is every machine assigned a module, in a stable order.
|
|
//
|
|
// **Assigned, not reported.** A machine that is assigned the module and has not applied it yet is
|
|
// exactly the machine an upgrade most needs to reach; waiting for it to report the old version
|
|
// first would mean the machines furthest behind are the last to be caught up.
|
|
func (i *Inventory) Running(ctx context.Context, module string) ([]string, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select n.name from assignment a join node n on n.id = a.node
|
|
where a.module = $1 order by n.name`, module)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []string
|
|
for rows.Next() {
|
|
var name string
|
|
if err := rows.Scan(&name); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, name)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// runsSomewhere refuses to put a module on a machine when it would put nothing there.
|
|
//
|
|
// **Not every thing the mesh builds is a thing a machine runs.** The image every module in the
|
|
// scripted toolchain is compiled on top of is built, versioned and depended upon like a module, and
|
|
// is registered as one so the mesh can track all three. It is not one: it declares no resources,
|
|
// because nothing about it belongs on a machine — its whole purpose is to be the starting point for
|
|
// other modules' builds.
|
|
//
|
|
// Without this, assigning it succeeds, the machine is sent a declaration containing nothing of it,
|
|
// and everything reports success. The operator has said "run this here" and the mesh has agreed to
|
|
// something it cannot do. Said at the assignment, which is where somebody is standing.
|
|
//
|
|
// A module whose resources are worked out per node declares none here and is still assignable —
|
|
// that is the point of it — so it is asked about separately rather than caught by the same test.
|
|
func (i *Inventory) runsSomewhere(ctx context.Context, module string) error {
|
|
var raw []byte
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select manifest from module where name = $1`, module).Scan(&raw)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
// Left to the insert, which already says this and says it the same way everywhere.
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var m catalogue.Manifest
|
|
if err := json.Unmarshal(raw, &m); err != nil {
|
|
// Unreadable is not the same as empty. A manifest the mesh cannot parse is a separate
|
|
// fault and refusing the assignment here would report it as the wrong one.
|
|
return nil
|
|
}
|
|
if !IsBuildInput(m) {
|
|
return nil
|
|
}
|
|
return fmt.Errorf(
|
|
"%s puts nothing on a machine, so there is nothing to assign. It is something this mesh "+
|
|
"builds and other modules are built on top of, not something a machine runs — "+
|
|
"`module list` shows what it produces", module)
|
|
}
|
|
|
|
// IsBuildInput reports whether a module exists to be built and never to be run.
|
|
//
|
|
// **The signal is that it builds something and places nothing** — not merely that it declares no
|
|
// resources. Those are different, and confusing them refused a module the mesh itself ships: the
|
|
// private network declares no resources either, because the control plane computes them when it
|
|
// composes a machine's declaration, and it is assigned to every machine that has to reach another
|
|
// one. Refusing it stopped a four-machine bed dead.
|
|
//
|
|
// Its own function because it is a judgement rather than a lookup, and a judgement with a wrong
|
|
// answer this expensive should be testable without a database.
|
|
func IsBuildInput(m catalogue.Manifest) bool {
|
|
builds := m.Build != nil && len(m.Build.Artifacts) > 0
|
|
return builds && m.Computed == "" && len(m.Resources) == 0
|
|
}
|