diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index b8e738f..74266f1 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -331,7 +331,7 @@ func runApply(ctx context.Context, opts options, d *declaration.Declaration, sou return err } - report, updated, applyErr := apply.Apply(ctx, sys, d, known, apply.ExecRunner, func(line string) { + report, updated, applyErr := apply.Apply(ctx, sys, d, known, store.OriginCarried, apply.ExecRunner, func(line string) { if !opts.json { fmt.Println(line) } @@ -538,7 +538,7 @@ func runLink(ctx context.Context, opts options) error { Fingerprint: mine.Membership.Fingerprint, Password: mine.Membership.Password, Signer: mine.Membership.Signer, - }, apply, opts.timeout) + }, apply, func(line string) { fmt.Println(line) }, opts.timeout) } // applyDeclared applies a declaration that has already been proved to come from the mesh. @@ -569,7 +569,10 @@ func applyDeclared(ctx context.Context, opts options, raw []byte) link.Report { return link.Report{Refused: err.Error()} } - outcome, updated, applyErr := apply.Apply(ctx, built, declared, known, apply.ExecRunner, nil) + // Declared, not carried. A declaration from the mesh removes only what the mesh previously + // declared — never what this machine raised for itself from its bundle (04-ISSUES/010). + outcome, updated, applyErr := apply.Apply(ctx, built, declared, known, store.OriginDeclared, + apply.ExecRunner, nil) // Saved whichever way it went. Recording only on success would lose the footprint of a // failed apply, and that footprint is on the machine either way. diff --git a/examples/substrate-first-node.lock b/examples/substrate-first-node.lock index cd85237..65839e1 100644 --- a/examples/substrate-first-node.lock +++ b/examples/substrate-first-node.lock @@ -74,7 +74,7 @@ "command": ["docker", "run", "--rm", "--network", "container:mesh-store", "-e", "MESH_STORE_INVENTORY=postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable", "-e", "MESH_STORE_IDENTITY=postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable", - "192.0.2.250:5000/mesh-control@sha256:1dfcf6a879e16e671d4d6459271fd2630e30560afa1211508ab3994494786893", + "192.0.2.250:5000/mesh-control@sha256:b14e0a9765445b56f21aa8cf25be6a614b46974de710ddad92ee577074065858", "migrate"], "verify": ["sh", "-c", "docker exec mesh-store psql -U postgres -d inventory -tAc \"select to_regclass('public.node')\" | grep -qx node && docker exec mesh-store psql -U postgres -d identity -tAc \"select to_regclass('public.signing_key')\" | grep -qx signing_key"] }, @@ -108,7 +108,7 @@ "id": "control-plane", "type": "container", "name": "mesh-control", - "image": "192.0.2.250:5000/mesh-control@sha256:1dfcf6a879e16e671d4d6459271fd2630e30560afa1211508ab3994494786893", + "image": "192.0.2.250:5000/mesh-control@sha256:b14e0a9765445b56f21aa8cf25be6a614b46974de710ddad92ee577074065858", "network": "host", "args": ["serve"], "volumes": ["mesh-broker-tls:/broker-tls:ro"], diff --git a/internal/apply/apply.go b/internal/apply/apply.go index 876a107..8681670 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -85,6 +85,7 @@ func Apply( sys system.System, d *declaration.Declaration, known store.State, + origin string, run Runner, log func(string), ) (Report, store.State, error) { @@ -98,7 +99,7 @@ func Apply( declared[r.Identity()] = true } - for _, orphan := range known.Orphans(declared) { + for _, orphan := range known.Orphans(declared, origin) { action, detail, err := remove(ctx, sys, orphan, run) if err != nil { return report, known, &Error{Resource: orphan.ID, Err: err, Done: report} @@ -119,7 +120,8 @@ func Apply( // Only now. The record follows the fact, never leads it. known.Record(store.Applied{ - ID: resource.Identity(), Type: string(resource.Kind()), + Origin: origin, + ID: resource.Identity(), Type: string(resource.Kind()), Target: outcome.Target, AppliedAt: time.Now().UTC(), }) report.Outcomes = append(report.Outcomes, outcome) diff --git a/internal/apply/apply_test.go b/internal/apply/apply_test.go index bc819f8..0181cc3 100644 --- a/internal/apply/apply_test.go +++ b/internal/apply/apply_test.go @@ -39,7 +39,7 @@ func TestApplyingTwiceChangesNothingTheSecondTime(t *testing.T) { {"id":"f","type":"file","path":"`+dir+`/etc/a.conf","content":"hello\n","mode":"0640"} ]}`) - first, state, err := Apply(context.Background(), archHost(t), d, store.State{}, noServices, nil) + first, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -47,7 +47,7 @@ func TestApplyingTwiceChangesNothingTheSecondTime(t *testing.T) { t.Fatal("the first apply on an empty machine changed nothing") } - second, _, err := Apply(context.Background(), archHost(t), d, state, noServices, nil) + second, _, err := Apply(context.Background(), archHost(t), d, state, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -65,7 +65,7 @@ func TestADriftedMachineIsReturned(t *testing.T) { {"id":"f","type":"file","path":"`+path+`","content":"correct\n","mode":"0644"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, noServices, nil) + _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -73,7 +73,7 @@ func TestADriftedMachineIsReturned(t *testing.T) { t.Fatal(err) } - report, _, err := Apply(context.Background(), archHost(t), d, state, noServices, nil) + report, _, err := Apply(context.Background(), archHost(t), d, state, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -97,7 +97,7 @@ func TestADroppedResourceIsRemoved(t *testing.T) { {"id":"keep","type":"file","path":"`+keep+`","content":"a\n"}, {"id":"drop","type":"file","path":"`+drop+`","content":"b\n"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), both, store.State{}, noServices, nil) + _, state, err := Apply(context.Background(), archHost(t), both, store.State{}, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -105,7 +105,7 @@ func TestADroppedResourceIsRemoved(t *testing.T) { one := parse(t, `{"declaration":1,"resources":[ {"id":"keep","type":"file","path":"`+keep+`","content":"a\n"} ]}`) - report, state, err := Apply(context.Background(), archHost(t), one, state, noServices, nil) + report, state, err := Apply(context.Background(), archHost(t), one, state, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -137,7 +137,7 @@ func TestNothingTheHostDidNotCreateIsTouched(t *testing.T) { d := parse(t, `{"declaration":1,"resources":[ {"id":"ours","type":"file","path":"`+filepath.Join(dir, "ours.conf")+`","content":"a\n"} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, noServices, nil); err != nil { + if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, noServices, nil); err != nil { t.Fatal(err) } @@ -156,7 +156,7 @@ func TestARenameToTheSamePathDoesNotDeleteTheNewFile(t *testing.T) { before := parse(t, `{"declaration":1,"resources":[ {"id":"old","type":"file","path":"`+path+`","content":"old\n"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), before, store.State{}, noServices, nil) + _, state, err := Apply(context.Background(), archHost(t), before, store.State{}, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -164,7 +164,7 @@ func TestARenameToTheSamePathDoesNotDeleteTheNewFile(t *testing.T) { after := parse(t, `{"declaration":1,"resources":[ {"id":"new","type":"file","path":"`+path+`","content":"new\n"} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), after, state, noServices, nil); err != nil { + if _, _, err := Apply(context.Background(), archHost(t), after, state, store.OriginCarried, noServices, nil); err != nil { t.Fatal(err) } @@ -192,7 +192,7 @@ func TestAFailedStepFailsTheApply(t *testing.T) { {"id":"never","type":"file","path":"`+filepath.Join(dir, "never.conf")+`","content":"b\n"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, noServices, nil) + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, noServices, nil) if err == nil { t.Fatal("an impossible resource did not fail the apply") } @@ -225,7 +225,7 @@ func TestNothingIsRecordedUntilItWorked(t *testing.T) { {"id":"doomed","type":"directory","path":"`+blocker+`"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, noServices, nil) + _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, noServices, nil) if err == nil { t.Fatal("expected a failure") } @@ -244,7 +244,7 @@ func TestAModeIsMaintainedNotJustSet(t *testing.T) { {"id":"f","type":"file","path":"`+path+`","content":"s\n","mode":"0600"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, noServices, nil) + _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -252,7 +252,7 @@ func TestAModeIsMaintainedNotJustSet(t *testing.T) { t.Fatal(err) } - report, _, err := Apply(context.Background(), archHost(t), d, state, noServices, nil) + report, _, err := Apply(context.Background(), archHost(t), d, state, store.OriginCarried, noServices, nil) if err != nil { t.Fatal(err) } @@ -283,7 +283,7 @@ func TestAServiceIsReadBackNotAssumed(t *testing.T) { d := parse(t, `{"declaration":1,"resources":[ {"id":"s","type":"service","unit":"doomed.service","state":"running"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err == nil { t.Fatal("a service that died immediately was reported as running") } @@ -299,7 +299,7 @@ func TestAnUnknownServiceStateIsRefusedNotGuessed(t *testing.T) { d := parse(t, `{"declaration":1,"resources":[ {"id":"s","type":"service","unit":"odd.service","state":"running"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err == nil || !strings.Contains(err.Error(), "neither running nor stopped") { t.Errorf("an unrecognised service state was not refused: %v", err) } @@ -323,7 +323,7 @@ func TestADroppedServiceIsStoppedNotDeleted(t *testing.T) { {"id":"other","type":"file","path":"`+filepath.Join(t.TempDir(), "a")+`","content":"a\n"} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), d, state, run, nil); err != nil { + if _, _, err := Apply(context.Background(), archHost(t), d, state, store.OriginCarried, run, nil); err != nil { t.Fatal(err) } joined := strings.Join(commands, "; ") @@ -349,7 +349,7 @@ func TestAUnitThatDoesNotExistIsNotStopped(t *testing.T) { {"id":"s","type":"service","unit":"never-installed.service","state":"stopped"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, absent, nil) + _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, absent, nil) if err == nil { t.Fatal("a unit that does not exist was reported as satisfactorily stopped") } @@ -370,7 +370,7 @@ func TestAMaskedUnitIsRefused(t *testing.T) { d := parse(t, `{"declaration":1,"resources":[ {"id":"s","type":"service","unit":"masked.service","state":"running"} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, masked, nil); err == nil { + if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, masked, nil); err == nil { t.Fatal("a masked unit was accepted") } } @@ -398,7 +398,7 @@ func TestForgettingAUnitThatIsGoneDoesNotStrandTheNode(t *testing.T) { {"id":"f","type":"file","path":"`+filepath.Join(t.TempDir(), "a")+`","content":"a\n"} ]}`) - report, state, err := Apply(context.Background(), archHost(t), d, known, run, nil) + report, state, err := Apply(context.Background(), archHost(t), d, known, store.OriginCarried, run, nil) if err != nil { t.Fatalf("a vanished unit stranded the apply: %v", err) } @@ -442,7 +442,7 @@ func TestABrokenPackageDatabaseIsNotReadAsNotInstalled(t *testing.T) { {"id":"rt","type":"package","package":"docker"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err == nil { t.Fatal("a broken package database was read as 'not installed'") } @@ -463,7 +463,7 @@ func TestAnInstalledPackageIsNotReinstalled(t *testing.T) { {"id":"rt","type":"package","package":"docker"} ]}`) - report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err != nil { t.Fatalf("apply failed: %v", err) } @@ -493,7 +493,7 @@ func TestAPackageIsNeverUninstalled(t *testing.T) { {"id":"f","type":"file","path":"`+filepath.Join(t.TempDir(), "a")+`","content":"a\n"} ]}`) - report, state, err := Apply(context.Background(), archHost(t), d, known, run, nil) + report, state, err := Apply(context.Background(), archHost(t), d, known, store.OriginCarried, run, nil) if err != nil { t.Fatalf("dropping a package stranded the apply: %v", err) } @@ -523,7 +523,7 @@ func TestAnActionThatIsAlreadyTrueDoesNotRun(t *testing.T) { {"id":"db","type":"action","command":["create-db","mesh"],"verify":["has-db","mesh"]} ]}`) - report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err != nil { t.Fatalf("apply failed: %v", err) } @@ -549,7 +549,7 @@ func TestAnActionThatSucceedsAndDoesNothingFails(t *testing.T) { {"id":"db","type":"action","command":["create-db","mesh"],"verify":["has-db","mesh"]} ]}`) - _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err == nil { t.Fatal("an action that reported success and did nothing was accepted") } @@ -576,7 +576,7 @@ func TestAnActionRunsInsideTheContainerItNames(t *testing.T) { {"id":"db","type":"action","in":"store","command":["createdb","mesh"],"verify":["psql","-lqt"]} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil); err != nil { + if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil); err != nil { t.Fatalf("apply failed: %v", err) } if !sawExec { @@ -603,7 +603,7 @@ func TestAContainerThatExitsImmediatelyFailsTheApply(t *testing.T) { {"id":"store","type":"container","name":"store","image":"`+pinned+`"} ]}`) - _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, state, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err == nil { t.Fatal("a container that exited immediately was reported as applied") } @@ -645,7 +645,7 @@ func TestAContainerWhoseDeclarationChangedIsReplaced(t *testing.T) { return "", nil } - report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err != nil { t.Fatalf("apply failed: %v", err) } @@ -675,7 +675,7 @@ func TestAContainerThatMatchesIsLeftAlone(t *testing.T) { return "", nil } - report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err != nil { t.Fatalf("apply failed: %v", err) } @@ -734,7 +734,7 @@ func TestAServiceIsEnabledAtBootWhenAsked(t *testing.T) { {"id":"rt","type":"service","unit":"docker.service","state":"running","boot":"enabled"} ]}`) - report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, systemctlStub(t, "loaded", "inactive", "disabled", &verbs), nil) if err != nil { t.Fatalf("apply failed: %v", err) @@ -754,7 +754,7 @@ func TestBootIsEnabledBeforeTheUnitIsStarted(t *testing.T) { d := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"rt","type":"service","unit":"docker.service","state":"running","boot":"enabled"} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, + if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, systemctlStub(t, "loaded", "inactive", "disabled", &verbs), nil); err != nil { t.Fatal(err) } @@ -769,7 +769,7 @@ func TestAlreadyEnabledAndRunningIsUnchanged(t *testing.T) { {"id":"rt","type":"service","unit":"docker.service","state":"running","boot":"enabled"} ]}`) - report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, + report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, systemctlStub(t, "loaded", "active", "enabled", &verbs), nil) if err != nil { t.Fatalf("apply failed: %v", err) @@ -789,7 +789,7 @@ func TestOmittingBootLeavesItAlone(t *testing.T) { d := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"rt","type":"service","unit":"docker.service","state":"running"} ]}`) - if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, + if _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, systemctlStub(t, "loaded", "inactive", "enabled", &verbs), nil); err != nil { t.Fatal(err) } @@ -809,7 +809,7 @@ func TestAStaticUnitCannotBeEnabled(t *testing.T) { {"id":"rt","type":"service","unit":"dbus.socket","state":"running","boot":"enabled"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, systemctlStub(t, "loaded", "active", "static", &verbs), nil) if err == nil { t.Fatal("a static unit was accepted as enable-able") @@ -824,7 +824,7 @@ func TestAnUnknownBootStateIsRefusedNotGuessed(t *testing.T) { d := parseTrusted(t, `{"declaration":1,"resources":[ {"id":"rt","type":"service","unit":"x.service","state":"running","boot":"enabled"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, systemctlStub(t, "loaded", "active", "indirect", &verbs), nil) if err == nil { t.Fatal("an unrecognised boot state was guessed at instead of refused") @@ -900,7 +900,7 @@ func TestAContainerUsesTheRuntimeTheMachineHas(t *testing.T) { // It will fail at read-back — the stub never reports it running — and what matters is // WHICH binary it used getting there. - _, _, _ = Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, _, _ = Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) for _, c := range calledWith { if c != "podman" { @@ -922,7 +922,7 @@ func TestNoRuntimeIsSaidPlainly(t *testing.T) { {"id":"store","type":"container","name":"store","image":"`+pinned+`"} ]}`) - _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, run, nil) + _, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil) if err == nil { t.Fatal("a machine with no container runtime applied a container") } diff --git a/internal/link/run.go b/internal/link/run.go index f190975..db68252 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -31,13 +31,20 @@ type Membership struct { // Applier is what the host does with a declaration that has been proved to come from the mesh. type Applier func(ctx context.Context, declaration []byte) Report +// Announce is how the link says what is happening, so a node running unattended leaves an +// account of it. Nil is allowed and means say nothing. +type Announce func(string) + // Run holds the link open, applying what arrives and reporting what happened. // // Outbound only, and nothing listens on this machine. The connection is the node's presence in // the mesh: while it is up the node is enrolled, and while it is down the node is disconnected — // which is an ordinary situation and not a failure, so this returns rather than panicking and // leaves restarting to whatever supervises it. -func Run(ctx context.Context, m Membership, apply Applier, timeout time.Duration) error { +func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error { + if say == nil { + say = func(string) {} + } config, err := PinnedConfig(m.Fingerprint) if err != nil { return err @@ -84,6 +91,18 @@ func Run(ctx context.Context, m Membership, apply Applier, timeout time.Duration } closed := conn.NotifyClose(make(chan *amqp.Error, 1)) + // Published mandatory, so the broker hands back anything it cannot route rather than + // dropping it. Without this a report goes to an exchange with no matching binding, the + // publisher is told nothing, and the mesh believes this node never answered while the node + // believes it did — which is what happened when `report` was left unbound on the other side. + returned := channel.NotifyReturn(make(chan amqp.Return, 4)) + go func() { + for r := range returned { + say(fmt.Sprintf("the broker could not route this node's %s: %s (%d %s)", + r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText)) + } + }() + for { select { case <-ctx.Done(): @@ -95,7 +114,15 @@ func Run(ctx context.Context, m Membership, apply Applier, timeout time.Duration return errors.New("the broker stopped delivering") } report := handle(ctx, m, apply, delivery) - publishReport(ctx, channel, m, report, timeout) + switch { + case report.Refused != "": + say("refused a declaration: " + report.Refused) + case len(report.Failed) > 0: + say(fmt.Sprintf("applied %d and failed: %v", len(report.Applied), report.Failed)) + default: + say(fmt.Sprintf("applied %d resource(s)", len(report.Applied))) + } + publishReport(ctx, channel, m, report, say, timeout) // Acknowledged after the report is published. A node that dies between applying and // reporting leaves the declaration on the broker and applies it again on return, // which is safe because applying is reconciliation — it converges rather than @@ -127,14 +154,21 @@ func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) R } func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report, - timeout time.Duration) { + say Announce, timeout time.Duration) { report.Node = m.Node body, err := json.Marshal(report) if err != nil { + say("cannot encode this node's own report: " + err.Error()) return } publish, cancel := context.WithTimeout(ctx, timeout) defer cancel() - _ = channel.PublishWithContext(publish, Exchange, KeyReport, false, false, - amqp.Publishing{ContentType: "application/json", Body: body}) + + // Said rather than swallowed. A report that fails to publish leaves the mesh believing this + // node never answered, while the node believes it did — and the two would go on disagreeing + // with nothing anywhere saying so. That shape of fault is the one this project keeps finding. + if err := channel.PublishWithContext(publish, Exchange, KeyReport, true, false, + amqp.Publishing{ContentType: "application/json", Body: body}); err != nil { + say(fmt.Sprintf("applied, and could not tell the mesh: %v", err)) + } } diff --git a/internal/store/store.go b/internal/store/store.go index e1f5165..060d3f3 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -33,6 +33,16 @@ const DefaultPath = "/var/lib/mesh-host/state.json" type Applied struct { ID string `json:"id"` Type string `json:"type"` + // Origin is who asked for this: the bundle this host carries, or the mesh. + // + // Recorded because the two must not remove each other. A node raises its own substrate from + // the bundle before any mesh exists, then enrols and is sent declarations — and a + // declaration naming two resources would otherwise remove the store, the broker and the + // control plane, which is 04-ISSUES/010 and happened on the first end-to-end run. + // + // Empty means carried, for state written before this field existed: everything a host had + // applied at that point came from its bundle. + Origin string `json:"origin,omitempty"` // Target is what was changed — a path, a unit — so removal knows what to undo without // re-reading a declaration that may no longer exist. Target string `json:"target"` @@ -165,12 +175,39 @@ func (s *State) Forget(id string) { // Reverse order because undoing in the order things were made undoes a directory before the // file inside it. Reversing is the only ordering the host can derive without deciding // anything, which is the line novox/hq ADR 0005 draws. -func (s State) Orphans(declared map[string]bool) []Applied { +func (s State) Orphans(declared map[string]bool, origin string) []Applied { var out []Applied for i := len(s.Resources) - 1; i >= 0; i-- { - if !declared[s.Resources[i].ID] { - out = append(out, s.Resources[i]) + r := s.Resources[i] + // Only this origin's own. A mesh declaration says nothing about what the bundle raised, + // and a bundle says nothing about what the mesh assigned — so neither may remove the + // other's by omission, which is the only way either could express removal. + if originOf(r) != origin { + continue + } + if !declared[r.ID] { + out = append(out, r) } } return out } + +// Origins a resource can have. +const ( + // Carried is the bundle this host was built with. + OriginCarried = "carried" + // Declared is the mesh, over the link. + OriginDeclared = "declared" +) + +// originOf reads a record's origin, treating absence as carried. +// +// State written before origins existed was all bundle-applied: a host had no other way to be +// told anything. Guessing wrong in the other direction would have a first upgrade remove the +// substrate, which is the fault this field exists to prevent. +func originOf(r Applied) string { + if r.Origin == "" { + return OriginCarried + } + return r.Origin +} diff --git a/internal/store/store_test.go b/internal/store/store_test.go index eab9a02..3867be9 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -116,7 +116,7 @@ func TestOrphansAreWhatWasAppliedAndIsNoLongerDeclared(t *testing.T) { {ID: "file", Type: "file", Target: "/etc/mesh/a.conf"}, {ID: "kept", Type: "file", Target: "/etc/mesh/b.conf"}, }} - orphans := s.Orphans(map[string]bool{"kept": true}) + orphans := s.Orphans(map[string]bool{"kept": true}, OriginCarried) if len(orphans) != 2 { t.Fatalf("expected two orphans, got %d: %+v", len(orphans), orphans) @@ -130,7 +130,59 @@ func TestOrphansAreWhatWasAppliedAndIsNoLongerDeclared(t *testing.T) { func TestNothingIsAnOrphanWhenEverythingIsDeclared(t *testing.T) { s := State{Resources: []Applied{{ID: "a", Type: "file", Target: "/a"}}} - if got := s.Orphans(map[string]bool{"a": true}); len(got) != 0 { + if got := s.Orphans(map[string]bool{"a": true}, OriginCarried); len(got) != 0 { t.Errorf("a declared resource was treated as an orphan: %+v", got) } } + +func TestADeclarationDoesNotOrphanWhatTheBundleRaised(t *testing.T) { + // 04-ISSUES/010. A first node raises its substrate from the bundle it carries, then enrols + // and is sent a declaration naming two resources. Before origins, that removed the store, the + // broker and the control plane that had sent it — the mesh deleting itself over the link the + // message arrived on, in under a second, on the first end-to-end run. + s := State{Resources: []Applied{ + {ID: "store", Type: "container", Target: "mesh-store", Origin: OriginCarried}, + {ID: "broker", Type: "container", Target: "mesh-broker", Origin: OriginCarried}, + {ID: "greeting", Type: "directory", Target: "/var/lib/demo", Origin: OriginDeclared}, + }} + + // The mesh declares nothing at all. Everything it previously declared is an orphan; nothing + // the bundle raised is. + orphans := s.Orphans(map[string]bool{}, OriginDeclared) + if len(orphans) != 1 || orphans[0].ID != "greeting" { + var got []string + for _, o := range orphans { + got = append(got, o.ID) + } + t.Fatalf("a declaration would remove %v; it may only remove what the mesh declared", got) + } +} + +func TestTheBundleDoesNotOrphanWhatTheMeshDeclared(t *testing.T) { + // The same rule the other way. A host reconciling its carried bundle must not remove what the + // mesh assigned to this node, or every restart would undo the node's actual work. + s := State{Resources: []Applied{ + {ID: "store", Type: "container", Target: "mesh-store", Origin: OriginCarried}, + {ID: "workload", Type: "container", Target: "some-app", Origin: OriginDeclared}, + }} + + orphans := s.Orphans(map[string]bool{"store": true}, OriginCarried) + if len(orphans) != 0 { + t.Errorf("reconciling the bundle would remove %s, which the mesh declared", orphans[0].ID) + } +} + +func TestStateWrittenBeforeOriginsExistedIsTreatedAsCarried(t *testing.T) { + // Every resource a host had applied before this field existed came from its bundle, because + // there was no other way to tell it anything. Guessing the other way would have the first + // declaration remove the substrate — which is the fault this exists to prevent, arriving + // through the upgrade that fixes it. + s := State{Resources: []Applied{{ID: "store", Type: "container", Target: "mesh-store"}}} + + if got := s.Orphans(map[string]bool{}, OriginDeclared); len(got) != 0 { + t.Errorf("a declaration would remove %s, recorded before origins existed", got[0].ID) + } + if got := s.Orphans(map[string]bool{}, OriginCarried); len(got) != 1 { + t.Error("the bundle cannot remove its own resource, so nothing could ever remove it") + } +}