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() }