A merge asks for a new module's build, and assign answered "no module of that name" until the build registered it, which read as a module nobody registered (hq issue 325). Keep every build request, tell a build in flight, a module known and not built, and an unknown name apart, and make an assignment made while the build runs when the build registers the module.
229 lines
8.2 KiB
Go
229 lines
8.2 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).
|
|
//
|
|
// 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.
|
|
|
|
// 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
|
|
|
|
// 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 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 string
|
|
// NotAsked is why the ask could not be made, 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.
|
|
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 an ask 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.
|
|
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.
|
|
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)
|
|
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`,
|
|
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))
|
|
return err
|
|
}
|
|
|
|
// RequestsNamed is every ask kept whose directory names the module, 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 an ask kept 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, ''), 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.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"
|
|
PendingApplied = "applied"
|
|
PendingRefused = "refused"
|
|
PendingExpired = "expired"
|
|
PendingWithdrawn = "withdrawn"
|
|
)
|
|
|
|
// PendingAssignment is an assignment kept until its module is registered.
|
|
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
|
|
}
|
|
|
|
// 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.
|
|
func (i *Inventory) RecordPending(ctx context.Context, p PendingAssignment) (PendingAssignment, error) {
|
|
if _, err := i.NodeByName(ctx, p.Node); 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 = 'waiting' do nothing
|
|
returning id, since`,
|
|
p.Node, 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
|
|
}
|
|
|
|
// 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")
|
|
}
|
|
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)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return tag.RowsAffected() == 1, nil
|
|
}
|