From 4bb19c9e4067d652293022019961e18e7fc1d9ba Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 18:05:31 +0200 Subject: [PATCH] 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) --- cmd/mesh-controller/plan.go | 14 +++ internal/catalogue/adoption_test.go | 8 ++ internal/catalogue/filtering.go | 13 +++ internal/inventory/catalogue.go | 152 +++++++++++++++++++++++++++- internal/inventory/ports_test.go | 61 +++++++++++ 5 files changed, 246 insertions(+), 2 deletions(-) diff --git a/cmd/mesh-controller/plan.go b/cmd/mesh-controller/plan.go index 0f540c7..4b5cbbf 100644 --- a/cmd/mesh-controller/plan.go +++ b/cmd/mesh-controller/plan.go @@ -384,6 +384,20 @@ func renderingFor(ctx context.Context, open *stores, node string, given[m.Module] = g } } + // And one holder per machine port across the node's modules: settings written before this was + // refused where they are set are refused here rather than composed into two containers on + // one port. + holders := map[int]string{} + for _, m := range plan.Modules { + for _, at := range given[m.Module] { + if other, twice := holders[at]; twice && other != m.Module { + return catalogue.Rendering{}, inventory.Node{}, fmt.Errorf("%w: %s and %s are "+ + "both given machine port %d on %s", inventory.ErrPortTaken, other, m.Module, + at, node) + } + holders[at] = m.Module + } + } ports := map[string]map[int]int{} for _, m := range plan.Modules { diff --git a/internal/catalogue/adoption_test.go b/internal/catalogue/adoption_test.go index 0d839e1..b53e956 100644 --- a/internal/catalogue/adoption_test.go +++ b/internal/catalogue/adoption_test.go @@ -317,6 +317,14 @@ func TestAGivenPortIsTheNodesAndReachesSomething(t *testing.T) { if _, err := GivenPorts(store, node(map[string]any{"5432": float64(70000)})); err == nil { t.Fatal("a machine port that is not a port was given") } + if _, err := GivenPorts(store, node(map[string]any{"5432": float64(22)})); err == nil { + t.Fatal("ssh's port was given") + } + broker := anAdoptedAnchor().Modules[2] + if _, err := GivenPorts(broker, node(map[string]any{"5671": float64(5700), + "5672": float64(5700)})); err == nil { + t.Fatal("one machine port was given for two of the module's ports") + } if stray := UnusedSettings(store, node(map[string]any{"5432": float64(5433)})); len(stray) != 0 { t.Fatalf("a given port is called stray: %v", stray) } diff --git a/internal/catalogue/filtering.go b/internal/catalogue/filtering.go index ae1354f..453f651 100644 --- a/internal/catalogue/filtering.go +++ b/internal/catalogue/filtering.go @@ -518,9 +518,22 @@ func GivenPorts(m Manifest, layers []Layer) (map[int]int, error) { return nil, fmt.Errorf("%s gives port %d the machine port %v, which is not a port", m.Module, port, value) } + if at == SSHPort { + return nil, fmt.Errorf("%s gives port %d the machine port %d, which is ssh's — the "+ + "one port a machine may never lose", m.Module, port, at) + } out[port] = at } } + // One holder per machine port, within the module too. + holder := map[int]int{} + for port, at := range out { + if other, twice := holder[at]; twice { + return nil, fmt.Errorf("%s gives machine port %d to both its %d and its %d", m.Module, + at, min(port, other), max(port, other)) + } + holder[at] = port + } if len(out) == 0 { return nil, nil } diff --git a/internal/inventory/catalogue.go b/internal/inventory/catalogue.go index 322de46..65a6e5d 100644 --- a/internal/inventory/catalogue.go +++ b/internal/inventory/catalogue.go @@ -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 { diff --git a/internal/inventory/ports_test.go b/internal/inventory/ports_test.go index b62f310..21a3761 100644 --- a/internal/inventory/ports_test.go +++ b/internal/inventory/ports_test.go @@ -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) + } +}