Make pending assignments safe to race, settle them on a tick, and say only what was checked

Review of #150: a withdrawal could land between the look and the act, a
failed ask read as a build in flight, a request kept the wrong asker, a
build being registered read as not built, build "true" could ask a build
nothing waited on, and a status read changed state. Claim a row under the
machine's hold before making it, keep a request only once asked, settle on
the controller's own tick, raise an assignment not made as a condition
until it is answered, and tie each row to its machine.
This commit is contained in:
jochen
2026-10-08 16:38:11 +02:00
parent 7d63d2e68c
commit 249d97d1c8
12 changed files with 1018 additions and 235 deletions
@@ -1,15 +1,15 @@
-- An assignment waits for its module's build (novox/hq issue 325).
-- An assignment waits for its module's build (novox/hq issue 325, ADR 0261).
--
-- A merge that adds a module asks for its build (issue 300), and the module is registered when the build's
-- outcome is taken in, minutes later. An `assign` in between was answered "no module of that name", which
-- read as "nobody registered it". The controller now keeps every build it asks for, so `assign` can tell a
-- build in flight from a module it never heard of, and an assignment made while the build runs is kept as
-- pending and made when the build registers the module, or ended with the reason when it does not.
-- pending and made by the controller when the build registers the module, or ended with the reason.
-- Every build the controller asked for, whatever asked it: a merge, a plan's tier, a person's `build`, an
-- `assign` with build. Its outcome is the `build` row of the same id, written when the outcome is taken in;
-- an ask without one is still running, or was lost. not_asked is why a merge's ask of a new module could
-- not be made, and the id is then not a build's.
-- Every build the controller asked for: a merge's new module, a plan's tier, a person's `build`, an `assign`
-- with build. Kept once the ask is made, never before. Its outcome is the `build` row of the same id;
-- an ask without one is still running, or was lost. not_asked is why the ask could not be made or waited
-- for; for a merge that could not ask, the id is the controller's and no build's.
create table build_request (
id text primary key,
repository text not null,
@@ -23,20 +23,27 @@ create table build_request (
);
create index build_request_asked_at on build_request (asked_at);
-- An assignment kept until its module is registered. waiting until the build registers the module, then
-- applied (the assignment was made), refused (the assignment was refused then, with the refusal), expired
-- (the build failed, was not registered, or said nothing within its bound) or withdrawn (`unassign`).
-- An assignment kept until its module is registered (ADR 0261). waiting until the build registers the
-- module; applying while the controller makes it, under the machine's hold; then applied, refused (the
-- assignment was refused then), expired (the build failed, was not registered, or said nothing within its
-- bound) or withdrawn (`unassign`). raised_at is when an expiry or refusal was raised as a condition,
-- acknowledged_at when a person took it back with `unassign`, cleared_at when its condition was cleared.
-- Gone with its machine. Ended rows are deleted after 30 days.
create table pending_assignment (
id bigserial primary key,
node text not null,
module text not null,
build_id text not null,
repository text not null default '',
source_path text not null default '',
since timestamptz not null default now(),
state text not null default 'waiting'
check (state in ('waiting', 'applied', 'refused', 'expired', 'withdrawn')),
note text not null default '',
settled_at timestamptz
id bigserial primary key,
node uuid not null references node (id) on delete cascade,
module text not null,
build_id text not null,
repository text not null default '',
source_path text not null default '',
since timestamptz not null default now(),
state text not null default 'waiting'
check (state in ('waiting', 'applying', 'applied', 'refused', 'expired', 'withdrawn')),
note text not null default '',
settled_at timestamptz,
raised_at timestamptz,
acknowledged_at timestamptz,
cleared_at timestamptz
);
create unique index pending_assignment_waiting on pending_assignment (node, module) where state = 'waiting';
create unique index pending_assignment_open on pending_assignment (node, module)
where state in ('waiting', 'applying');
+190 -64
View File
@@ -10,36 +10,36 @@ import (
"github.com/jackc/pgx/v5"
)
// What the controller asked to be built, and the assignments waiting for it (novox/hq issue 325).
// 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, an ask lets `assign` say which case it is in,
// and an assignment made while the build runs is kept until the build registers the 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.
// requestKeptFor is how long an ask is kept: long past any build's bound, short enough to stay small.
const requestKeptFor = 30 * 24 * time.Hour
// 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 an ask that could not be made, an id of
// the controller's making that no build carries.
// 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 what asked: "merge", "build", "plan", "assign".
// For says who asked: "merge", "plan", "build" (a person), "assign".
For string
// NotAsked is why the ask could not be made, empty when it was.
// NotAsked is why the ask could not be made, or not waited for; empty when it was.
NotAsked string
At time.Time
}
// Name is the module this ask 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.
// 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)
@@ -47,19 +47,19 @@ func (a BuildRequest) Name() string {
return strings.TrimSuffix(path.Base(strings.TrimRight(a.Repository, "/")), ".git")
}
// RequestOutcome is an ask beside what came of it, where anything did.
// 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 it registered as, HeardAt when it was taken in.
// 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 ask. Idempotent on the id: a merge's ask recorded by the merge and by the
// asking itself is one ask. Asks older than requestKeptFor are let go in the same act.
// 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() {
@@ -72,17 +72,25 @@ func (i *Inventory) RecordBuildRequest(ctx context.Context, a BuildRequest) erro
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)
values ($1, $2, $3, $4, $5, $6, $7, $8, $9)
on conflict (id) do update set
commit_hash = case when build_request.commit_hash = '' then excluded.commit_hash else build_request.commit_hash end,
asked_for = case when build_request.asked_for = '' then excluded.asked_for else build_request.asked_for end`,
on conflict (id) do nothing`,
a.ID, a.Repository, a.Seat, strings.Trim(a.Path, "/"), a.Ref, a.Commit, a.For, notAsked, at); err != nil {
return err
}
_, err := i.store.Pool().Exec(ctx, `delete from build_request where asked_at < $1`, time.Now().Add(-requestKeptFor))
_, err := i.store.Pool().Exec(ctx, `delete from build_request where asked_at < $1`, time.Now().Add(-KeptFor))
return err
}
// RequestsNamed is every ask kept whose directory names the module, newest first, each with its outcome.
// MarkNotAsked says a build request was asked and its asker could not hand it over, or stopped waiting for
// it: 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
}
// 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 {
@@ -97,8 +105,8 @@ func (i *Inventory) RequestsNamed(ctx context.Context, name string) ([]RequestOu
return out, nil
}
// RequestedNames is the name of every module an ask kept is expected to register: what `assign` offers as a
// closest name beside the modules registered.
// 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 {
@@ -141,13 +149,14 @@ func (i *Inventory) requests(ctx context.Context) ([]RequestOutcome, error) {
// 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.
// PendingAssignment is an assignment kept until its module is registered (ADR 0261).
type PendingAssignment struct {
ID int64
Node string
@@ -159,22 +168,48 @@ type PendingAssignment struct {
// 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")
// RecordPending keeps an assignment until its module is registered. One waits per machine and module.
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) {
if _, err := i.NodeByName(ctx, p.Node); err != nil {
node, err := i.NodeByName(ctx, p.Node)
if err != nil {
return PendingAssignment{}, err
}
err := i.store.Pool().QueryRow(ctx,
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 = 'waiting' do nothing
on conflict (node, module) where state in ('waiting', 'applying') do nothing
returning id, since`,
p.Node, p.Module, p.Build, p.Repository, strings.Trim(p.Path, "/")).Scan(&p.ID, &p.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
}
@@ -185,44 +220,135 @@ func (i *Inventory) RecordPending(ctx context.Context, p PendingAssignment) (Pen
return p, nil
}
// Pending is every pending assignment still waiting, and every one ended since the time given, oldest
// first. A zero time is the waiting ones alone.
func (i *Inventory) Pending(ctx context.Context, endedSince time.Time) ([]PendingAssignment, error) {
rows, err := i.store.Pool().Query(ctx,
`select id, node, module, build_id, repository, source_path, since, state, note, settled_at
from pending_assignment
where state = 'waiting' or (settled_at is not null and settled_at >= $1)
order by since, id`, endedSince)
if err != nil {
return nil, err
}
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); err != nil {
return nil, err
}
if endedSince.IsZero() && p.State != PendingWaiting {
continue
}
out = append(out, p)
}
return out, rows.Err()
}
// SettlePending ends a waiting pending assignment. False when it was no longer waiting: another process
// settled it first, and its word stands.
func (i *Inventory) SettlePending(ctx context.Context, id int64, state, note string) (bool, error) {
if state == PendingWaiting {
return false, errors.New("a pending assignment is settled to applied, refused, expired or withdrawn")
}
// 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 state = $2, note = $3, settled_at = now()
where id = $1 and state = 'waiting'`, id, state, note)
`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.
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`,
now.Add(-KeptFor))
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
+25
View File
@@ -0,0 +1,25 @@
package inventory
import (
"testing"
"time"
)
// A pending assignment goes with its machine (novox/hq issue 325, review of mesh-controller#150 point 7).
func TestAPendingAssignmentGoesWithItsMachine(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
if _, err := inv.AddNode(ctx, "leaving"); err != nil {
t.Fatal(err)
}
if _, err := inv.RecordPending(ctx, PendingAssignment{Node: "leaving", Module: "sensors", Build: "b-1"}); err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil {
t.Fatalf("a machine with a pending assignment could not be removed: %v", err)
}
rows, err := inv.Pending(ctx, time.Unix(0, 0))
if err != nil || len(rows) != 0 {
t.Fatalf("a pending assignment outlived its machine: %+v %v", rows, err)
}
}