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 }