Files
mesh-controller/internal/inventory/sends.go
T
jschoubben e2801773c3
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 delivered
Take an ordered engine's report as a send's answer only by its epoch and sequence (novox/hq issue 485)
A periodic report about an identical earlier send names the same digest and
can arrive after a later send the machine never received, so the digest and
the time read it as applied. The digest-and-time rule stays for engines that
report no order. Tests now hold the digest condition and the sequence one.
2026-10-11 20:44:54 +02:00

341 lines
15 KiB
Go

package inventory
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/jackc/pgx/v5"
)
// The assignment generation and the record of every send (novox/hq issue 234, migration 0094).
//
// A declaration carries the generation of the assignments it was composed from, and a machine refuses one
// composed from an older generation than it applied — the sequence orders arrival and cannot tell a later
// send carrying an older view. And every send is written down with who sent it, from which generation and
// naming which modules, so what a machine was last told, by whom, is read back rather than guessed at.
// sendsKept is how many sends are kept per machine: the newest. Enough to read back a morning of pushes
// on a busy machine; the record is for finding who sent what, not an archive.
const sendsKept = 200
// AssignmentGeneration is the mesh's assignment generation now: raised by the store itself in the same
// transaction as every change to what is assigned where (a trigger on the assignment table).
func (i *Inventory) AssignmentGeneration(ctx context.Context) (int64, error) {
var g int64
if err := i.store.Pool().QueryRow(ctx, `select generation from assignment_generation`).Scan(&g); err != nil {
return 0, fmt.Errorf("reading the assignment generation: %w", err)
}
return g, nil
}
// RaiseAssignmentGeneration raises the generation past one a machine applied, and answers it: one more
// than that, or what it already was when it is past it. It never lowers it.
//
// For one case only: a machine refused a send composed from the generation the mesh holds now, so the
// machine applied a higher one than the mesh's counter — a store put back from a backup. Every send after
// would be refused for ever; raised past it, the next send carries a generation the machine takes, composed
// from the assignments the store holds now. A stale sender never meets this: its generation is below the
// counter, and the counter is not touched for it.
func (i *Inventory) RaiseAssignmentGeneration(ctx context.Context, past int64) (int64, error) {
var g int64
err := i.store.Pool().QueryRow(ctx,
`update assignment_generation set generation = greatest(generation, $1 + 1) returning generation`, past).Scan(&g)
if err != nil {
return 0, fmt.Errorf("raising the assignment generation past %d: %w", past, err)
}
return g, nil
}
// SentGeneration is the generation a machine was last sent, by its id; zero for one sent without. What the
// mesh WOULD send is stamped with it, so it reads as byte for byte what it DID send when nothing else changed.
func (i *Inventory) SentGeneration(ctx context.Context, id string) (int64, error) {
var g *int64
if err := i.store.Pool().QueryRow(ctx, `select sent_generation from node where id = $1`, id).Scan(&g); err != nil {
return 0, fmt.Errorf("reading the generation %s was last sent: %w", id, err)
}
if g == nil {
return 0, nil
}
return *g, nil
}
// ReadsGeneration says a machine's node-engine said it reads a generation in a declaration, by its id.
func (i *Inventory) ReadsGeneration(ctx context.Context, id string) (bool, error) {
var reads bool
if err := i.store.Pool().QueryRow(ctx, `select reads_generation from node where id = $1`, id).Scan(&reads); err != nil {
return false, fmt.Errorf("reading whether %s reads a generation: %w", id, err)
}
return reads, nil
}
// RecordReadsGeneration keeps what a machine's latest report said of reading a generation.
func (i *Inventory) RecordReadsGeneration(ctx context.Context, id string, reads bool) error {
_, err := i.store.Pool().Exec(ctx, `update node set reads_generation = $2 where id = $1`, id, reads)
return err
}
// Send is one declaration sent to a machine, as the mesh records it.
type Send struct {
// Node is the machine's id; NodeName its name, filled where it is read back.
Node string
NodeName string
Sequence int64
// Epoch is the lease epoch its sender acted under, zero for one that acted under none.
Epoch int64
// Sender is who sent it, in words: the caller, and the process and build that composed it.
Sender string
// Generation is the assignment generation it was composed from, zero when not known.
Generation int64
Digest string
// Modules are the modules it named.
Modules []string
SentAt time.Time
// RefusedAt is when the machine refused it for its generation, nil when it did not; RefusedApplied the
// generation the machine said it had applied then; CounterRaisedTo what the mesh's counter was raised to
// because the refusal showed it behind the machine, zero when it was not.
RefusedAt *time.Time
RefusedApplied int64
CounterRaisedTo int64
// Recorded is false for a refusal of a sequence the mesh has no send for: its sender is unknown.
Recorded bool
// Answer is what the machine last reported of this send; filled only where a last send is read back
// (LastSend, LastSends), the zero answer elsewhere.
Answer SendAnswer
}
// SendAnswer is what a machine reported of one send, read from its last report when that report is about the
// send's declaration (novox/hq issue 485): whether it applied it, failed part of it or refused it, and why.
// The zero answer is a machine that has not reported on this send — yet, or since a later report replaced it.
type SendAnswer struct {
// Outcome is applied, failed or refused, as the machine reported it; empty when it has not reported on it.
Outcome string
// Refused is the machine's own words when it refused the declaration whole.
Refused string
// Failed is how many of its resources failed, when some did.
Failed int
// At is when the machine reported it; nil when it has not.
At *time.Time
}
// unknownSender is the sender of a refused sequence the mesh has no record of sending.
const unknownSender = "a sender the mesh has no record of (no send of this sequence was recorded: sent by " +
"hand, before sends were recorded, or by a controller whose record was not written)"
// SendNotRecordedError is a send whose machine's record was written — its digest and the generation it
// carried — and whose record of who sent it was not. The send is away and the machine's record stands; the
// caller says this loudly and does not take the send for failed.
type SendNotRecordedError struct{ Err error }
func (e *SendNotRecordedError) Error() string {
return "the record of the send was not written: " + e.Err.Error()
}
func (e *SendNotRecordedError) Unwrap() error { return e.Err }
// RecordSentWith writes down what a machine was just sent — its digest, the builds it carried, the epoch and
// the generation on the wire — and the send itself, who sent it and from which generation (novox/hq issue
// 234), in one transaction. The digest and the generation are written together, so the would-send is never
// stamped with a generation the machine was not sent. The send's own record is written under a savepoint: if
// it cannot be, the machine's record is still committed and a *SendNotRecordedError says so. The oldest
// sends beyond the newest 200 for the machine are let go.
func (i *Inventory) RecordSentWith(ctx context.Context, node, digest string, builds map[string]string, epoch uint64,
generation int64, s Send) error {
var sentEpoch *int64
if epoch > 0 {
e := int64(epoch)
sentEpoch = &e
}
var carried *string
if builds != nil {
raw, err := json.Marshal(builds)
if err != nil {
return err
}
text := string(raw)
carried = &text
}
tx, err := i.store.Pool().Begin(ctx)
if err != nil {
return err
}
defer func() { _ = tx.Rollback(ctx) }()
if _, err := tx.Exec(ctx,
`update node set sent = $2, sent_at = now(), sent_builds = $3::jsonb, sent_epoch = $4, sent_generation = $5
where id = $1`, node, digest, carried, sentEpoch, nullIfZero(generation)); err != nil {
return err
}
recordErr := recordSend(ctx, tx, node, s)
if err := tx.Commit(ctx); err != nil {
return err
}
if recordErr != nil {
return &SendNotRecordedError{Err: recordErr}
}
return nil
}
// recordSend writes one send under a savepoint of tx, which it rolls back when the send cannot be written.
func recordSend(ctx context.Context, tx pgx.Tx, node string, s Send) error {
sp, err := tx.Begin(ctx)
if err != nil {
return err
}
defer func() { _ = sp.Rollback(ctx) }()
modules := s.Modules
if modules == nil {
modules = []string{}
}
if _, err := sp.Exec(ctx,
`insert into declaration_send (node, sequence, epoch, sender, generation, digest, modules)
values ($1, $2, $3, $4, $5, $6, $7)`,
node, nullIfZero(s.Sequence), nullIfZero(s.Epoch), s.Sender, nullIfZero(s.Generation), s.Digest,
modules); err != nil {
return err
}
if err := pruneSends(ctx, sp, node); err != nil {
return err
}
return sp.Commit(ctx)
}
// pruneSends lets go of a machine's sends beyond the newest 200.
func pruneSends(ctx context.Context, q queries, node string) error {
_, err := q.Exec(ctx,
`delete from declaration_send where node = $1 and id not in
(select id from declaration_send where node = $1 order by id desc limit $2)`, node, sendsKept)
return err
}
// sendColumns are a send's columns as scanSends reads them.
const sendColumns = `s.node, n.name, coalesce(s.sequence, 0), coalesce(s.epoch, 0), s.sender, coalesce(s.generation, 0),
s.digest, s.modules, s.sent_at, s.refused_at, coalesce(s.refused_applied, 0), coalesce(s.counter_raised_to, 0),
s.recorded`
func scanSends(rows pgx.Rows) ([]Send, error) {
defer rows.Close()
var out []Send
for rows.Next() {
var s Send
if err := rows.Scan(&s.Node, &s.NodeName, &s.Sequence, &s.Epoch, &s.Sender, &s.Generation, &s.Digest,
&s.Modules, &s.SentAt, &s.RefusedAt, &s.RefusedApplied, &s.CounterRaisedTo, &s.Recorded); err != nil {
return nil, err
}
out = append(out, s)
}
return out, rows.Err()
}
// LastSend is the last declaration a machine was sent, by its name, with what the machine answered of it; false
// when none was recorded. A refusal of a sequence the mesh has no send for is not a send, and is not it.
func (i *Inventory) LastSend(ctx context.Context, name string) (Send, bool, error) {
sends, err := i.lastSends(ctx, `n.name = $1`, name)
if err != nil {
return Send{}, false, err
}
s, found := sends[name]
return s, found, nil
}
// LastSends is every machine's last send, by machine name, each with what the machine answered of it: what
// `status` says per machine (novox/hq issue 485). A machine never sent anything has no entry.
func (i *Inventory) LastSends(ctx context.Context) (map[string]Send, error) {
return i.lastSends(ctx, `true`)
}
// lastSends reads the newest recorded send of each machine matching where, beside the machine's last report
// when that report is about the send (novox/hq issue 485). It names the send's declaration (its digest), and:
// - from an ordered node-engine, one that reports the epoch and sequence it acted on, those are the send's
// own. A periodic report about an identical earlier send names the same digest and may arrive after a later
// send the machine never received; only the sequence tells them apart. It also holds for a fast machine
// whose report lands before the send's row is stamped.
// - from an engine that reports no order, the report came after the send — both times the store's own — so
// a machine sent again a declaration it once applied does not read as having applied the new send.
func (i *Inventory) lastSends(ctx context.Context, where string, args ...any) (map[string]Send, error) {
rows, err := i.store.Pool().Query(ctx, `select distinct on (n.name) `+sendColumns+`,
coalesce(r.outcome, ''), coalesce(r.refused, ''), case when jsonb_typeof(r.failed) = 'array' then jsonb_array_length(r.failed) else 0 end, r.at
from declaration_send s join node n on n.id = s.node
left join node_report r on r.node = s.node and r.declared <> '' and r.declared = s.digest
and case when r.reported_sequence is not null
then r.reported_sequence = s.sequence and coalesce(r.reported_epoch, 0) = coalesce(s.epoch, 0)
else r.at >= s.sent_at end
where s.recorded and `+where+`
order by n.name, s.id desc`, args...)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]Send{}
for rows.Next() {
var s Send
if err := rows.Scan(&s.Node, &s.NodeName, &s.Sequence, &s.Epoch, &s.Sender, &s.Generation, &s.Digest,
&s.Modules, &s.SentAt, &s.RefusedAt, &s.RefusedApplied, &s.CounterRaisedTo, &s.Recorded,
&s.Answer.Outcome, &s.Answer.Refused, &s.Answer.Failed, &s.Answer.At); err != nil {
return nil, err
}
out[s.NodeName] = s
}
return out, rows.Err()
}
// RefusedSend keeps a machine's refusal of the send of a sequence for the generation it came from, and
// answers that send, sender and all; false when no send of that sequence was recorded (RecordUnrecordedRefusal
// keeps that one). The newest of that sequence when there are several: a declaration a person sent by hand may
// repeat one. raisedTo is what the mesh's counter was raised to for it, zero for none.
func (i *Inventory) RefusedSend(ctx context.Context, node string, sequence, applied, raisedTo int64) (Send, bool, error) {
rows, err := i.store.Pool().Query(ctx, `with refused as (
update declaration_send set refused_at = now(), refused_applied = $3, counter_raised_to = $4
where id = (select id from declaration_send where node = $1 and sequence = $2 and recorded
order by id desc limit 1)
returning *)
select `+sendColumns+` from refused s join node n on n.id = s.node`, node, sequence, applied,
nullIfZero(raisedTo))
if err != nil {
return Send{}, false, err
}
sends, err := scanSends(rows)
if err != nil || len(sends) == 0 {
return Send{}, false, err
}
return sends[0], true, nil
}
// RecordUnrecordedRefusal keeps a machine's refusal of a sequence the mesh has no send for, as a row of its
// own whose sender is unknown, and answers it: S20 raises it like any other, because a refusal nobody can
// attribute is no quieter for that. s carries the machine, the sequence, epoch, generation and digest the
// refused declaration carried, and the refusal's applied generation and counter raise.
func (i *Inventory) RecordUnrecordedRefusal(ctx context.Context, s Send) (Send, error) {
rows, err := i.store.Pool().Query(ctx, `with refused as (
insert into declaration_send (node, sequence, epoch, sender, generation, digest, refused_at, refused_applied,
counter_raised_to, recorded)
values ($1, $2, $3, $4, $5, $6, now(), $7, $8, false) returning *)
select `+sendColumns+` from refused s join node n on n.id = s.node`,
s.Node, nullIfZero(s.Sequence), nullIfZero(s.Epoch), unknownSender, nullIfZero(s.Generation), s.Digest,
s.RefusedApplied, nullIfZero(s.CounterRaisedTo))
if err != nil {
return Send{}, err
}
sends, err := scanSends(rows)
if err != nil {
return Send{}, err
}
if len(sends) == 0 {
return Send{}, fmt.Errorf("the refusal of sequence %d was not kept", s.Sequence)
}
if err := pruneSends(ctx, i.store.Pool(), s.Node); err != nil {
return Send{}, err
}
return sends[0], nil
}
// RefusedSendsSince is every send refused for its generation since a moment, newest first.
func (i *Inventory) RefusedSendsSince(ctx context.Context, since time.Time) ([]Send, error) {
rows, err := i.store.Pool().Query(ctx, `select `+sendColumns+`
from declaration_send s join node n on n.id = s.node
where s.refused_at >= $1 order by s.refused_at desc`, since)
if err != nil {
return nil, err
}
return scanSends(rows)
}