Files
mesh-controller/internal/inventory/batches.go
T
jochen 96fc4209d3
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
Assemble merges in a rolling window and walk each batch once (hq ADR 0276, issue 362)
Every merge opened a walk and the next merge of the branch superseded it: two
catalogue merges 18 s apart left a walk no delivery held, and the operator
started it by hand 58 minutes later. A merge now joins the open batch, kept in
the store (migration 0089), which is cut into one walk when no merge came for
merge-window (90 s) or at merge-window-at-most (10 min): one commit per
repository, the latest of its branch, with every file the batch's merges
changed. One walk at a time; a started walk is never superseded, a waiting one
is folded into the next. The walk names every merge it answers on the wire
(delivery.merges, taken_over_by, batch). A failed walk walks its earlier merges
alone, newest first, until one is delivered. A delivery group's order becomes
tier edges inside the walk. plans shows the batch assembling; S18 and S19
bound its waits; S16 names the merges a waiting walk answers.
2026-10-10 13:20:04 +02:00

118 lines
4.5 KiB
Go

package inventory
import (
"context"
"encoding/json"
"errors"
"strings"
"time"
"github.com/jackc/pgx/v5"
)
// The merges a batch or walk answers (novox/hq ADR 0276): every merge the controller heard into a branch a
// module follows is kept here once, with the batch or walk that answers it. A merge is never dropped and never
// walked twice: its key is its repository and commit, and it belongs to one plan at a time.
// BatchedMerge is one merge as the controller heard it.
type BatchedMerge struct {
Repository string
Branch string
Commit string
// Merged is when the forge made it: the branch's order. Heard is when the controller heard it: what the
// window is measured from.
Merged time.Time
Heard time.Time
// Event is the forge's announcement, what the cut plans from.
Event json.RawMessage
// Plan is the batch or walk answering it; empty while it waits to be walked alone.
Plan string
// Alone says it is walked on its own commit, after a failed walk that carried it in a later one.
Alone bool
}
// AddMerge keeps a merge for a plan; false when the merge was kept already, for whatever plan.
func (i *Inventory) AddMerge(ctx context.Context, m BatchedMerge) (bool, error) {
tag, err := i.store.Pool().Exec(ctx,
`insert into batched_merge (repository, branch, commit_hash, merged_at, heard_at, event, plan_id, alone)
values (lower($1), $2, $3, $4, $5, $6, nullif($7, ''), $8) on conflict do nothing`,
m.Repository, m.Branch, m.Commit, mergedAt(m.Merged), m.Heard, []byte(m.Event), m.Plan, m.Alone)
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil
}
// MergeOf is one kept merge; false when it was never heard.
func (i *Inventory) MergeOf(ctx context.Context, repository, commit string) (BatchedMerge, bool, error) {
ms, err := i.merges(ctx, `where repository = lower($1) and commit_hash = $2`, repository, commit)
if err != nil || len(ms) == 0 {
return BatchedMerge{}, false, err
}
return ms[0], true, nil
}
// MergesOf is every merge a plan answers, oldest merge first.
func (i *Inventory) MergesOf(ctx context.Context, plan string) ([]BatchedMerge, error) {
return i.merges(ctx, `where plan_id = $1 order by merged_at nulls last, heard_at`, plan)
}
// LaterMergesOf is every merge of a repository's branch made after a moment and answered by a plan, the
// newest first.
func (i *Inventory) LaterMergesOf(ctx context.Context, repository, branch string, after time.Time) ([]BatchedMerge, error) {
return i.merges(ctx, `where repository = lower($1) and branch = $2 and merged_at > $3 and plan_id is not null
order by merged_at desc`, repository, branch, after)
}
// AloneMerges is every merge waiting to be walked on its own commit, the newest first (ADR 0276 decision 3).
func (i *Inventory) AloneMerges(ctx context.Context) ([]BatchedMerge, error) {
return i.merges(ctx, `where plan_id is null and alone order by merged_at desc nulls last, heard_at desc`)
}
// AnswerMerge moves a merge to the plan that answers it now; an empty plan with alone leaves it waiting to
// be walked on its own commit.
func (i *Inventory) AnswerMerge(ctx context.Context, repository, commit, plan string, alone bool) error {
_, err := i.store.Pool().Exec(ctx,
`update batched_merge set plan_id = nullif($3, ''), alone = $4 where repository = lower($1) and commit_hash = $2`,
repository, commit, plan, alone)
return err
}
// Batches is every batch not yet cut: assembling or queued, oldest first. One at a time is the rule; more
// are read so a fault is seen.
func (i *Inventory) Batches(ctx context.Context) ([]Plan, error) {
return i.plans(ctx, `where state in ('assembling', 'queued') order by created`)
}
func (i *Inventory) merges(ctx context.Context, tail string, args ...any) ([]BatchedMerge, error) {
rows, err := i.store.Pool().Query(ctx,
`select repository, branch, commit_hash, merged_at, heard_at, event, coalesce(plan_id, ''), alone
from batched_merge `+tail, args...)
if err != nil {
return nil, err
}
defer rows.Close()
var out []BatchedMerge
for rows.Next() {
var m BatchedMerge
var merged *time.Time
var event []byte
if err := rows.Scan(&m.Repository, &m.Branch, &m.Commit, &merged, &m.Heard, &event, &m.Plan, &m.Alone); err != nil {
return nil, err
}
if merged != nil {
m.Merged = merged.UTC()
}
m.Heard = m.Heard.UTC()
m.Event = event
out = append(out, m)
}
if errors.Is(rows.Err(), pgx.ErrNoRows) {
return nil, nil
}
return out, rows.Err()
}
// SameRepository says two spellings of a repository are one.
func SameRepository(a, b string) bool { return strings.EqualFold(a, b) }