247 lines
7.4 KiB
Go
247 lines
7.4 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
|
|
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"`
|
|
}
|
|
|
|
// 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
|
|
// 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 {
|
|
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,
|
|
adoption_reported = now(), last_seen = now()
|
|
where id = $1`, node, heldRaw, firewall, reachRaw)
|
|
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 []byte
|
|
var firewall *string
|
|
var at *time.Time
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select held, firewall, reachable, adoption_reported from node where name = $1`, name).
|
|
Scan(&heldRaw, &firewall, &reachRaw, &at)
|
|
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(reachRaw) > 0 {
|
|
if err := json.Unmarshal(reachRaw, &out.Reachable); err != nil {
|
|
return Adoption{}, err
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|