package inventory import ( "context" "encoding/json" "errors" "fmt" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgconn" ) // Which port a machine uses for what a module needs reachable. // // **A module cannot choose this** (novox/hq ADR 0038). It is written once and assigned anywhere, // so any number it picks is a guess about a machine it has never seen. The mesh is the only thing // that knows what else is there, so the mesh chooses — and a module says only that something must // be reachable, and what the software itself calls it. // Assigned is where the machine puts one of a module's ports. type Assigned struct { Module string // Wanted is the port the software uses — what the module wrote down. Wanted int // Machine is where this machine publishes it. Machine int // Fixed means the protocol chose it, not the mesh. Fixed bool } // The range the mesh assigns from. // // High and unprivileged, so an assignment never needs root and never lands on something a person // would recognise. Below the range Linux uses for outgoing connections, so an assignment cannot // collide with a port the kernel handed to something else while the machine was working. const ( firstAssignable = 20000 lastAssignable = 29999 ) // ErrPortTaken is returned when a port the protocol fixes is already held by another module. var ErrPortTaken = errors.New("that port is already held on this machine") // PortFor is where one module's port lives on one machine, choosing it the first time. // // **Kept once chosen.** A port that moved on every declaration would restart both ends each time, // and would hand a consumer a number that was true when it was read — which is the same argument // that makes a credential stable. func (i *Inventory) PortFor( ctx context.Context, node, module string, wanted int, fixed bool, ) (Assigned, error) { record, err := i.NodeByName(ctx, node) if err != nil { return Assigned{}, err } var held Assigned err = i.store.Pool().QueryRow(ctx, `select machine, fixed from port_assignment where node = $1 and module = $2 and wanted = $3`, record.ID, module, wanted).Scan(&held.Machine, &held.Fixed) if err == nil { held.Module, held.Wanted = module, wanted // A module that has become fixed since it was assigned must move to the port the protocol // requires. Said rather than silently kept: a mail system on the port it was lent is a // mail system nothing can deliver to. if fixed && held.Machine != wanted { if err := i.freePort(ctx, record.ID, module, wanted); err != nil { return Assigned{}, err } return i.assignPort(ctx, record.ID, node, module, wanted, true) } return held, nil } if !errors.Is(err, pgx.ErrNoRows) { return Assigned{}, err } return i.assignPort(ctx, record.ID, node, module, wanted, fixed) } // RecordCarried keeps what a machine says it already holds, replacing whatever it said before. // // **Replaced whole, not merged.** A machine that gave a port back must be believed about that too, // and a set the mesh only ever adds to would keep a port reserved for something that is no longer // there. func (i *Inventory) RecordCarried(ctx context.Context, node string, ports []int) error { record, err := i.NodeByName(ctx, node) if err != nil { return err } if ports == nil { ports = []int{} } _, err = i.store.Pool().Exec(ctx, `update node set carried_ports = $2 where id = $1`, record.ID, ports) return err } // carriedOn is what a machine said it already holds. func (i *Inventory) carriedOn(ctx context.Context, nodeID any) ([]int, error) { var ports []int err := i.store.Pool().QueryRow(ctx, `select carried_ports from node where id = $1`, nodeID).Scan(&ports) if err != nil { return nil, err } return ports, nil } func (i *Inventory) freePort(ctx context.Context, node any, module string, wanted int) error { _, err := i.store.Pool().Exec(ctx, `delete from port_assignment where node = $1 and module = $2 and wanted = $3`, node, module, wanted) return err } func (i *Inventory) assignPort( ctx context.Context, nodeID any, node, module string, wanted int, fixed bool, ) (Assigned, error) { taken, err := i.portsOn(ctx, nodeID) if err != nil { return Assigned{}, err } // And what the machine itself says it already holds — the foundation it raised before there // was a mesh to ask (novox/hq ADR 0038). Not assignments: nothing here chose them, and // nothing here can move them. carried, err := i.carriedOn(ctx, nodeID) if err != nil { return Assigned{}, err } for _, port := range carried { if _, mine := taken[port]; !mine { taken[port] = "something this machine already runs" } } // And every port this machine was given for a module (novox/hq ADR 0100): the foundation's // ports, as genesis chose them, are the node's settings and never the mesh's to hand out. given, err := i.givenOn(ctx, nodeID) if err != nil { return Assigned{}, err } for port, by := range given { if _, mine := taken[port]; !mine { taken[port] = by } } machine := wanted if !fixed { // The lowest free one, so a machine's assignments are stable and readable rather than // scattered — and so the same set of modules on two machines gets the same numbers, which // makes a difference between two machines mean something. machine = 0 for candidate := firstAssignable; candidate <= lastAssignable; candidate++ { if _, held := taken[candidate]; !held { machine = candidate break } } if machine == 0 { return Assigned{}, fmt.Errorf( "%s has no free port left between %d and %d, which is ten thousand of them — "+ "something is assigning ports it never gives back", node, firstAssignable, lastAssignable) } } if by, held := taken[machine]; held { return Assigned{}, fmt.Errorf( "%w: %s needs %d and %s already has it on %s. A port the protocol fixes can have one "+ "holder per machine, so one of them has to go somewhere else", ErrPortTaken, module, machine, by, node) } _, err = i.store.Pool().Exec(ctx, `insert into port_assignment (node, module, wanted, machine, fixed) values ($1, $2, $3, $4, $5)`, nodeID, module, wanted, machine, fixed) // Two allocations at once can both pick the same lowest free port; the unique index lets one // through and hands the other a constraint violation in SQL. Said in the mesh's words instead // — and for a port the mesh chose, simply chosen again: the free list has moved, the retry // reads it fresh, and the caller never learns the race happened. var collided *pgconn.PgError if errors.As(err, &collided) && collided.Code == "23505" { // Which race decides what happens next. The table has two keys, so this is one of two // collisions: the racer was *this same assignment* (the primary key), in which case its // answer is the answer — kept-once-chosen does not care who did the choosing — or it was // another module taking the machine port (the unique index), in which case the free list // has moved and an unfixed pick is simply made again. Asking the table tells them apart; // branching on the constraint's name would couple this to the migration's spelling. var held Assigned reread := i.store.Pool().QueryRow(ctx, `select machine, fixed from port_assignment where node = $1 and module = $2 and wanted = $3`, nodeID, module, wanted).Scan(&held.Machine, &held.Fixed) if reread == nil { held.Module, held.Wanted = module, wanted return held, nil } if !errors.Is(reread, pgx.ErrNoRows) { return Assigned{}, reread } if !fixed { return i.assignPort(ctx, nodeID, node, module, wanted, false) } return Assigned{}, fmt.Errorf( "%w: %s needs %d on %s and something else was given it at the same moment — "+ "two assignments raced, and the port the protocol fixes went to the other one", ErrPortTaken, module, wanted, node) } if err != nil { return Assigned{}, err } return Assigned{Module: module, Wanted: wanted, Machine: machine, Fixed: fixed}, nil } /** Which machine ports are spoken for, and by whom. */ func (i *Inventory) portsOn(ctx context.Context, nodeID any) (map[int]string, error) { rows, err := i.store.Pool().Query(ctx, `select machine, module from port_assignment where node = $1`, nodeID) if err != nil { return nil, err } defer rows.Close() out := map[int]string{} for rows.Next() { var machine int var module string if err := rows.Scan(&machine, &module); err != nil { return nil, err } out[machine] = module } return out, rows.Err() } // PortsFor is every assignment a node holds, for composing its declaration. func (i *Inventory) PortsFor(ctx context.Context, node string) ([]Assigned, error) { record, err := i.NodeByName(ctx, node) if err != nil { return nil, err } rows, err := i.store.Pool().Query(ctx, `select module, wanted, machine, fixed from port_assignment where node = $1 order by module, wanted`, record.ID) if err != nil { return nil, err } defer rows.Close() var out []Assigned for rows.Next() { var a Assigned if err := rows.Scan(&a.Module, &a.Wanted, &a.Machine, &a.Fixed); err != nil { return nil, err } out = append(out, a) } return out, rows.Err() } // ReleasePorts gives back everything a module held on a machine, for when it is unassigned. func (i *Inventory) ReleasePorts(ctx context.Context, node, module string) error { record, err := i.NodeByName(ctx, node) if err != nil { return err } _, err = i.store.Pool().Exec(ctx, `delete from port_assignment where node = $1 and module = $2`, record.ID, module) return err } // givenOn is every machine port a module was given on this node by its `ports` setting, and which // module it was given to. func (i *Inventory) givenOn(ctx context.Context, nodeID any) (map[int]string, error) { rows, err := i.store.Pool().Query(ctx, `select module, values->'ports' from settings where node = $1 and jsonb_typeof(values->'ports') = 'object'`, nodeID) if err != nil { return nil, err } defer rows.Close() out := map[int]string{} for rows.Next() { var module string var raw []byte if err := rows.Scan(&module, &raw); err != nil { return nil, err } var given map[string]any if err := json.Unmarshal(raw, &given); err != nil { return nil, err } for _, v := range given { if at, ok := v.(float64); ok { out[int(at)] = module } } } return out, rows.Err() }