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.
150 lines
5.6 KiB
Go
150 lines
5.6 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"
|
|
)
|
|
|
|
// DurationKinds are every kind, in the order `durations` shows them.
|
|
var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild}
|
|
|
|
// 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()
|
|
}
|
|
|
|
// 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
|
|
}
|