mesh/delivery superseded: a newer head of the same pull request
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 fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.62s)
The operator's budget (a leaf module running everywhere within five minutes of its merge, a core module within ten) cannot be held to without knowing where a walk's time goes. A walk now keeps when its batch's window closed, when it was cut and its class (migration 0093), each gate reading, and what each machine of the rest was sent; 'delivery walks' and plan-moved say its phases from the merge to every machine of the rest reporting the build applied, and ended walks keep them as walk-phase durations per class. Probe D16 reads mesh-delivery's 'times' and raises delivery.<class>.over-budget. A phase that cannot be measured is said unknown, never zero. Measurement only: no walk is held, sent or judged differently.
161 lines
6.3 KiB
Go
161 lines
6.3 KiB
Go
package inventory
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
)
|
|
|
|
// The durations the core's bounds are set from (novox/hq to-be 45 Phase 0): see migration 0066.
|
|
|
|
// The kinds of duration recorded.
|
|
const (
|
|
DurationApply = "apply"
|
|
DurationHeartbeatGap = "heartbeat-gap"
|
|
DurationPlanTier = "plan-tier"
|
|
DurationBuild = "build"
|
|
// DurationWalkPhase is one phase of an ended walk, per class (novox/hq ADR 0282 decision 6): subject
|
|
// "<class>/<phase>", and "<class>/total" for its merge to every machine running it.
|
|
DurationWalkPhase = "walk-phase"
|
|
)
|
|
|
|
// DurationKinds are every kind, in the order `durations` shows them.
|
|
var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild, DurationWalkPhase}
|
|
|
|
// DurationsKeptFor is how long a duration is kept: long enough to set a bound from, and to correct it
|
|
// in Phase 1's first live week.
|
|
const DurationsKeptFor = 30 * 24 * time.Hour
|
|
|
|
// Duration is one measurement.
|
|
type Duration struct {
|
|
Kind string `json:"kind"`
|
|
Subject string `json:"subject"`
|
|
Node string `json:"node,omitempty"`
|
|
Ref string `json:"ref"`
|
|
Started time.Time `json:"started"`
|
|
Took time.Duration `json:"took"`
|
|
Detail string `json:"detail,omitempty"`
|
|
}
|
|
|
|
// RecordDuration keeps one measurement, once: the same thing measured again is not a second row.
|
|
func (i *Inventory) RecordDuration(ctx context.Context, d Duration) error {
|
|
if d.Took < 0 {
|
|
return nil // a clock that went back measures nothing
|
|
}
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`insert into duration (kind, subject, node, ref, started, took_ms, detail)
|
|
values ($1, $2, $3, $4, $5, $6, $7) on conflict do nothing`,
|
|
d.Kind, d.Subject, d.Node, d.Ref, d.Started, d.Took.Milliseconds(), d.Detail)
|
|
return err
|
|
}
|
|
|
|
// RecordApplyDuration measures a machine's report of the declaration it was last sent: from the send
|
|
// to the first report of it. A report of anything else, or of a send already measured, measures
|
|
// nothing — a machine reconciling reports the same declaration every few minutes.
|
|
func (i *Inventory) RecordApplyDuration(ctx context.Context, node, declared, outcome string) error {
|
|
if declared == "" {
|
|
return nil
|
|
}
|
|
_, err := i.store.Pool().Exec(ctx,
|
|
`insert into duration (kind, subject, node, ref, started, took_ms, detail)
|
|
select $1, n.name, n.name, n.sent || '@' || to_char(n.sent_at at time zone 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US'),
|
|
n.sent_at, (extract(epoch from (now() - n.sent_at)) * 1000)::bigint, $4
|
|
from node n
|
|
where n.name = $2 and n.sent = $3 and n.sent_at is not null
|
|
on conflict do nothing`,
|
|
DurationApply, node, declared, outcome)
|
|
return err
|
|
}
|
|
|
|
// RecordHeartbeatGap measures the silence before a machine's word: from the last word heard before
|
|
// it, which the caller read before recording this one.
|
|
func (i *Inventory) RecordHeartbeatGap(ctx context.Context, node string, before time.Time) error {
|
|
if before.IsZero() {
|
|
return nil
|
|
}
|
|
now := time.Now()
|
|
return i.RecordDuration(ctx, Duration{Kind: DurationHeartbeatGap, Subject: node, Node: node,
|
|
Ref: before.UTC().Format(time.RFC3339Nano), Started: before, Took: now.Sub(before)})
|
|
}
|
|
|
|
// Durations is every measurement of a kind since a moment, oldest first; every kind when kind is empty.
|
|
func (i *Inventory) Durations(ctx context.Context, kind string, since time.Time) ([]Duration, error) {
|
|
rows, err := i.store.Pool().Query(ctx,
|
|
`select kind, subject, node, ref, started, took_ms, detail from duration
|
|
where ($1 = '' or kind = $1) and recorded >= $2 order by recorded`, kind, since)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []Duration
|
|
for rows.Next() {
|
|
var d Duration
|
|
var ms int64
|
|
if err := rows.Scan(&d.Kind, &d.Subject, &d.Node, &d.Ref, &d.Started, &ms, &d.Detail); err != nil {
|
|
return nil, err
|
|
}
|
|
d.Took = time.Duration(ms) * time.Millisecond
|
|
out = append(out, d)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// WalkPhasesKept says whether any phase of a walk is kept as a walk-phase duration, and whether its total is.
|
|
func (i *Inventory) WalkPhasesKept(ctx context.Context, plan string) (anyKept, totalKept bool, err error) {
|
|
err = i.store.Pool().QueryRow(ctx,
|
|
`select count(*) > 0, count(*) filter (where ref = $2) > 0 from duration where kind = $1 and ref like $3`,
|
|
DurationWalkPhase, plan+"/total", plan+"/%").Scan(&anyKept, &totalKept)
|
|
return anyKept, totalKept, err
|
|
}
|
|
|
|
// ForgetOldDurations removes what is older than DurationsKeptFor, and says how many.
|
|
func (i *Inventory) ForgetOldDurations(ctx context.Context) (int64, error) {
|
|
tag, err := i.store.Pool().Exec(ctx, `delete from duration where recorded < $1`,
|
|
time.Now().Add(-DurationsKeptFor))
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return tag.RowsAffected(), nil
|
|
}
|
|
|
|
// planTierLeft is the measurement a plan's save makes when it leaves a tier: the tier moved on, or
|
|
// the plan ended. Read in the save's own transaction, so two saves cannot both measure one tier.
|
|
func planTierLeft(ctx context.Context, tx pgx.Tx, p Plan, now time.Time) (entered time.Time, err error) {
|
|
var oldTier int
|
|
var oldState string
|
|
var since time.Time
|
|
err = tx.QueryRow(ctx,
|
|
`select tier, state, coalesce(tier_entered, created) from release_plan where id = $1 for update`,
|
|
p.ID).Scan(&oldTier, &oldState, &since)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return now, nil // a new plan enters its first tier now
|
|
}
|
|
if err != nil {
|
|
return time.Time{}, err
|
|
}
|
|
if oldState == PlanAssembling || oldState == PlanQueued {
|
|
return now, nil // a batch cut into a walk enters its first tier now (novox/hq ADR 0276)
|
|
}
|
|
wasOpen := oldState == PlanBuilding || oldState == PlanRolling
|
|
if !wasOpen || (oldTier == p.Tier && p.Open()) {
|
|
return since, nil // still in the tier, or already ended
|
|
}
|
|
var modules []string
|
|
if oldTier >= 0 && oldTier < len(p.Tiers) {
|
|
modules = p.Tiers[oldTier]
|
|
}
|
|
detail := fmt.Sprintf("plan %s, tier %d of %d (%v), left %s", p.ID, oldTier, len(p.Tiers), modules, p.State)
|
|
if p.Open() {
|
|
detail = fmt.Sprintf("plan %s, tier %d of %d (%v), moved on to tier %d", p.ID, oldTier, len(p.Tiers), modules, p.Tier)
|
|
}
|
|
_, err = tx.Exec(ctx,
|
|
`insert into duration (kind, subject, node, ref, started, took_ms, detail)
|
|
values ($1, $2, '', $3, $4, $5, $6) on conflict do nothing`,
|
|
DurationPlanTier, p.Repository, fmt.Sprintf("%s/tier-%d", p.ID, oldTier), since,
|
|
now.Sub(since).Milliseconds(), detail)
|
|
return now, err
|
|
}
|