Asking what the mesh would send must not change it
`status` hung. It composes a declaration for every node to answer *is this machine running what I would send it*, and composing one assigns each module a machine port — so the question wrote to the database, and wrote to the same rows as the machine it was asking about. `port_assignment` is unique on (node, machine). Two transactions inserting the same port do not race, they queue: the second waits on the index until the first commits. A status polled every two seconds while a node applies is two writers on those rows, and the poll stopped returning rather than returning something wrong — which is the better failure of the two, and still a failure. The latent version of this was there before anything polled: two compositions running at once could both allocate. So allocation belongs to the send path alone. The mesh chooses a port when it commits to sending one; every other caller reads what was chosen. A module with nothing assigned has never been sent, which is precisely what "waiting" means — the read needs no number to be right about that, and inventing one would make the answer worse. Named rather than passed as a bare bool: at three call sites, `true` and `false` say nothing about which of these two things is meant. Checked by the lab, which now polls status throughout an apply.
This commit is contained in:
@@ -244,14 +244,35 @@ func declarationFor(ctx context.Context, open *stores, node string,
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return declarationWith(ctx, open, node, plan, settings, gens)
|
||||
return declarationWith(ctx, open, node, plan, settings, gens, Reading)
|
||||
}
|
||||
|
||||
// declarationWith is the same, for a caller that has already worked out the generators once and
|
||||
// is about to use them for every node.
|
||||
// Choosing says whether this composition may allocate what has not been allocated yet.
|
||||
//
|
||||
// **Asking what the mesh would send must not change what the mesh would send.** Composing a
|
||||
// declaration assigns each module a machine port, and `status` composes one for every node to
|
||||
// answer *is this machine running what I would send it* — so the question allocated, wrote, and
|
||||
// contended with the very machine it was asking about. A status command that polls every two
|
||||
// seconds while a node is applying is then two writers on the same rows, which is how it came to
|
||||
// hang rather than answer.
|
||||
//
|
||||
// So the mesh chooses a port when it commits to sending one, and every other caller reads what
|
||||
// was chosen. A module with nothing assigned yet has never been sent, which is exactly what a
|
||||
// machine "waiting" means — the read needs no number to be right about that.
|
||||
type Choosing bool
|
||||
|
||||
const (
|
||||
// Allocating is the send path: what is not assigned yet is assigned now and kept.
|
||||
Allocating Choosing = true
|
||||
// Reading is every question: what is assigned is used, and nothing is created.
|
||||
Reading Choosing = false
|
||||
)
|
||||
|
||||
func declarationWith(ctx context.Context, open *stores, node string,
|
||||
plan catalogue.Resolution, settings catalogue.SettingsBy,
|
||||
gens map[string]catalogue.Generator) ([]map[string]any, error) {
|
||||
gens map[string]catalogue.Generator, choosing Choosing) ([]map[string]any, error) {
|
||||
inv := open.inventory
|
||||
grants, err := grantsFor(ctx, open, node)
|
||||
if err != nil {
|
||||
@@ -263,6 +284,21 @@ func declarationWith(ctx context.Context, open *stores, node string,
|
||||
// assigned anywhere: any number it picks is a guess about a machine it has never seen. Made
|
||||
// before the declaration is composed, because the container's mapping, the rule set and what a
|
||||
// consumer is told are all derived from it.
|
||||
// What this machine was already given, for a composition that may not allocate.
|
||||
already := map[string]map[int]int{}
|
||||
if choosing == Reading {
|
||||
held, err := inv.PortsFor(ctx, node)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, a := range held {
|
||||
if already[a.Module] == nil {
|
||||
already[a.Module] = map[int]int{}
|
||||
}
|
||||
already[a.Module][a.Wanted] = a.Machine
|
||||
}
|
||||
}
|
||||
|
||||
ports := map[string]map[int]int{}
|
||||
for _, m := range plan.Modules {
|
||||
for _, l := range m.Listens {
|
||||
@@ -273,7 +309,8 @@ func declarationWith(ctx context.Context, open *stores, node string,
|
||||
// *where this module's port is on this machine* and every reader of it needs that
|
||||
// answer whether or not the mesh was the one who chose it.
|
||||
where, mayAssign := m.MachineSide(l.Port)
|
||||
if mayAssign {
|
||||
switch {
|
||||
case mayAssign && choosing == Allocating:
|
||||
at, err := inv.PortFor(ctx, node, m.Module, l.Port, l.Fixed)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf(
|
||||
@@ -281,6 +318,11 @@ func declarationWith(ctx context.Context, open *stores, node string,
|
||||
m.Module, l.Port, node, err)
|
||||
}
|
||||
where = at.Machine
|
||||
case mayAssign:
|
||||
// Whatever was chosen last time, and nothing if there was no last time.
|
||||
if at, known := already[m.Module][l.Port]; known {
|
||||
where = at
|
||||
}
|
||||
}
|
||||
if ports[m.Module] == nil {
|
||||
ports[m.Module] = map[int]int{}
|
||||
|
||||
@@ -265,7 +265,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
// The private network is in here with everything else. It used to be composed separately
|
||||
// and prepended, which meant every machine with an address was on it and no machine could
|
||||
// be kept off. It is a module now, so it arrives the way a module does.
|
||||
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens)
|
||||
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Allocating)
|
||||
if err != nil {
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err))
|
||||
continue
|
||||
@@ -335,7 +335,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||
continue
|
||||
}
|
||||
resources, err := declarationWith(ctx, open, name, plan, settings, gens)
|
||||
resources, err := declarationWith(ctx, open, name, plan, settings, gens, Allocating)
|
||||
if err != nil {
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||
continue
|
||||
@@ -399,7 +399,7 @@ func wouldSend(ctx context.Context, open *stores,
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens)
|
||||
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user