mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery covered: a later merge that contains it was delivered: novox/mesh-controller@eca6390d6fe4 (merged as 33dc85d8 into main, walk plan-17915444…
A source repository's name is not its identity, and a trunk anyone may push to makes the trunk rule mean nothing (the review of 2026-10-09). Through any verb a module now registers only from a repository on the mesh's forge whose trunk refuses direct pushes, requires a status and lets no administrator merge past one, asked of the forge's own tools; and only from the repository by the forge's id, recorded at registration (migration 0084), so one deleted and made again under the name is refused. The serving controller marks its environment, so nothing it runs or starts reads as the terminal, and a terminal request covers only the repository and path it asked.
406 lines
16 KiB
Go
406 lines
16 KiB
Go
package inventory
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"path"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
)
|
|
|
|
// What the controller asked to be built, and the assignments waiting for it (novox/hq issue 325, ADR 0261).
|
|
//
|
|
// A build's outcome was kept, and its asking was not: between a merge asking for a new module and the
|
|
// outcome registering it, the mesh had no record that anything was coming, and `assign` answered "no module
|
|
// of that name" for a module minutes from existing. Kept here, a build request lets `assign` say which case
|
|
// it is in, and an assignment made while the build runs is kept until the build registers the module.
|
|
|
|
// KeptFor is how long a build request, and an ended pending assignment, is kept.
|
|
const KeptFor = 30 * 24 * time.Hour
|
|
|
|
// BuildRequest is one build the controller asked for.
|
|
type BuildRequest struct {
|
|
// ID is the build's id, as its outcome will carry it; for a merge whose ask could not be made, an id
|
|
// of the controller's making that no build carries.
|
|
ID string
|
|
// Repository and Seat are the source as the mesh spells it (novox/hq ADR 0111): a path on the seat's
|
|
// holder, or a URL with no seat. Path is the module's directory in it, Ref the branch or commit asked.
|
|
Repository, Seat, Path, Ref string
|
|
// Commit is the commit that made the ask, where one did: a merge's.
|
|
Commit string
|
|
// For says who asked: "merge", "plan", "build" (a person), "assign".
|
|
For string
|
|
// NotAsked is why the ask could not be handed over; empty when it was.
|
|
NotAsked string
|
|
// OutcomeUnknown is why an asker that handed the build over stopped waiting for it: asked, outcome
|
|
// unknown. Read as in flight until its outcome or its bound.
|
|
OutcomeUnknown string
|
|
At time.Time
|
|
// AtTerminal says the request was asked at the controller's terminal, never through a verb (novox/hq ADR
|
|
// 0266): what lets its outcome register a module from a repository the catalogue does not build it from.
|
|
AtTerminal bool
|
|
}
|
|
|
|
// Name is the module this request is expected to register, read from its directory: the last element of
|
|
// its path, or of its repository for a module at the repository's root. A guess until the build says, and
|
|
// the one a person reads too: `modules/sensors` is sensors.
|
|
func (a BuildRequest) Name() string {
|
|
if p := strings.Trim(a.Path, "/"); p != "" && p != "." {
|
|
return path.Base(p)
|
|
}
|
|
return strings.TrimSuffix(path.Base(strings.TrimRight(a.Repository, "/")), ".git")
|
|
}
|
|
|
|
// RequestOutcome is a build request beside what came of it, where anything did.
|
|
type RequestOutcome struct {
|
|
BuildRequest
|
|
// Heard is whether an outcome of this id was taken in; Failed is its builder's words when it failed,
|
|
// Module what the build said it built, HeardAt when it was taken in.
|
|
Heard bool
|
|
Failed string
|
|
Module string
|
|
HeardAt time.Time
|
|
}
|
|
|
|
// RecordBuildRequest keeps one build request, once it is asked. Idempotent on the id. Requests older than
|
|
// KeptFor are deleted in the same act.
|
|
func (i *Inventory) RecordBuildRequest(ctx context.Context, a BuildRequest) error {
|
|
at := a.At
|
|
if at.IsZero() {
|
|
at = time.Now().UTC()
|
|
}
|
|
var notAsked *string
|
|
if a.NotAsked != "" {
|
|
notAsked = &a.NotAsked
|
|
}
|
|
if _, err := i.store.Pool().Exec(ctx,
|
|
`insert into build_request (id, repository, seat, source_path, ref, commit_hash, asked_for, not_asked, asked_at,
|
|
at_terminal)
|
|
values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
|
|
on conflict (id) do nothing`,
|
|
a.ID, a.Repository, a.Seat, strings.Trim(a.Path, "/"), a.Ref, a.Commit, a.For, notAsked, at,
|
|
a.AtTerminal); err != nil {
|
|
return err
|
|
}
|
|
_, err := i.store.Pool().Exec(ctx, `delete from build_request where asked_at < $1`, time.Now().Add(-KeptFor))
|
|
return err
|
|
}
|
|
|
|
// BuildRequestByID is the build request kept under this id, and whether one was kept.
|
|
func (i *Inventory) BuildRequestByID(ctx context.Context, id string) (BuildRequest, bool, error) {
|
|
var a BuildRequest
|
|
err := i.store.Pool().QueryRow(ctx,
|
|
`select id, repository, seat, source_path, ref, commit_hash, asked_for, at_terminal, asked_at
|
|
from build_request where id = $1`, id).Scan(&a.ID, &a.Repository, &a.Seat, &a.Path, &a.Ref, &a.Commit,
|
|
&a.For, &a.AtTerminal, &a.At)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return BuildRequest{}, false, nil
|
|
}
|
|
return a, err == nil, err
|
|
}
|
|
|
|
// MarkNotAsked says a kept build request was never handed over: the words are kept, unless its outcome was
|
|
// heard first.
|
|
func (i *Inventory) MarkNotAsked(ctx context.Context, id, why string) error {
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`update build_request set not_asked = $2
|
|
where id = $1 and not exists (select 1 from build where build.id = $1)`, id, why)
|
|
return err
|
|
}
|
|
|
|
// MarkOutcomeUnknown says the asker of a build it handed over stopped waiting for it: asked, outcome unknown.
|
|
// The build may still run; its outcome, when heard, is the last word.
|
|
func (i *Inventory) MarkOutcomeUnknown(ctx context.Context, id, why string) error {
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`update build_request set outcome_unknown = $2
|
|
where id = $1 and not exists (select 1 from build where build.id = $1)`, id, why)
|
|
return err
|
|
}
|
|
|
|
// RequestsNamed is every build request kept whose directory names the module, or whose build said it built
|
|
// it, newest first, each with its outcome.
|
|
func (i *Inventory) RequestsNamed(ctx context.Context, name string) ([]RequestOutcome, error) {
|
|
all, err := i.requests(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var out []RequestOutcome
|
|
for _, a := range all {
|
|
if a.Name() == name || a.Module == name {
|
|
out = append(out, a)
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// RequestedNames is the name of every module a kept build request is expected to register: what `assign`
|
|
// offers as a closest name beside the modules registered.
|
|
func (i *Inventory) RequestedNames(ctx context.Context) ([]string, error) {
|
|
all, err := i.requests(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
seen := map[string]bool{}
|
|
var out []string
|
|
for _, a := range all {
|
|
if n := a.Name(); n != "" && !seen[n] {
|
|
seen[n] = true
|
|
out = append(out, n)
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (i *Inventory) requests(ctx context.Context) ([]RequestOutcome, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select a.id, a.repository, a.seat, a.source_path, a.ref, a.commit_hash, a.asked_for,
|
|
coalesce(a.not_asked, ''), coalesce(a.outcome_unknown, ''), a.asked_at,
|
|
b.id is not null, coalesce(b.failed, ''), coalesce(b.module, ''), coalesce(b.at, a.asked_at)
|
|
from build_request a left join build b on b.id = a.id
|
|
order by a.asked_at desc, a.id desc`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []RequestOutcome
|
|
for rows.Next() {
|
|
var a RequestOutcome
|
|
if err := rows.Scan(&a.ID, &a.Repository, &a.Seat, &a.Path, &a.Ref, &a.Commit, &a.For, &a.NotAsked,
|
|
&a.OutcomeUnknown, &a.At, &a.Heard, &a.Failed, &a.Module, &a.HeardAt); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// The states of a pending assignment.
|
|
const (
|
|
PendingWaiting = "waiting"
|
|
PendingApplying = "applying"
|
|
PendingApplied = "applied"
|
|
PendingRefused = "refused"
|
|
PendingExpired = "expired"
|
|
PendingWithdrawn = "withdrawn"
|
|
)
|
|
|
|
// PendingAssignment is an assignment kept until its module is registered (ADR 0261).
|
|
type PendingAssignment struct {
|
|
ID int64
|
|
Node string
|
|
Module string
|
|
// Build is the build it waits for; Repository and Path where that build is from.
|
|
Build, Repository, Path string
|
|
Since time.Time
|
|
State string
|
|
// Note is what ended it, in the words said when it ended.
|
|
Note string
|
|
Settled *time.Time
|
|
// Raised is when an expiry or refusal was raised as a condition; Acknowledged when a person took it
|
|
// back; Cleared when its condition was cleared.
|
|
Raised, Acknowledged, Cleared *time.Time
|
|
}
|
|
|
|
// Open is whether the pending assignment is still to be made.
|
|
func (p PendingAssignment) Open() bool {
|
|
return p.State == PendingWaiting || p.State == PendingApplying
|
|
}
|
|
|
|
// ErrAlreadyPending is a second pending assignment of one module to one machine.
|
|
var ErrAlreadyPending = errors.New("that assignment is already waiting for its build")
|
|
|
|
const pendingColumns = `p.id, n.name, p.module, p.build_id, p.repository, p.source_path, p.since, p.state, p.note,
|
|
p.settled_at, p.raised_at, p.acknowledged_at, p.cleared_at`
|
|
|
|
func scanPending(rows pgx.Rows) ([]PendingAssignment, error) {
|
|
defer rows.Close()
|
|
var out []PendingAssignment
|
|
for rows.Next() {
|
|
var p PendingAssignment
|
|
if err := rows.Scan(&p.ID, &p.Node, &p.Module, &p.Build, &p.Repository, &p.Path, &p.Since, &p.State,
|
|
&p.Note, &p.Settled, &p.Raised, &p.Acknowledged, &p.Cleared); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, p)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// RecordPending keeps an assignment until its module is registered. One is open per machine and module.
|
|
func (i *Inventory) RecordPending(ctx context.Context, p PendingAssignment) (PendingAssignment, error) {
|
|
node, err := i.NodeByName(ctx, p.Node)
|
|
if err != nil {
|
|
return PendingAssignment{}, err
|
|
}
|
|
err = i.store.Pool().QueryRow(ctx,
|
|
`insert into pending_assignment (node, module, build_id, repository, source_path)
|
|
values ($1, $2, $3, $4, $5)
|
|
on conflict (node, module) where state in ('waiting', 'applying') do nothing
|
|
returning id, since`,
|
|
node.ID, p.Module, p.Build, p.Repository, strings.Trim(p.Path, "/")).Scan(&p.ID, &p.Since)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return PendingAssignment{}, ErrAlreadyPending
|
|
}
|
|
if err != nil {
|
|
return PendingAssignment{}, err
|
|
}
|
|
p.State = PendingWaiting
|
|
return p, nil
|
|
}
|
|
|
|
// RepointPending makes a waiting pending assignment wait for another build. False when it no longer waits.
|
|
func (i *Inventory) RepointPending(ctx context.Context, id int64, build, repository, sourcePath string) (bool, error) {
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`update pending_assignment set build_id = $2, repository = $3, source_path = $4
|
|
where id = $1 and state = 'waiting'`, id, build, repository, strings.Trim(sourcePath, "/"))
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return tag.RowsAffected() == 1, nil
|
|
}
|
|
|
|
// Pending is every pending assignment still waiting, oldest first, when endedSince is zero. Otherwise it is
|
|
// every open one and every one ended since then.
|
|
func (i *Inventory) Pending(ctx context.Context, endedSince time.Time) ([]PendingAssignment, error) {
|
|
if endedSince.IsZero() {
|
|
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
|
|
from pending_assignment p join node n on n.id = p.node
|
|
where p.state = 'waiting' order by p.since, p.id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return scanPending(rows)
|
|
}
|
|
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
|
|
from pending_assignment p join node n on n.id = p.node
|
|
where p.state in ('waiting', 'applying') or p.settled_at >= $1
|
|
order by p.since, p.id`, endedSince)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return scanPending(rows)
|
|
}
|
|
|
|
// PendingFor is every pending assignment of one module to one machine, newest first.
|
|
func (i *Inventory) PendingFor(ctx context.Context, node, module string) ([]PendingAssignment, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
|
|
from pending_assignment p join node n on n.id = p.node
|
|
where n.name = $1 and p.module = $2 order by p.since desc, p.id desc`, node, module)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return scanPending(rows)
|
|
}
|
|
|
|
// Stuck is every pending assignment left applying since before the time given: a controller that claimed
|
|
// it and stopped before it said what came of it.
|
|
func (i *Inventory) Stuck(ctx context.Context, before time.Time) ([]PendingAssignment, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
|
|
from pending_assignment p join node n on n.id = p.node
|
|
where p.state = 'applying' and p.settled_at < $1 order by p.id`, before)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return scanPending(rows)
|
|
}
|
|
|
|
// ToRaise is every pending assignment that expired or was refused and has not been raised as a condition;
|
|
// ToClear every one raised and not cleared.
|
|
func (i *Inventory) ToRaise(ctx context.Context) ([]PendingAssignment, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
|
|
from pending_assignment p join node n on n.id = p.node
|
|
where p.state in ('expired', 'refused') and p.raised_at is null order by p.id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return scanPending(rows)
|
|
}
|
|
|
|
func (i *Inventory) ToClear(ctx context.Context) ([]PendingAssignment, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
|
|
from pending_assignment p join node n on n.id = p.node
|
|
where p.raised_at is not null and p.cleared_at is null order by p.id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return scanPending(rows)
|
|
}
|
|
|
|
// MarkPending stamps a pending assignment raised, acknowledged or cleared, now.
|
|
func (i *Inventory) MarkPending(ctx context.Context, id int64, what string) error {
|
|
column := map[string]string{"raised": "raised_at", "acknowledged": "acknowledged_at", "cleared": "cleared_at"}[what]
|
|
if column == "" {
|
|
return errors.New("a pending assignment is marked raised, acknowledged or cleared")
|
|
}
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`update pending_assignment set `+column+` = coalesce(`+column+`, now()) where id = $1`, id)
|
|
return err
|
|
}
|
|
|
|
// ClaimPending takes a waiting pending assignment for making, so nothing else settles or withdraws it while
|
|
// it is made. False when it no longer waits: another process took it, or it was withdrawn.
|
|
func (i *Inventory) ClaimPending(ctx context.Context, id int64) (bool, error) {
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`update pending_assignment set state = 'applying', settled_at = now() where id = $1 and state = 'waiting'`, id)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return tag.RowsAffected() == 1, nil
|
|
}
|
|
|
|
// ReleaseClaim puts a claimed pending assignment back to waiting.
|
|
func (i *Inventory) ReleaseClaim(ctx context.Context, id int64) error {
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`update pending_assignment set state = 'waiting', settled_at = null where id = $1 and state = 'applying'`, id)
|
|
return err
|
|
}
|
|
|
|
// SettlePending ends a pending assignment in the state it is read in (from): waiting for an expiry or a
|
|
// withdrawal, applying for the outcome of making it. False when it was no longer in that state.
|
|
func (i *Inventory) SettlePending(ctx context.Context, id int64, from, state, note string) (bool, error) {
|
|
if state == PendingWaiting || state == PendingApplying {
|
|
return false, errors.New("a pending assignment is settled to applied, refused, expired or withdrawn")
|
|
}
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`update pending_assignment set state = $3, note = $4, settled_at = now()
|
|
where id = $1 and state = $2`, id, from, state, note)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return tag.RowsAffected() == 1, nil
|
|
}
|
|
|
|
// ForgetEndedPending deletes ended pending assignments settled more than KeptFor ago, and says how many. One
|
|
// raised as a condition is kept until that condition is cleared: deleting it would leave the condition
|
|
// with nothing to answer it.
|
|
func (i *Inventory) ForgetEndedPending(ctx context.Context, now time.Time) (int64, error) {
|
|
tag, err := i.store.Pool().Exec(ctx,
|
|
`delete from pending_assignment where state not in ('waiting', 'applying') and settled_at < $1
|
|
and (raised_at is null or cleared_at is not null)`,
|
|
now.Add(-KeptFor))
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return tag.RowsAffected(), nil
|
|
}
|
|
|
|
// PendingKnown says which of these pending assignments are still on record.
|
|
func (i *Inventory) PendingKnown(ctx context.Context, ids []int64) (map[int64]bool, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select id from pending_assignment where id = any($1)`, ids)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := map[int64]bool{}
|
|
for rows.Next() {
|
|
var id int64
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
out[id] = true
|
|
}
|
|
return out, rows.Err()
|
|
}
|