package inventory import ( "context" "errors" "fmt" "sort" "strings" "time" "github.com/jackc/pgx/v5" ) // Where a walk's time went (novox/hq ADR 0282 decision 6, issue 382): from its merge to every machine of its // rest running its builds, phase by phase, worked out from the walk's own moments. Measurement only: nothing // here decides anything about the walk. // // **A phase that cannot be measured is said unknown, never zero.** Where a moment is missing (a walk kept // before it was recorded, a machine not yet reported), the phases around it are unknown, and the span between // the moments on either side is counted as unknown time: the measured phases and the unknown time always add // up to the total. A phase that did not happen (no window for a merge walked alone, no first machine for a // module nobody runs) is said none, with no time. // The phases, in the order a walk passes them (ADR 0282's table). const ( PhaseWindow = "window" // the merge to its batch's window closing PhaseQueued = "queued" // the window closed to the walk being cut: waiting behind another walk PhaseWord = "word" // the cut to the delivery's word, for a walk that waits for one PhaseBetween = "between-tiers" // one tier's end to the next tier's ask PhaseBuild = "build" // a tier asked to its last module built PhaseSend = "send-first" // built to sent to the first machines PhaseJudge = "judgement" // sent first to the first-node gate's last verdict: its readings PhaseRest = "send-rest" // judged to sent to the rest PhaseApply = "apply" // the last send to every machine of the rest reporting the build applied PhaseOpen = "open" // a walk still running: since its last moment PhaseMeasured = "measured" PhaseUnknown = "unknown" PhaseNone = "none" ) // The classes of a walk (ADR 0282 decision 1) and their delivery budgets. const ( ClassCore = "core" ClassLeaf = "leaf" ) // ApplySilentAfter is how long after a send to the rest a machine that has said nothing is left out of the // walk's end, as a machine not heard from (ADR 0282 decision 1: a sleeping laptop is listed, not waited for). const ApplySilentAfter = 15 * time.Minute // WalkPhases is a walk's time, phase by phase. type WalkPhases struct { Class string `json:"class,omitempty"` // From is the walk's earliest merge; End when its last phase ended: every machine of its rest reported the // build applied. Nil End while that is not known. From *time.Time `json:"from,omitempty"` End *time.Time `json:"end,omitempty"` // TotalMS is End − From; zero while either is unknown, Total its words. TotalMS int64 `json:"total_ms,omitempty"` Total string `json:"total,omitempty"` // UnknownMS is the time within the walk no phase could be measured over. UnknownMS int64 `json:"unknown_ms,omitempty"` Phases []WalkPhase `json:"phases"` Silent []string `json:"silent,omitempty"` Said string `json:"said"` Merges []MergeStart `json:"merges,omitempty"` } // MergeStart is one merge the walk answers and when it was made: where that merge's delivery time starts. type MergeStart struct { Repository string `json:"repository"` Commit string `json:"commit"` Merged time.Time `json:"merged"` } // WalkPhase is one phase of a walk. type WalkPhase struct { Name string `json:"name"` // Tier is the tier a tier's phase belongs to; -1 for a phase of the walk. Tier int `json:"tier"` State string `json:"state"` Start *time.Time `json:"start,omitempty"` End *time.Time `json:"end,omitempty"` TookMS int64 `json:"took_ms,omitempty"` Took string `json:"took,omitempty"` Said string `json:"said,omitempty"` } // AppliedReport is a machine's first report, after a send, that it applied what it was sent. type AppliedReport struct { At time.Time Outcome string } // AppliedLookup answers a machine's first report after a send to it; false when it has not reported. type AppliedLookup func(node string, sent time.Time) (AppliedReport, bool) // point is one moment of the walk, ending the phase named. type point struct { phase string tier int at *time.Time none bool // the phase did not happen said string // why it is none, or unknown } // Phases is the walk's time phase by phase, at now; applied answers the rest's reports (nil: none read). func (p Plan) WalkPhases(now time.Time, applied AppliedLookup) *WalkPhases { if p.Release != nil || p.Batch() { return nil } w := &WalkPhases{} if p.Times != nil { w.Class = p.Times.Class } from := p.firstMerge(w) if from != nil { f := from.Truncate(time.Millisecond) from = &f } if from == nil { w.Said = "when its merge was made is not kept: its delivery time is unknown" } else { w.From = from } points := p.points(now, applied, w) // Walk the moments: each known moment after a known one is a measured phase; a missing moment makes the // phases up to the next known one unknown, and their span unknown time. last := from var pending []int for _, pt := range points { if pt.none { w.Phases = append(w.Phases, WalkPhase{Name: pt.phase, Tier: pt.tier, State: PhaseNone, Said: pt.said}) continue } ph := WalkPhase{Name: pt.phase, Tier: pt.tier, Said: pt.said} if pt.at == nil { ph.State = PhaseUnknown w.Phases = append(w.Phases, ph) pending = append(pending, len(w.Phases)-1) continue } at := pt.at.Truncate(time.Millisecond) ph.End = &at if last != nil && at.Before(*last) { // **Out of order** (issue 382): this moment was recorded before the one before it — a module's first // send kept before its build was, two clocks, a late record. Neither phase can be measured honestly: // the one before it is said unknown and its time joins the unknown time, and this one is said unknown // with no time. The walk goes on from the later moment, so nothing is negative and the sum holds. for i := len(w.Phases) - 1; i >= 0; i-- { prev := &w.Phases[i] if prev.End == nil { continue // a phase that did not happen, or is unknown with no moment: the one before carries last } if prev.State == PhaseMeasured { w.UnknownMS += prev.TookMS } prev.State, prev.Start, prev.TookMS, prev.Took = PhaseUnknown, nil, 0, "" prev.Said = "out of order: the phase after it ended first" break } ph.State, ph.Said = PhaseUnknown, "out of order: it ended before the phase before it" w.Phases = append(w.Phases, ph) pending = nil continue } if last != nil && len(pending) == 0 { start := *last ph.State, ph.Start = PhaseMeasured, &start ph.TookMS = at.Sub(start).Milliseconds() ph.Took = words(at.Sub(start)) } else { ph.State = PhaseUnknown if last != nil { w.UnknownMS += at.Sub(*last).Milliseconds() } } w.Phases = append(w.Phases, ph) pending = nil last = &at } if len(pending) > 0 { // The walk's last moments are unknown: it has no end. if w.Said == "" { w.Said = "its end is unknown: " + w.Phases[pending[0]].Name + " " + orNot(w.Phases[pending[0]].Said) } return w } if last == nil || from == nil { return w } if p.Open() { since := *last w.Phases = append(w.Phases, WalkPhase{Name: PhaseOpen, Tier: -1, State: PhaseOpen, Start: &since, TookMS: now.Sub(since).Milliseconds(), Took: words(now.Sub(since)), Said: "the walk is " + p.State}) w.Said = "open: " + p.State + ", " + words(now.Sub(*from)) + " since its merge" return w } if p.State != PlanDone { w.Said = "ended " + p.State + ": no delivery time" return w } end := *last w.End = &end w.TotalMS = end.Sub(*from).Milliseconds() w.Total = words(end.Sub(*from)) if w.Said == "" { if w.UnknownMS > 0 { w.Said = fmt.Sprintf("%s from its merge to running everywhere, %s of it unknown", w.Total, words(time.Duration(w.UnknownMS)*time.Millisecond)) } else { w.Said = w.Total + " from its merge to running everywhere" } } return w } // firstMerge is when the walk's earliest merge was made, and every merge's start kept on w. func (p Plan) firstMerge(w *WalkPhases) *time.Time { var first *time.Time if p.Delivery != nil { for _, m := range p.Delivery.Merges { if m.Merged.IsZero() { continue } w.Merges = append(w.Merges, MergeStart{Repository: m.Repository, Commit: m.Commit, Merged: m.Merged}) if first == nil || m.Merged.Before(*first) { at := m.Merged first = &at } } } if first == nil { for _, c := range p.Carried() { if c.Merged.IsZero() { continue } w.Merges = append(w.Merges, MergeStart{Repository: c.Repository, Commit: c.Commit, Merged: c.Merged}) if first == nil || c.Merged.Before(*first) { at := c.Merged first = &at } } } return first } // points are the walk's moments in order. func (p Plan) points(now time.Time, applied AppliedLookup, w *WalkPhases) []point { var out []point // The window and the wait behind another walk. cut := p.Created if p.Times != nil && p.Times.Cut != nil { cut = *p.Times.Cut } switch { case p.Times == nil: out = append(out, point{phase: PhaseWindow, tier: -1, said: "not kept for a walk made before ADR 0282"}, point{phase: PhaseQueued, tier: -1, at: &cut, said: "the window and the wait behind another walk together"}) case p.Times.WindowClosed == nil: out = append(out, point{phase: PhaseWindow, tier: -1, none: true, said: "no window: walked on its own"}, point{phase: PhaseQueued, tier: -1, at: &cut}) default: closed := *p.Times.WindowClosed out = append(out, point{phase: PhaseWindow, tier: -1, at: &closed}, point{phase: PhaseQueued, tier: -1, at: &cut}) } if p.Delivery != nil && p.Delivery.Awaits != "" { out = append(out, point{phase: PhaseWord, tier: -1, at: p.Delivery.Go, said: "waiting for " + p.Delivery.Awaits + "'s word"}) } else { out = append(out, point{phase: PhaseWord, tier: -1, none: true, said: "waits for no word"}) } var lastSends []restSend for t, tier := range p.Tiers { if t > p.Tier || (t == p.Tier && p.Tier < len(p.Tiers) && !askedAnyOf(p, tier)) { break } var asked, built, first, judged, rest *time.Time allBuilt, anyFirst, allJudged, allRest := true, false, true, true for _, m := range tier { s := p.Modules[m] if s == nil { allBuilt, allRest = false, false continue } asked = earliest(asked, s.AskedAt) if s.BuiltAt == nil { if s.State != "deleted" { allBuilt = false } } else { built = latest(built, s.BuiltAt) } if s.FirstAt != nil { anyFirst = true first = earliest(first, s.FirstAt) if s.Gate != nil { if s.Gate.JudgedAt == nil { allJudged = false } else { judged = latest(judged, s.Gate.JudgedAt) } } } if s.SentAt == nil { if s.State != "deleted" { allRest = false } } else { rest = latest(rest, s.SentAt) for node := range s.Rest { lastSends = append(lastSends, restSend{module: m, node: node, at: *s.SentAt}) } } } if t > 0 { out = append(out, point{phase: PhaseBetween, tier: t, at: asked}) } if !allBuilt { built = nil } out = append(out, point{phase: PhaseBuild, tier: t, at: built}) if anyFirst { if !allJudged { judged = nil } out = append(out, point{phase: PhaseSend, tier: t, at: first}, point{phase: PhaseJudge, tier: t, at: judged}) } else { out = append(out, point{phase: PhaseSend, tier: t, none: true, said: "no first machine: nothing to judge"}, point{phase: PhaseJudge, tier: t, none: true, said: "no first machine: nothing to judge"}) } if !allRest { rest = nil } out = append(out, point{phase: PhaseRest, tier: t, at: rest}) } if p.State != PlanDone { return out } // Every machine of the rest running the build: its first report after the send, applied. if p.Times == nil { return append(out, point{phase: PhaseApply, tier: -1, said: "the rest's sends are not kept for a walk made before ADR 0282"}) } if len(lastSends) == 0 { return append(out, point{phase: PhaseApply, tier: -1, none: true, said: "no machine beyond the first: the first-node gate's readings were its run"}) } if applied == nil { return append(out, point{phase: PhaseApply, tier: -1, said: "the rest's reports were not read"}) } var end *time.Time var waiting, failed []string for _, s := range lastSends { r, ok := applied(s.node, s.at) switch { case !ok && now.Sub(s.at) > ApplySilentAfter, ok && r.At.Sub(s.at) > ApplySilentAfter: w.Silent = appendOnce(w.Silent, s.node) case !ok: waiting = appendOnce(waiting, s.node) case r.Outcome != OutcomeApplied: failed = appendOnce(failed, s.node+" ("+r.Outcome+")") default: at := r.At end = latest(end, &at) } } switch { case len(failed) > 0: return append(out, point{phase: PhaseApply, tier: -1, said: "not applied on " + strings.Join(failed, ", ")}) case len(waiting) > 0: return append(out, point{phase: PhaseApply, tier: -1, said: "waiting for " + strings.Join(waiting, ", ") + " to report it applied"}) case end == nil: return append(out, point{phase: PhaseApply, tier: -1, said: "no machine of the rest heard from: " + strings.Join(w.Silent, ", ")}) } // A machine may have reported before the last tier ended: the walk ends at whichever is later. for i := len(out) - 1; i >= 0; i-- { if out[i].at != nil { if out[i].at.After(*end) { e := *out[i].at end = &e } break } } said := "" if len(w.Silent) > 0 { sort.Strings(w.Silent) said = "not waited for, not heard from: " + strings.Join(w.Silent, ", ") } return append(out, point{phase: PhaseApply, tier: -1, at: end, said: said}) } type restSend struct { module, node string at time.Time } func askedAnyOf(p Plan, tier []string) bool { for _, m := range tier { if s := p.Modules[m]; s != nil && s.AskedAt != nil { return true } } return false } func earliest(a, b *time.Time) *time.Time { if b == nil { return a } if a == nil || b.Before(*a) { t := *b return &t } return a } func latest(a, b *time.Time) *time.Time { if b == nil { return a } if a == nil || b.After(*a) { t := *b return &t } return a } func appendOnce(to []string, s string) []string { for _, x := range to { if x == s { return to } } return append(to, s) } func orNot(s string) string { if s == "" { return "is not known" } return "(" + s + ")" } // words is a duration as a person reads it. func words(d time.Duration) string { if d < 0 { return "-" + words(-d) } return d.Round(100 * time.Millisecond).String() } // FirstAppliedAfter is a machine's first report of a send made at or after a walk's send to it, within // ApplySilentAfter of it — the controller's `apply` durations, which measure every send to its first report (to-be // 45 Phase 0). A declaration sent later carries the build too, so its report counts; one sent past the bound is a // machine that slept, listed as silent and never moving the walk's end (ADR 0282 decision 1). func (i *Inventory) FirstAppliedAfter(ctx context.Context, node string, sent time.Time) (AppliedReport, bool, error) { var started time.Time var ms int64 var outcome string err := i.store.Pool().QueryRow(ctx, `select started, took_ms, detail from duration where kind = $1 and subject = $2 and started >= $3 and started < $4 order by started limit 1`, DurationApply, node, sent.Add(-appliedSlack), sent.Add(ApplySilentAfter)).Scan(&started, &ms, &outcome) if errors.Is(err, pgx.ErrNoRows) { return AppliedReport{}, false, nil } if err != nil { return AppliedReport{}, false, err } return AppliedReport{At: started.Add(time.Duration(ms) * time.Millisecond).UTC(), Outcome: outcome}, true, nil } // appliedSlack is how much earlier than a walk's record of its send the machine's own record of it may be: // the send is recorded on the machine first, then on the walk. const appliedSlack = 2 * time.Second