package link import ( "context" "crypto/sha256" "encoding/hex" "encoding/json" "fmt" "sync" "time" ) // Queue is the node-engine's one apply queue (novox/hq to-be 45 §6, ADR 0227 rule 1). // // **One worker applies; everything else asks it to.** A declaration delivered over the link and the // five-minute reconcile are two reasons to *enqueue* the same act, never two paths that apply. Before // this they were two, each taking an apply lock, and every order between them was a fault met on the // live mesh: the reconcile read what was kept before a delivery and applied it over the newer one // (novox/hq issues 257, 261); the reconcile went first and its report about the older declaration // reached the mesh after the delivery's (267). Each fix closed one order and left the class open. // // The queue holds at most one pending request, coalesced: whatever has been delivered and not yet // taken, and whether a reconcile is due. When the worker starts it takes the **newest declaration held // at that moment**, sets the rest aside — each reported as set aside, then settled unapplied — applies // it once, and makes **one report**, naming the declaration's order and its own report sequence. A // reconcile due while a delivery waits is that delivery's apply: applying the newest thing the mesh // said is what reconciling is. Only with nothing delivered does a reconcile apply what was kept, read // once its turn has come. // // **Reports leave in the order they are made**, because the worker makes and says them, one at a time. // A report made while the link is down waits: an apply's in the unsaid store (issue 264), a // reconcile's in one slot replaced by anything newer, and both are said, in that order, when the next // link opens — before anything newly delivered is applied. type Queue struct { // Membership is whose declarations these are: their signer, and the node the reports are from. Membership Membership // Apply applies a declaration the mesh signed. Apply Applier // Reconcile holds the machine to what it was last told, read once it is this act's turn. False // when there is nothing to hold it to, or the reconcile could not run; it says why itself. Reconcile Reconcile // News is whether a reconcile's report is worth saying unasked (novox/hq ADR 0100, ADR 0140). Nil // is never: a reconcile is otherwise silent. News func(Report) bool // Heard is told every account of the machine the mesh has taken — an apply's or a reconcile's, // not a refusal or a declaration set aside, which say nothing about the machine. Nil is allowed. Heard func(Report) // Unsaid keeps an apply's report until the mesh has taken it (issue 264). Nil keeps nothing. Unsaid Unsaid // Numbers keeps the report sequence and the refusals counted, across restarts and self-updates. // Nil counts in memory, which is only right where nothing outlives the process — a test. Numbers Numbers Say Announce Timeout time.Duration once sync.Once wake chan struct{} mu sync.Mutex idle *sync.Cond // delivering is a delivered declaration in hand: its report is owed on the link it came over, so // that link is not let go until it is said (detach). A reconcile in hand holds no link — its report // waits for the next one if this one goes. delivering bool waiting []Declaration // delivered and not yet taken, in the order they arrived due bool // a reconcile asked for bus Bus // the link open now; nil while there is none unasked *Report // a reconcile's report not yet said, made while there was no link // saying is held while a report is made and said, so the order they leave in is the order they // are made in — whoever says them, the worker or a link opening. saying sync.Mutex counted bool numbers numbered } // Reconcile is the reconcile's act, run by the queue's worker. See Queue.Reconcile. type Reconcile func(ctx context.Context) (Report, bool) // Numbers is where the queue keeps what it counts, so a successor goes on from where it stopped. type Numbers interface { // Read is what was kept: zero for a node that never reported, an error when it cannot be read — // never zero for "could not tell". Read() (reportSequence, refusedOlder int64, err error) // Save keeps them, on disk when it returns. Save(reportSequence, refusedOlder int64) error } type numbered struct{ sequence, refused int64 } func (q *Queue) init() { q.once.Do(func() { q.wake = make(chan struct{}, 1) q.idle = sync.NewCond(&q.mu) if q.Say == nil { q.Say = func(string) {} } }) } // Deliver enqueues what the link delivered, in the order it arrived. It never waits for an apply. func (q *Queue) Deliver(arrived ...Declaration) { q.init() q.mu.Lock() q.waiting = append(q.waiting, arrived...) q.mu.Unlock() q.signal() } // ReconcileDue enqueues a reconcile. One asked for while another is waiting is the same one. func (q *Queue) ReconcileDue() { q.init() q.mu.Lock() q.due = true q.mu.Unlock() q.signal() } func (q *Queue) signal() { select { case q.wake <- struct{}{}: default: // One is already waiting to be read, and the worker takes everything pending when it does. } } // Run is the worker. It returns once the context ends and the act in hand is finished and said: // asked to stop — the host standing aside for its successor — nothing further is applied. func (q *Queue) Run(ctx context.Context) { q.init() for { for q.step(ctx) { } select { case <-ctx.Done(): return case <-q.wake: } } } // step takes what is pending and does it, once. False when there was nothing to do, or the queue // was asked to stop. func (q *Queue) step(ctx context.Context) bool { q.mu.Lock() if ctx.Err() != nil || (len(q.waiting) == 0 && !q.due) { q.mu.Unlock() return false } batch, due := q.waiting, q.due q.waiting, q.due = nil, false q.delivering = len(batch) > 0 q.mu.Unlock() defer func() { q.mu.Lock() q.delivering = false q.idle.Broadcast() q.mu.Unlock() }() if len(batch) == 0 { q.reconcile(ctx) return true } latest, aside := pick(batch) for _, old := range aside { if d := declaredIn(old.Body()); d != "" && d == declaredIn(latest.Body()) { // The same declaration again — delivered twice while the worker was busy. Nothing is set // aside: it is about to be applied. _ = old.Handled() continue } q.Say(fmt.Sprintf("set aside declaration %s (%s): a newer one arrived with it", short(declaredIn(old.Body())), orderOf(old.Body()).Words())) q.tell(ctx, Report{Declared: declaredIn(old.Body()), Order: orderOf(old.Body()), Superseded: declaredIn(latest.Body())}, setAside, old) } report := handleBody(ctx, q.Membership, latest.Body(), q.Apply) switch { case report.OlderThan != nil: // **Refused, counted and said** (rule 2): one line here, the count on this and every later // report, and the report itself naming what was refused and what is held. total := q.countRefusal() q.Say(fmt.Sprintf("refused declaration %s (%s): this node applied %s, from a newer lease "+ "holder — %d refused as older so far", short(report.Declared), report.Order.Words(), report.OlderThan.Words(), total)) case report.Refused != "": q.Say("refused a declaration: " + report.Refused) case len(report.Failed) > 0: q.Say(fmt.Sprintf("applied %d and failed: %v%s", len(report.Applied), report.Failed, heldNote(report.Held))) default: q.Say(fmt.Sprintf("applied %d resource(s)%s", len(report.Applied), heldNote(report.Held))) } if due && report.Refused != "" { // Nothing was applied, so the reconcile this stood in for has not happened. It runs next. q.mu.Lock() q.due = true q.mu.Unlock() } // A refusal as older is not kept to be said again: what the mesh waits for from this node is the // account of the last apply, and that is still what is kept. kind := applied if report.OlderThan != nil { kind = refusedOlder } q.tell(ctx, report, kind, latest) return true } func (q *Queue) reconcile(ctx context.Context) { if q.Reconcile == nil { return } report, ok := q.Reconcile(ctx) if !ok || q.News == nil || !q.News(report) { return } q.tell(ctx, report, unasked) } // What a report is, for what tell does with it besides saying it. type kindOf int const ( // applied is an apply's account: kept until said (issue 264), and fresher than any reconcile's // report still waiting. applied kindOf = iota // setAside is a declaration a newer one took the place of, unapplied. setAside // refusedOlder is a declaration refused as older than what was applied. Not kept: what the mesh // waits for from this node is still the account of the last apply. refusedOlder // unasked is a reconcile's report. Not said — no link, or the bus refused it — the newest waits // for the next link, and anything the worker says before then replaces it. unasked ) // tell makes a report this node's — its node, its report sequence, the refusals counted — and says it // on the link open now. Kept first, when it is an apply's, so a report lost between here and the bus // is said by whichever host links next. Each declaration it settles is settled after it is said, // either way. True when the broker took it. func (q *Queue) tell(ctx context.Context, r Report, kind kindOf, settle ...Declaration) bool { q.saying.Lock() defer q.saying.Unlock() r = q.stamp(r) keep := kind == applied if keep { // What this says supersedes any reconcile's report still waiting to be said. q.unasked = nil if q.Unsaid != nil { if err := q.Unsaid.Keep(r); err != nil { q.Say("cannot keep this apply's report until it is said: " + err.Error()) } } } q.mu.Lock() bus := q.bus q.mu.Unlock() published := false if bus != nil { // **Said even when the link is being let go** (issue 264): not cancelled with it, still // bounded by the timeout, and the link is not closed until this returns (detach). published = publishReport(context.WithoutCancel(ctx), bus, q.Membership, r, q.Say, q.Timeout) } if !published && kind == unasked { q.unasked = &r } if published { if keep && q.Unsaid != nil { if err := q.Unsaid.Said(r.Declared); err != nil { q.Say("said this apply's report, and cannot forget it: " + err.Error()) } } if q.Heard != nil && (kind == applied || kind == unasked) && r.Refused == "" { q.Heard(r) } } // Settled after the report is published. A node that dies between applying and reporting // leaves the declaration with the mesh and applies it again on return, which is safe because // applying is reconciliation — it converges rather than repeating. for _, d := range settle { _ = d.Handled() } return published } // stamp gives a report its node, the next report sequence and the refusals counted. Called holding // saying. func (q *Queue) stamp(r Report) Report { q.load() q.numbers.sequence++ q.keepNumbers() r.Node = q.Membership.Node r.ReportSequence = q.numbers.sequence r.RefusedOlder = q.numbers.refused return r } func (q *Queue) countRefusal() int64 { q.saying.Lock() defer q.saying.Unlock() q.load() q.numbers.refused++ q.keepNumbers() return q.numbers.refused } // load reads what was counted, once. **Unreadable is said, and the sequence goes on from above // anything it can have reached** — the time in milliseconds — rather than from one, which would make // every report this host says next older than what it already said, and refused. func (q *Queue) load() { if q.counted { return } q.counted = true if q.Numbers == nil { return } sequence, refused, err := q.Numbers.Read() if err != nil { floor := time.Now().UnixMilli() q.Say(fmt.Sprintf("cannot read this node's report sequence: %v — going on from %d, above any "+ "it can have reached, and counting refusals again from none", err, floor)) q.numbers = numbered{sequence: floor} return } q.numbers = numbered{sequence: sequence, refused: refused} } func (q *Queue) keepNumbers() { if q.Numbers == nil { return } if err := q.Numbers.Save(q.numbers.sequence, q.numbers.refused); err != nil { // Said, and not fatal: the report is still worth saying. What is lost is that a successor // could number one again. q.Say("cannot keep this node's report sequence: " + err.Error()) } } // attach is a link opening. **What the last apply did, if the mesh never heard it**, is said first, // exactly as it was made; then a reconcile's report that waited; and only then may the worker say // anything on it — so a report said again can never land after, and read as newer than, a later one. func (q *Queue) attach(ctx context.Context, bus Bus) { q.init() q.saying.Lock() defer q.saying.Unlock() reporting := context.WithoutCancel(ctx) q.load() said := true if q.Unsaid != nil { switch kept, ok, err := q.Unsaid.Pending(); { case err != nil: q.Say("cannot read the report this node kept unsaid: " + err.Error()) case ok: q.Say("saying again what the last apply did: its report never reached the mesh") if said = publishReport(reporting, bus, q.Membership, kept, q.Say, q.Timeout); said { if err := q.Unsaid.Said(kept.Declared); err != nil { q.Say("said the kept report, and cannot forget it: " + err.Error()) } if q.Heard != nil && kept.Refused == "" { q.Heard(kept) } } } } if said && q.unasked != nil { if publishReport(reporting, bus, q.Membership, *q.unasked, q.Say, q.Timeout) { if q.Heard != nil { q.Heard(*q.unasked) } q.unasked = nil } } q.mu.Lock() q.bus = bus q.mu.Unlock() } // detach is the link closing: it waits for a delivery in hand to be applied and said — an apply that // stood aside for its successor is reported on this link before it is let go (issue 264) — and then // nothing more is said on it. A reconcile in hand is not waited for: a link that was dropped or roused // is opened again at once, and the reconcile's report, if it is news, waits for it. func (q *Queue) detach(bus Bus) { q.init() q.mu.Lock() for q.delivering { q.idle.Wait() } q.mu.Unlock() q.saying.Lock() defer q.saying.Unlock() q.mu.Lock() if q.bus == bus { q.bus = nil } q.mu.Unlock() } // pick is the newest of what was delivered and the rest, in arrival order, to set aside: by order // when the declarations claim one, by arrival when they do not (Order.Supersedes). func pick(batch []Declaration) (Declaration, []Declaration) { latest := batch[0] var aside []Declaration for _, next := range batch[1:] { if orderOf(next.Body()).Supersedes(orderOf(latest.Body())) { aside = append(aside, latest) latest = next continue } aside = append(aside, next) } return latest, aside } // orderOf is the order a signed declaration claims, or none when it claims none or cannot be read. // Read from the envelope alone; the signature is verified when the winner is applied, and a forged // message that lied about its order would only set aside real ones — which are reported as set aside, // and the next push sends the current one again. func orderOf(body []byte) Order { var signed Signed if err := json.Unmarshal(body, &signed); err != nil { return Order{} } var o Order if err := json.Unmarshal(signed.Declaration, &o); err != nil || o.Epoch < 0 || o.Sequence < 0 { return Order{} } return o } // declaredIn names a signed declaration as the mesh does — the digest of the declaration's bytes, the // same the apply's report carries — for a report about one that was not applied. Empty if the message // is not one: a forged or garbled message is refused when its turn comes; here it is only named. func declaredIn(body []byte) string { var signed Signed if err := json.Unmarshal(body, &signed); err != nil || len(signed.Declaration) == 0 { return "" } sum := sha256.Sum256(signed.Declaration) return hex.EncodeToString(sum[:]) }