The host applies the newest declaration, a file may be created once, the foundation filters first
031: a window of unacknowledged declarations is drained to the newest; the rest are set aside and reported as superseded. 035: a file resource may say create-once — written when absent, kept untouched when present (ADR 0087). 054: the bundle installs nftables and loads a base ruleset before the store and broker, in the table the filter module later replaces (ADR 0088).
This commit is contained in:
@@ -57,6 +57,14 @@ type Report struct {
|
||||
// Refused is set when the declaration was rejected whole rather than applied in part.
|
||||
Refused string `json:"refused,omitempty"`
|
||||
|
||||
// Superseded names the newer declaration this one was set aside for, unapplied.
|
||||
//
|
||||
// A machine asked to be five successive things becomes the last one (novox/hq issue 031):
|
||||
// when several declarations are waiting, the host applies the newest and acknowledges the
|
||||
// rest without applying them. Each of those is still reported, because silence reads as a
|
||||
// machine that ignored an instruction and "applied" would be a lie — this is the third word.
|
||||
Superseded string `json:"superseded,omitempty"`
|
||||
|
||||
// Carried are the machine's ports held by what this host raised from its own bundle.
|
||||
//
|
||||
// **So the mesh can assign around what it did not put here** (novox/hq ADR 0038). A node
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// A machine asked to be five things becomes the last one: what is already waiting supersedes what
|
||||
// arrived first, and everything set aside is named so it can be reported.
|
||||
func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) {
|
||||
deliveries := make(chan amqp.Delivery, 8)
|
||||
for _, id := range []string{"two", "three", "four"} {
|
||||
deliveries <- amqp.Delivery{Body: []byte(id)}
|
||||
}
|
||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("one")}, 50*time.Millisecond)
|
||||
if string(apply.Body) != "four" {
|
||||
t.Fatalf("applied %q, not the newest", apply.Body)
|
||||
}
|
||||
if len(superseded) != 3 || string(superseded[0].Body) != "one" || string(superseded[2].Body) != "three" {
|
||||
t.Fatalf("set aside %d: %v", len(superseded), superseded)
|
||||
}
|
||||
}
|
||||
|
||||
// One declaration with nothing behind it is applied as it always was, after the window.
|
||||
func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) {
|
||||
deliveries := make(chan amqp.Delivery, 1)
|
||||
began := time.Now()
|
||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("only")}, 30*time.Millisecond)
|
||||
if string(apply.Body) != "only" || len(superseded) != 0 {
|
||||
t.Fatalf("got %q with %d set aside", apply.Body, len(superseded))
|
||||
}
|
||||
if time.Since(began) < 30*time.Millisecond {
|
||||
t.Fatal("did not wait the window for a straggler")
|
||||
}
|
||||
}
|
||||
|
||||
// A straggler within the window is taken; one after it is the next push.
|
||||
func TestAStragglerWithinTheWindowIsTaken(t *testing.T) {
|
||||
deliveries := make(chan amqp.Delivery, 2)
|
||||
go func() {
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
deliveries <- amqp.Delivery{Body: []byte("late")}
|
||||
}()
|
||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("first")}, 100*time.Millisecond)
|
||||
if string(apply.Body) != "late" || len(superseded) != 1 {
|
||||
t.Fatalf("got %q with %d set aside", apply.Body, len(superseded))
|
||||
}
|
||||
}
|
||||
+61
-4
@@ -208,10 +208,13 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
return fmt.Errorf("cannot declare this node's queue %s: %w", queue, err)
|
||||
}
|
||||
|
||||
// One at a time. A declaration is applied to a machine, and applying two at once would race
|
||||
// on the same filesystem — so the broker holds the next one until this one is finished,
|
||||
// where it survives a restart.
|
||||
if err := channel.Qos(1, 0, false); err != nil {
|
||||
// Applying is one at a time — two at once would race on the same filesystem — but SEEING is
|
||||
// not: with a prefetch of one the host could never know that a newer declaration was already
|
||||
// waiting, and so applied every one of a backlog in turn, at the better part of a minute each,
|
||||
// becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged
|
||||
// deliveries lets it drain to the newest; each declaration still survives a restart on the
|
||||
// broker until it is acknowledged, which happens only after it is applied or set aside.
|
||||
if err := channel.Qos(drainDepth, 0, false); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -257,6 +260,15 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
if !ok {
|
||||
return errors.New("the broker stopped delivering")
|
||||
}
|
||||
// Whatever else is already waiting supersedes this one. Each set-aside declaration
|
||||
// is reported as such, then acknowledged unapplied.
|
||||
delivery, superseded := newest(deliveries, delivery, drainWindow)
|
||||
for _, old := range superseded {
|
||||
say("set aside a declaration: a newer one arrived with it")
|
||||
publishReport(ctx, channel, m, Report{Node: m.Node, Declared: declaredIn(old.Body),
|
||||
Superseded: declaredIn(delivery.Body)}, say, timeout)
|
||||
_ = old.Ack(false)
|
||||
}
|
||||
report := handle(ctx, m, apply, delivery)
|
||||
switch {
|
||||
case report.Refused != "":
|
||||
@@ -276,6 +288,51 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
}
|
||||
}
|
||||
|
||||
// drainDepth is how many declarations the host will hold unacknowledged while it looks for a newer
|
||||
// one; drainWindow is how long it waits for another to follow the one it has. Both small: a push is
|
||||
// rare and a backlog is the exception this exists for, not the shape of ordinary traffic.
|
||||
const (
|
||||
drainDepth = 16
|
||||
drainWindow = 750 * time.Millisecond
|
||||
)
|
||||
|
||||
// newest takes what is already waiting behind `first` and returns the last of them to apply, and
|
||||
// the rest to set aside. It waits `window` for a straggler after each arrival and no longer: a
|
||||
// declaration in flight from the mesh arrives within that; one that does not is the next push.
|
||||
func newest(deliveries <-chan amqp.Delivery, first amqp.Delivery, window time.Duration) (amqp.Delivery, []amqp.Delivery) {
|
||||
latest := first
|
||||
var superseded []amqp.Delivery
|
||||
for {
|
||||
select {
|
||||
case next, ok := <-deliveries:
|
||||
if !ok {
|
||||
return latest, superseded
|
||||
}
|
||||
superseded = append(superseded, latest)
|
||||
latest = next
|
||||
case <-time.After(window):
|
||||
return latest, superseded
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// declaredIn is the id a signed declaration carries, for a report about one that was not applied.
|
||||
// Empty if the message is not one — a forged or garbled message is refused by handleBody 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 {
|
||||
return ""
|
||||
}
|
||||
var d struct {
|
||||
Declared string `json:"declared"`
|
||||
}
|
||||
if err := json.Unmarshal(signed.Declaration, &d); err != nil {
|
||||
return ""
|
||||
}
|
||||
return d.Declared
|
||||
}
|
||||
|
||||
func handle(ctx context.Context, m Membership, apply Applier, delivery amqp.Delivery) Report {
|
||||
return handleBody(ctx, m, delivery.Body, apply)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user