diff --git a/examples/foundation-first-node.lock b/examples/foundation-first-node.lock index cbdf3e3..5903648 100644 --- a/examples/foundation-first-node.lock +++ b/examples/foundation-first-node.lock @@ -46,6 +46,37 @@ "state": "running", "boot": "enabled" }, + // **A filter before anything listens** (novox/hq issue 054, ADR 0088). The store and the + // broker are adopted as modules later and so bind to every interface from the moment they + // start; the packet filter that governs who may reach them is a module too, installed a + // dozen steps later. Between the two, a control-node facing the network had its store and + // its bus open to anyone who could reach the machine. So the foundation carries a filter of + // its own — the same table the filter module will replace wholesale once it can derive one: + // drop by default, keep loopback, replies, ssh and the mesh's own ports (the bus a node + // enrols over, the registry a node pulls from), and let the container runtime's own + // networks through the forward chain so containers keep working. A published container port + // is forwarded, never input (issue 047), which is why the forward chain is where the store's + // and broker's ports are refused from outside. + { + "id": "base-filter-package", + "type": "package", + "package": "nftables" + }, + { + "id": "base-filter", + "type": "file", + "path": "/etc/nftables.conf", + "mode": "0644", + "content": "#!/usr/sbin/nft -f\n# the foundation's own filter, until the mesh derives one (novox/hq issue 054)\ntable inet mesh {}\ndelete table inet mesh\n\ntable inet mesh {\n\tchain input {\n\t\ttype filter hook input priority filter; policy drop;\n\t\tct state established,related accept\n\t\tct state invalid drop\n\t\tiif lo accept\n\t\ticmp type echo-request accept\n\t\ticmpv6 type { echo-request, nd-neighbor-solicit, nd-neighbor-advert, nd-router-advert } accept\n\t\t# ssh, from anywhere — never closed\n\t\ttcp dport 22 accept\n\t}\n\tchain output {\n\t\ttype filter hook output priority filter; policy accept;\n\t}\n\tchain forward {\n\t\ttype filter hook forward priority filter; policy drop;\n\t\tct state established,related accept\n\t\tct state invalid drop\n\t\t# the container runtime's bridge networks, and the networks its compose files are given\n\t\tip saddr 172.16.0.0/12 accept\n\t\tip saddr 192.168.128.0/17 accept\n\t\t# the mesh's own: the bus a node enrols over, the registry a node pulls from\n\t\tct original proto-dst 5671 accept\n\t\tct original proto-dst 5000 accept\n\t}\n}\n" + }, + { + "id": "base-filter-loaded", + "type": "service", + "unit": "nftables.service", + "state": "running", + "boot": "enabled", + "restart-on": ["base-filter"] + }, { "id": "store", "type": "container", diff --git a/internal/apply/apply.go b/internal/apply/apply.go index a440d23..893aa17 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -470,6 +470,18 @@ func applyFile(r *declaration.File, previous store.Applied, unseal Unseal) (Outc if readErr != nil && !errors.Is(readErr, os.ErrNotExist) { return out, readErr } + // A seed that is already there is left exactly as it is — whatever has grown in it since is + // not the mesh's to put back (novox/hq issue 035). What this host records is what it once + // wrote, so a later declaration that changes the seed is not mistaken for drift either. + if r.CreateOnce && existed { + out.wrote = previous.Wrote + if out.wrote == "" { + out.wrote = digestOf(string(existing)) + } + out.Action = "kept" + out.Detail = "created once, and present; what is in it now is not the mesh's to change" + return out, nil + } var beforeMode os.FileMode if existed { diff --git a/internal/apply/apply_test.go b/internal/apply/apply_test.go index 3dfc524..52a179c 100644 --- a/internal/apply/apply_test.go +++ b/internal/apply/apply_test.go @@ -1558,3 +1558,40 @@ func TestAContainerStaleFromAnEarlierApplyIsReplaced(t *testing.T) { "(removed=%v created=%v)", removed, created) } } + +// A seed is written once. What grows in it afterwards is somebody else's work the mesh asked for, +// and a reconcile leaves it alone — content, mode and owner — and says so (novox/hq issue 035). +func TestASeedIsCreatedOnceAndWhatGrowsInItIsKept(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "acl.conf") + d := parse(t, fmt.Sprintf(`{"declaration":1,"resources":[ + {"id":"acl","type":"file","path":%q,"content":"user default on\n","mode":"0600","create-once":true} + ]}`, path)) + + report, state, err := Apply(context.Background(), archHost(t), d, store.State{}, + store.OriginCarried, noServices, nil, nil) + if err != nil { + t.Fatal(err) + } + if got := report.Outcomes[0].Action; got != "created" { + t.Fatalf("first apply: %q", got) + } + + // The program persists into it. + if err := os.WriteFile(path, []byte("user default on\nuser app-one on >secret\n"), 0o600); err != nil { + t.Fatal(err) + } + + report, _, err = Apply(context.Background(), archHost(t), d, state, + store.OriginCarried, noServices, nil, nil) + if err != nil { + t.Fatal(err) + } + if got := report.Outcomes[0].Action; got != "kept" { + t.Fatalf("second apply: %q — a seed was reconciled", got) + } + grown, _ := os.ReadFile(path) + if string(grown) != "user default on\nuser app-one on >secret\n" { + t.Fatalf("what grew in the seed was wiped: %q", grown) + } +} diff --git a/internal/bootstrap/rootsecrets_test.go b/internal/bootstrap/rootsecrets_test.go index c23f864..e31814f 100644 --- a/internal/bootstrap/rootsecrets_test.go +++ b/internal/bootstrap/rootsecrets_test.go @@ -165,3 +165,49 @@ func TestTheBrokerAdminMarkerHasALineEnding(t *testing.T) { } } } + +// The foundation filters before anything listens: nftables is in place before the store, and its +// rules refuse the store's and broker's client ports from outside while keeping ssh, the bus a node +// enrols over and the registry a node pulls from (novox/hq issue 054). +func TestTheFoundationFiltersBeforeAnythingListens(t *testing.T) { + template, err := os.ReadFile("../../examples/foundation-first-node.lock") + if err != nil { + t.Skip("no example bundle beside this checkout") + } + r, err := Rewrite(template, "sha256:"+strings.Repeat("ab", 32)) + if err != nil { + t.Fatal(err) + } + if _, err := RewriteRoot(&r, RootCredentials{Store: "s", Broker: "b"}); err != nil { + t.Fatal(err) + } + var filterAt, loadedAt, storeAt, brokerAt = -1, -1, -1, -1 + var rules string + for i, res := range r.Declaration.Resources { + switch res.Identity() { + case "base-filter": + filterAt = i + rules = res.(*declaration.File).Content + case "base-filter-loaded": + loadedAt = i + case "store": + storeAt = i + case "broker": + brokerAt = i + } + } + if !(filterAt >= 0 && filterAt < loadedAt && loadedAt < storeAt && storeAt < brokerAt) { + t.Fatalf("order: filter %d loaded %d store %d broker %d", filterAt, loadedAt, storeAt, brokerAt) + } + for _, want := range []string{"policy drop", "tcp dport 22 accept", "ct original proto-dst 5671 accept", + "ct original proto-dst 5000 accept", "ip saddr 172.16.0.0/12 accept", "table inet mesh"} { + if !strings.Contains(rules, want) { + t.Errorf("the base filter lacks %q", want) + } + } + for _, mustNot := range []string{"5432", "5672"} { + if strings.Contains(rules, mustNot) { + t.Errorf("the base filter names %s, which must be reached from the machine and its containers only", mustNot) + } + } +} diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index d1ffc9c..f4831fb 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -135,6 +135,19 @@ type File struct { Content string `json:"content"` Mode string `json:"mode,omitempty"` + // CreateOnce says the content is a seed: written when the file is absent, and left alone — + // content, mode and owner — whenever it is present. + // + // **Two intentions had one vocabulary** (novox/hq issue 035, ADR 0087). "This file has this + // content, for ever" is what an ordinary file says, and the host holds the machine to it. A + // module that needs a file to exist before a program first starts — an access list the + // program then persists into, a bootstrap configuration it rewrites — needs the other thing, + // and with only the first available, every reconcile restored the seed behind the running + // program and erased what had grown in it, reporting success. What grows in a seeded file + // is somebody else's work the mesh asked for; the mesh removes nothing it did not create + // (ADR 0030), and it does not overwrite that either. + CreateOnce bool `json:"create-once,omitempty"` + // Sealed is content encrypted to this node's sealing key, for a file the mesh must deliver // without being able to read. // diff --git a/internal/link/messages.go b/internal/link/messages.go index 67b63c5..6a48e41 100644 --- a/internal/link/messages.go +++ b/internal/link/messages.go @@ -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 diff --git a/internal/link/newest_test.go b/internal/link/newest_test.go new file mode 100644 index 0000000..7ad6cbb --- /dev/null +++ b/internal/link/newest_test.go @@ -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)) + } +} diff --git a/internal/link/run.go b/internal/link/run.go index dc62e5c..247d015 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -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) }