Give a machine port one holder: refuse ssh's, another module's and a doubled one, and release the assignment a given port replaces (hq ADR 0100)
This commit is contained in:
@@ -5,6 +5,8 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
@@ -593,11 +595,157 @@ func (i *Inventory) SetSettings(ctx context.Context, nodeName, module string, va
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = i.store.Pool().Exec(ctx,
|
||||
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)
|
||||
return wrapModule(err, module)
|
||||
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 {
|
||||
|
||||
@@ -257,3 +257,64 @@ func TestAGivenPortIsNeverAssigned(t *testing.T) {
|
||||
t.Fatalf("a fixed port given to another module was handed over: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// novox/hq ADR 0100: a given machine port has one holder. It is refused when another module has it,
|
||||
// assigned or given, when the same layer gives it twice, and when it is ssh's; and when it replaces
|
||||
// what the mesh assigned for that port, the assignment is given back.
|
||||
func TestAGivenPortHasOneHolderAndReplacesTheAssignment(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t, "postgres", "web", "cache")
|
||||
ctx := t.Context()
|
||||
give := func(module string, ports map[string]any) error {
|
||||
return inv.SetSettings(ctx, node, module, map[string]any{catalogue.PortsSetting: ports})
|
||||
}
|
||||
|
||||
web, err := inv.PortFor(ctx, node, "web", 8080, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := give("postgres", map[string]any{"5432": web.Machine}); !errors.Is(err, ErrPortTaken) {
|
||||
t.Fatalf("a port the mesh assigned to web was given to postgres: %v", err)
|
||||
}
|
||||
if err := give("postgres", map[string]any{"5432": 22}); err == nil {
|
||||
t.Fatal("ssh's port was given")
|
||||
}
|
||||
if err := give("postgres", map[string]any{"5432": 5433, "5433": 5433}); err == nil {
|
||||
t.Fatal("one machine port was given for two ports")
|
||||
}
|
||||
if err := give("cache", map[string]any{"6379": 6380}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := give("postgres", map[string]any{"5432": 6380}); !errors.Is(err, ErrPortTaken) {
|
||||
t.Fatalf("a port given to cache was given to postgres: %v", err)
|
||||
}
|
||||
|
||||
// Postgres was assigned a port for 5432; given one, the assignment is released and the number
|
||||
// is free again.
|
||||
assigned, err := inv.PortFor(ctx, node, "postgres", 5432, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := give("postgres", map[string]any{"5432": 5433}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Given again, the same: its own given port is not a collision with itself.
|
||||
if err := give("postgres", map[string]any{"5432": 5433}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
held, err := inv.PortsFor(ctx, node)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, a := range held {
|
||||
if a.Module == "postgres" {
|
||||
t.Fatalf("the assignment a given port replaced is still held: %+v", a)
|
||||
}
|
||||
}
|
||||
other, err := inv.PortFor(ctx, node, "cache", 11211, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if other.Machine != assigned.Machine {
|
||||
t.Fatalf("the released port %d was not free again (got %d)", assigned.Machine, other.Machine)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user