Files
mesh-controller/internal/inventory/sends.go
T
jschoubben 313efa826c
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 pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
Say in every declaration the assignment generation it came from, and record every send (hq issue 234)
A declaration newer in sequence than every other named four fewer modules
than the assignments held, a machine applied it, and nothing the mesh kept
said who sent it or from what view. The store now raises an assignment
generation in the same transaction as every assignment change (a trigger,
so a cascade counts too); a send reads it before composing, carries it to
engines that said they read it, marks a gate's put-back, and is recorded
with its sequence, sender, generation and modules. The would-send is
stamped with what was last sent, so nothing reads behind for it; a refusal
is kept on the send it refused and raised as S20 naming the sender; plan
says what the machine was last told.
2026-10-11 11:30:16 +02:00

203 lines
8.6 KiB
Go

package inventory
import (
"context"
"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 and whether that send was a put-back, by its
// id; zero for one sent without. What the mesh WOULD send is stamped with these, 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, bool, error) {
var g *int64
var putBack bool
if err := i.store.Pool().QueryRow(ctx, `select sent_generation, sent_put_back from node where id = $1`, id).
Scan(&g, &putBack); err != nil {
return 0, false, fmt.Errorf("reading the generation %s was last sent: %w", id, err)
}
if g == nil {
return 0, putBack, nil
}
return *g, putBack, 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; PutBack whether a
// gate sent it to put the machine back.
Generation int64
PutBack bool
// ToldGeneration and ToldPutBack are what the machine was told of them on the wire: zero and false for a
// machine whose node-engine has not said it reads a generation. What the would-send is stamped with.
ToldGeneration int64
ToldPutBack bool
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.
RefusedAt *time.Time
RefusedApplied int64
}
// RecordSend writes one send down, and what the machine was last sent of its generation, in one
// transaction; the oldest beyond the newest 200 for the machine are let go.
func (i *Inventory) RecordSend(ctx context.Context, s Send) error {
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_generation = $2, sent_put_back = $3 where id = $1`,
s.Node, nullIfZero(s.ToldGeneration), s.ToldPutBack); err != nil {
return err
}
modules := s.Modules
if modules == nil {
modules = []string{}
}
if _, err := tx.Exec(ctx,
`insert into declaration_send (node, sequence, epoch, sender, generation, put_back, digest, modules)
values ($1, $2, $3, $4, $5, $6, $7, $8)`,
s.Node, nullIfZero(s.Sequence), nullIfZero(s.Epoch), s.Sender, nullIfZero(s.Generation), s.PutBack,
s.Digest, modules); err != nil {
return err
}
if _, err := tx.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)`, s.Node, sendsKept); err != nil {
return err
}
return tx.Commit(ctx)
}
// 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.put_back, s.digest, s.modules, s.sent_at, s.refused_at, coalesce(s.refused_applied, 0)`
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.PutBack,
&s.Digest, &s.Modules, &s.SentAt, &s.RefusedAt, &s.RefusedApplied); 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; false when none was recorded.
func (i *Inventory) LastSend(ctx context.Context, name string) (Send, bool, error) {
rows, err := i.store.Pool().Query(ctx, `select `+sendColumns+`
from declaration_send s join node n on n.id = s.node
where n.name = $1 order by s.id desc limit 1`, name)
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
}
// 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 — sent before the
// record, or by a hand the mesh did not see. The newest of that sequence when there are several: a
// declaration a person sent by hand may repeat one.
func (i *Inventory) RefusedSend(ctx context.Context, node string, sequence, applied int64) (Send, bool, error) {
rows, err := i.store.Pool().Query(ctx, `with refused as (
update declaration_send set refused_at = now(), refused_applied = $3
where id = (select id from declaration_send where node = $1 and sequence = $2 order by id desc limit 1)
returning *)
select `+sendColumns+` from refused s join node n on n.id = s.node`, node, sequence, applied)
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
}
// 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)
}