Report what an adopted node holds, its firewall and what is reachable, and speak unasked when that changes (hq ADR 0100)
This commit is contained in:
+99
-3
@@ -20,6 +20,7 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"syscall"
|
"syscall"
|
||||||
"text/tabwriter"
|
"text/tabwriter"
|
||||||
"time"
|
"time"
|
||||||
@@ -27,10 +28,12 @@ import (
|
|||||||
"github.com/novox/mesh-host/internal/apply"
|
"github.com/novox/mesh-host/internal/apply"
|
||||||
"github.com/novox/mesh-host/internal/bundle"
|
"github.com/novox/mesh-host/internal/bundle"
|
||||||
"github.com/novox/mesh-host/internal/declaration"
|
"github.com/novox/mesh-host/internal/declaration"
|
||||||
|
"github.com/novox/mesh-host/internal/firewall"
|
||||||
"github.com/novox/mesh-host/internal/identity"
|
"github.com/novox/mesh-host/internal/identity"
|
||||||
"github.com/novox/mesh-host/internal/inventory"
|
"github.com/novox/mesh-host/internal/inventory"
|
||||||
"github.com/novox/mesh-host/internal/link"
|
"github.com/novox/mesh-host/internal/link"
|
||||||
"github.com/novox/mesh-host/internal/profile"
|
"github.com/novox/mesh-host/internal/profile"
|
||||||
|
"github.com/novox/mesh-host/internal/reachable"
|
||||||
"github.com/novox/mesh-host/internal/store"
|
"github.com/novox/mesh-host/internal/store"
|
||||||
"github.com/novox/mesh-host/internal/system"
|
"github.com/novox/mesh-host/internal/system"
|
||||||
"github.com/novox/mesh-host/internal/upgrade"
|
"github.com/novox/mesh-host/internal/upgrade"
|
||||||
@@ -437,6 +440,23 @@ func enrol(ctx context.Context, opts options) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// An adopted node keeps the firewall it was found with (novox/hq ADR 0100), so a host that
|
||||||
|
// cannot speak that firewall must say so now — before the mesh records a node it could never
|
||||||
|
// open anything on.
|
||||||
|
if token.Adopted {
|
||||||
|
kind, name, err := firewall.Detect(ctx, apply.ExecRunner)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if kind == firewall.Unsupported {
|
||||||
|
return fmt.Errorf(
|
||||||
|
"this token joins this machine adopted, keeping the firewall found on it, and it is "+
|
||||||
|
"filtered by %s, which no host speaks yet. Nothing was enrolled", name)
|
||||||
|
}
|
||||||
|
fmt.Printf("joining adopted: what is on this machine is kept, and its firewall (%s) stays in force\n",
|
||||||
|
string(kind))
|
||||||
|
}
|
||||||
|
|
||||||
fmt.Printf("token for broker %s\n", token.Broker)
|
fmt.Printf("token for broker %s\n", token.Broker)
|
||||||
fmt.Printf(" pinned certificate %s\n", token.Fingerprint)
|
fmt.Printf(" pinned certificate %s\n", token.Fingerprint)
|
||||||
fmt.Printf(" signing key %s\n",
|
fmt.Printf(" signing key %s\n",
|
||||||
@@ -616,7 +636,22 @@ func runLink(ctx context.Context, opts options) error {
|
|||||||
// new declarations; this holds the machine in the last one whether the link is up or not. A
|
// new declarations; this holds the machine in the last one whether the link is up or not. A
|
||||||
// laptop shut for a week comes back and reconciles — it does not come back and ask what it is
|
// laptop shut for a week comes back and reconciles — it does not come back and ask what it is
|
||||||
// (novox/hq ADR 0004).
|
// (novox/hq ADR 0004).
|
||||||
go holdTheMachine(ctx, opts, mine, say, sched)
|
// Reports a reconcile has to make unasked — what an adopted node holds changed, or its
|
||||||
|
// firewall did — go out over the link when it is up (novox/hq ADR 0100).
|
||||||
|
outbox := make(chan link.Report, 1)
|
||||||
|
watch := &adoptionWatch{}
|
||||||
|
applier = watch.noting(applier)
|
||||||
|
go holdTheMachine(ctx, opts, mine, say, sched, func(r link.Report) {
|
||||||
|
if !watch.changed(r) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-outbox:
|
||||||
|
// An older one nobody has published yet; this one says everything it did.
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
outbox <- r
|
||||||
|
})
|
||||||
|
|
||||||
return link.HoldRoused(ctx, link.Membership{
|
return link.HoldRoused(ctx, link.Membership{
|
||||||
Node: mine.Node,
|
Node: mine.Node,
|
||||||
@@ -624,7 +659,46 @@ func runLink(ctx context.Context, opts options) error {
|
|||||||
Fingerprint: mine.Membership.Fingerprint,
|
Fingerprint: mine.Membership.Fingerprint,
|
||||||
Password: mine.Membership.Password,
|
Password: mine.Membership.Password,
|
||||||
Signer: mine.Membership.Signer,
|
Signer: mine.Membership.Signer,
|
||||||
}, applier, say, opts.timeout, rousedBySignal(ctx))
|
}, applier, say, opts.timeout, rousedBySignal(ctx), outbox)
|
||||||
|
}
|
||||||
|
|
||||||
|
// adoptionWatch remembers what the node last said about what it holds and its firewall, so a
|
||||||
|
// reconcile speaks unasked only when that changed.
|
||||||
|
type adoptionWatch struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
last string
|
||||||
|
}
|
||||||
|
|
||||||
|
// fingerprint is what a report says about adoption: each hold and whether it changed, and the
|
||||||
|
// firewall.
|
||||||
|
func adoptionFingerprint(r link.Report) string {
|
||||||
|
parts := []string{"firewall=" + r.Firewall}
|
||||||
|
for _, h := range r.Held {
|
||||||
|
parts = append(parts, h.ID+"="+h.Changed)
|
||||||
|
}
|
||||||
|
sort.Strings(parts[1:])
|
||||||
|
return strings.Join(parts, "\n")
|
||||||
|
}
|
||||||
|
|
||||||
|
// changed records a report and says whether it differs from the last one that went out.
|
||||||
|
func (w *adoptionWatch) changed(r link.Report) bool {
|
||||||
|
w.mu.Lock()
|
||||||
|
defer w.mu.Unlock()
|
||||||
|
now := adoptionFingerprint(r)
|
||||||
|
if now == w.last {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
w.last = now
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// noting wraps the applier, so a report the link publishes after a delivery counts as said.
|
||||||
|
func (w *adoptionWatch) noting(apply link.Applier) link.Applier {
|
||||||
|
return func(ctx context.Context, raw, signature []byte) link.Report {
|
||||||
|
r := apply(ctx, raw, signature)
|
||||||
|
w.changed(r)
|
||||||
|
return r
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// rousedBySignal is the machine telling this process that its link is probably stale.
|
// rousedBySignal is the machine telling this process that its link is probably stale.
|
||||||
@@ -672,7 +746,7 @@ func rousedBySignal(ctx context.Context) link.Roused {
|
|||||||
const ReconcileEvery = 5 * time.Minute
|
const ReconcileEvery = 5 * time.Minute
|
||||||
|
|
||||||
func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, say link.Announce,
|
func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, say link.Announce,
|
||||||
sched *apply.Scheduler) {
|
sched *apply.Scheduler, publish func(link.Report)) {
|
||||||
ticker := time.NewTicker(ReconcileEvery)
|
ticker := time.NewTicker(ReconcileEvery)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
@@ -695,6 +769,12 @@ func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, s
|
|||||||
}
|
}
|
||||||
|
|
||||||
report := applyDeclared(ctx, opts, declared, sched)
|
report := applyDeclared(ctx, opts, declared, sched)
|
||||||
|
// A reconcile is otherwise silent. On an adopted node it speaks when what it holds or
|
||||||
|
// its firewall changed, because that is how a predecessor still writing is caught
|
||||||
|
// (novox/hq ADR 0100); publish decides whether anything did.
|
||||||
|
if publish != nil && report.Refused == "" && (len(report.Held) > 0 || report.Firewall != "") {
|
||||||
|
publish(report)
|
||||||
|
}
|
||||||
switch {
|
switch {
|
||||||
case report.Refused != "":
|
case report.Refused != "":
|
||||||
say("what this node was last told no longer applies: " + report.Refused)
|
say("what this node was last told no longer applies: " + report.Refused)
|
||||||
@@ -761,6 +841,22 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D
|
|||||||
}
|
}
|
||||||
|
|
||||||
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
|
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
|
||||||
|
// What this node found and holds, its firewall, and what is reachable on it — so an adopted
|
||||||
|
// node never reads as converged (novox/hq ADR 0100).
|
||||||
|
for _, h := range updated.Held {
|
||||||
|
report.Held = append(report.Held, link.Held{ID: h.ID, Module: h.Module, Kind: h.Kind,
|
||||||
|
Target: h.Target, Since: h.Since, Changed: h.Changed, Kept: h.Kept})
|
||||||
|
}
|
||||||
|
if declared.Adoption != nil {
|
||||||
|
if updated.Firewall != nil {
|
||||||
|
report.Firewall = updated.Firewall.Kind
|
||||||
|
}
|
||||||
|
reached, err := reachable.Collect(ctx, apply.ExecRunner)
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "mesh-host: applied, and could not read what is reachable here: %v\n", err)
|
||||||
|
}
|
||||||
|
report.Reachable = reached
|
||||||
|
}
|
||||||
for _, change := range outcome.Outcomes {
|
for _, change := range outcome.Outcomes {
|
||||||
// What is held is not what this machine owns: it was found, and is kept as it was until
|
// What is held is not what this machine owns: it was found, and is kept as it was until
|
||||||
// its module is taken (novox/hq ADR 0100).
|
// its module is taken (novox/hq ADR 0100).
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/novox/mesh-host/internal/link"
|
||||||
"github.com/novox/mesh-host/internal/store"
|
"github.com/novox/mesh-host/internal/store"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -135,3 +137,35 @@ func TestAFlagAfterAPositionalIsRead(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Defends novox/hq ADR 0100: a reconcile on an adopted node speaks unasked only when what it holds
|
||||||
|
// or its firewall changed — which is how a predecessor still writing is caught, without a report
|
||||||
|
// every five minutes saying nothing new.
|
||||||
|
func TestAReconcileSpeaksOnlyWhenWhatIsHeldChanged(t *testing.T) {
|
||||||
|
w := &adoptionWatch{}
|
||||||
|
held := link.Report{Firewall: "ufw", Held: []link.Held{{ID: "hello-web.page"}, {ID: "hello-web.server"}}}
|
||||||
|
if !w.changed(held) {
|
||||||
|
t.Fatal("the first report of a hold was not said")
|
||||||
|
}
|
||||||
|
again := link.Report{Firewall: "ufw", Held: []link.Held{{ID: "hello-web.server"}, {ID: "hello-web.page"}}}
|
||||||
|
if w.changed(again) {
|
||||||
|
t.Error("the same holds in another order were said again")
|
||||||
|
}
|
||||||
|
rewritten := link.Report{Firewall: "ufw", Held: []link.Held{{ID: "hello-web.page", Changed: "rewritten"}, {ID: "hello-web.server"}}}
|
||||||
|
if !w.changed(rewritten) {
|
||||||
|
t.Error("a held file rewritten by something else was not said")
|
||||||
|
}
|
||||||
|
if !w.changed(link.Report{Firewall: "none", Held: rewritten.Held}) {
|
||||||
|
t.Error("a changed firewall was not said")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestWhatTheLinkPublishedCountsAsSaid(t *testing.T) {
|
||||||
|
w := &adoptionWatch{}
|
||||||
|
report := link.Report{Firewall: "ufw", Held: []link.Held{{ID: "a"}}}
|
||||||
|
applier := w.noting(func(context.Context, []byte, []byte) link.Report { return report })
|
||||||
|
applier(context.Background(), nil, nil)
|
||||||
|
if w.changed(report) {
|
||||||
|
t.Error("a reconcile repeated what the link had just published")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -220,6 +220,19 @@ func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) {
|
|||||||
if len(fields) != 5 {
|
if len(fields) != 5 {
|
||||||
t.Errorf("the token has %d fields, expected 5: %v", len(fields), fields)
|
t.Errorf("the token has %d fields, expected 5: %v", len(fields), fields)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// novox/hq ADR 0100: an adopted node's token says so, and a converged one's is unchanged.
|
||||||
|
raw, err = json.Marshal(Token{Version: 1, Secret: "s", Adopted: true})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
fields = map[string]any{}
|
||||||
|
if err := json.Unmarshal(raw, &fields); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if fields["adopted"] != true {
|
||||||
|
t.Errorf("an adopted token does not say \"adopted\": %v", fields)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestACompleteTokenParses(t *testing.T) {
|
func TestACompleteTokenParses(t *testing.T) {
|
||||||
|
|||||||
@@ -30,6 +30,11 @@ type Token struct {
|
|||||||
Fingerprint string `json:"fingerprint,omitempty"`
|
Fingerprint string `json:"fingerprint,omitempty"`
|
||||||
Signer []byte `json:"signer,omitempty"`
|
Signer []byte `json:"signer,omitempty"`
|
||||||
Secret string `json:"secret"`
|
Secret string `json:"secret"`
|
||||||
|
|
||||||
|
// Adopted says this node joins adopted (novox/hq ADR 0100). The host checks it speaks the
|
||||||
|
// firewall found here before enrolling, because an adopted node keeps that firewall in force.
|
||||||
|
// Absent for a converged node.
|
||||||
|
Adopted bool `json:"adopted,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// ParseToken reads a token a person pasted.
|
// ParseToken reads a token a person pasted.
|
||||||
|
|||||||
@@ -1,5 +1,7 @@
|
|||||||
package link
|
package link
|
||||||
|
|
||||||
|
import "time"
|
||||||
|
|
||||||
// The wire formats shared with the control plane, which defines them separately because this
|
// The wire formats shared with the control plane, which defines them separately because this
|
||||||
// binary requires nothing present and does not import it. A test on each side asserts the field
|
// binary requires nothing present and does not import it. A test on each side asserts the field
|
||||||
// names, so a rename breaks both at once rather than on a real machine months later.
|
// names, so a rename breaks both at once rather than on a real machine months later.
|
||||||
@@ -85,4 +87,44 @@ type Report struct {
|
|||||||
// than the send, and the machine reads as caught up with words it has not read yet. Clocks
|
// than the send, and the machine reads as caught up with words it has not read yet. Clocks
|
||||||
// cannot answer "which"; the digest is the answer itself.
|
// cannot answer "which"; the digest is the answer itself.
|
||||||
Declared string `json:"declared,omitempty"`
|
Declared string `json:"declared,omitempty"`
|
||||||
|
|
||||||
|
// Held is what this adopted node found and is keeping as it was until its module is taken
|
||||||
|
// (novox/hq ADR 0100). Without it an adopted node reads as converged.
|
||||||
|
Held []Held `json:"held,omitempty"`
|
||||||
|
|
||||||
|
// Firewall is the firewall found on this machine — "ufw" or "none" — and empty on a node that
|
||||||
|
// was never asked, which is every converged one.
|
||||||
|
Firewall string `json:"firewall,omitempty"`
|
||||||
|
|
||||||
|
// Reachable is what can be reached on this machine now: every listening socket and every
|
||||||
|
// published container port. Only an adopted node reports it; it is what converging the node
|
||||||
|
// previews, so nothing closes without being named first.
|
||||||
|
Reachable []Reach `json:"reachable,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Held is one file or container found on an adopted node and kept as it was.
|
||||||
|
type Held struct {
|
||||||
|
ID string `json:"id"`
|
||||||
|
Module string `json:"module"`
|
||||||
|
Kind string `json:"kind"`
|
||||||
|
Target string `json:"target"`
|
||||||
|
Since time.Time `json:"since"`
|
||||||
|
// Changed is what something other than the mesh did to it since — rewritten, stopped,
|
||||||
|
// replaced or gone — and empty while it is as found.
|
||||||
|
Changed string `json:"changed,omitempty"`
|
||||||
|
// Kept is where a file's original was kept.
|
||||||
|
Kept string `json:"kept,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Reach is one thing reachable on the machine: a listening socket, or a published container port.
|
||||||
|
type Reach struct {
|
||||||
|
Protocol string `json:"protocol"`
|
||||||
|
Address string `json:"address"`
|
||||||
|
Port int `json:"port"`
|
||||||
|
// By is what holds it — a process, or a container's name.
|
||||||
|
By string `json:"by,omitempty"`
|
||||||
|
// Published is a container port the runtime publishes, reached on the forwarded path; its
|
||||||
|
// container's own port is ContainerPort.
|
||||||
|
Published bool `json:"published,omitempty"`
|
||||||
|
ContainerPort int `json:"container-port,omitempty"`
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -120,6 +120,14 @@ func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) {
|
|||||||
{Signed{Declaration: []byte("{}"), Signature: []byte("x")}, []string{"declaration", "signature"}},
|
{Signed{Declaration: []byte("{}"), Signature: []byte("x")}, []string{"declaration", "signature"}},
|
||||||
{Report{Node: "n", Applied: []string{"a"}, Failed: map[string]string{"k": "v"}, Refused: "r"},
|
{Report{Node: "n", Applied: []string{"a"}, Failed: map[string]string{"k": "v"}, Refused: "r"},
|
||||||
[]string{"node", "applied", "failed", "refused"}},
|
[]string{"node", "applied", "failed", "refused"}},
|
||||||
|
// novox/hq ADR 0100: what an adopted node holds, the firewall it was found with, and what
|
||||||
|
// is reachable on it.
|
||||||
|
{Report{Node: "n", Held: []Held{{ID: "i"}}, Firewall: "ufw", Reachable: []Reach{{Port: 1}}},
|
||||||
|
[]string{"node", "held", "firewall", "reachable"}},
|
||||||
|
{Held{ID: "i", Module: "m", Kind: "file", Target: "/t", Changed: "rewritten", Kept: "/k"},
|
||||||
|
[]string{"id", "module", "kind", "target", "since", "changed", "kept"}},
|
||||||
|
{Reach{Protocol: "tcp", Address: "0.0.0.0", Port: 8080, By: "c", Published: true, ContainerPort: 80},
|
||||||
|
[]string{"protocol", "address", "port", "by", "published", "container-port"}},
|
||||||
} {
|
} {
|
||||||
raw, err := json.Marshal(c.value)
|
raw, err := json.Marshal(c.value)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
+16
-5
@@ -70,15 +70,21 @@ type Announce func(string)
|
|||||||
type Roused <-chan struct{}
|
type Roused <-chan struct{}
|
||||||
|
|
||||||
func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error {
|
func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error {
|
||||||
return HoldRoused(ctx, m, apply, say, timeout, nil)
|
return HoldRoused(ctx, m, apply, say, timeout, nil, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
// HoldRoused is Hold, told when the machine has reason to think its link is stale.
|
// Outbox carries reports the node has to say without having been sent anything — what a
|
||||||
|
// reconcile found changed on an adopted node (novox/hq ADR 0100). Published while the link is up;
|
||||||
|
// a report made while it is down waits in the channel for the next one. Nil is allowed.
|
||||||
|
type Outbox <-chan Report
|
||||||
|
|
||||||
|
// HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed
|
||||||
|
// reports to publish between deliveries.
|
||||||
func HoldRoused(ctx context.Context, m Membership, apply Applier, say Announce,
|
func HoldRoused(ctx context.Context, m Membership, apply Applier, say Announce,
|
||||||
timeout time.Duration, roused Roused) error {
|
timeout time.Duration, roused Roused, outbox Outbox) error {
|
||||||
|
|
||||||
return holdWith(ctx, func(ctx context.Context) error {
|
return holdWith(ctx, func(ctx context.Context) error {
|
||||||
return Run(ctx, m, apply, say, timeout)
|
return Run(ctx, m, apply, say, timeout, outbox)
|
||||||
}, say, roused)
|
}, say, roused)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -171,7 +177,8 @@ func holdWith(ctx context.Context, run attempt, say Announce, roused Roused) err
|
|||||||
//
|
//
|
||||||
// Outbound only, and nothing listens on this machine. Returns when the link ends, for any reason;
|
// Outbound only, and nothing listens on this machine. Returns when the link ends, for any reason;
|
||||||
// Hold is what decides whether to open it again.
|
// Hold is what decides whether to open it again.
|
||||||
func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error {
|
func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration,
|
||||||
|
outbox Outbox) error {
|
||||||
if say == nil {
|
if say == nil {
|
||||||
say = func(string) {}
|
say = func(string) {}
|
||||||
}
|
}
|
||||||
@@ -254,6 +261,10 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
|||||||
return nil
|
return nil
|
||||||
case <-beat.C:
|
case <-beat.C:
|
||||||
publishAlive(ctx, channel, m, say, timeout)
|
publishAlive(ctx, channel, m, say, timeout)
|
||||||
|
case report := <-outbox:
|
||||||
|
// Said without having been asked: a reconcile found what an adopted node holds, or
|
||||||
|
// its firewall, changed since it last said.
|
||||||
|
publishReport(ctx, channel, m, report, say, timeout)
|
||||||
case reason := <-closed:
|
case reason := <-closed:
|
||||||
return fmt.Errorf("the link closed: %v", reason)
|
return fmt.Errorf("the link closed: %v", reason)
|
||||||
case delivery, ok := <-deliveries:
|
case delivery, ok := <-deliveries:
|
||||||
|
|||||||
@@ -0,0 +1,181 @@
|
|||||||
|
// Package reachable reads what can be reached on this machine now: every listening socket, and
|
||||||
|
// every container port the runtime publishes (novox/hq ADR 0100).
|
||||||
|
//
|
||||||
|
// It is what converging an adopted node previews — each port, whether a module declares it or it
|
||||||
|
// will close — and what a converged genesis counts before refusing a machine in use. It reads; it
|
||||||
|
// never decides what is the mesh's.
|
||||||
|
package reachable
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"regexp"
|
||||||
|
"sort"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-host/internal/link"
|
||||||
|
"github.com/novox/mesh-host/internal/system"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Runner executes a command.
|
||||||
|
type Runner = system.Runner
|
||||||
|
|
||||||
|
// Reach is one thing reachable on this machine, in the words the report carries.
|
||||||
|
type Reach = link.Reach
|
||||||
|
|
||||||
|
// Collect reads the machine's listening sockets and the runtime's published ports. A published
|
||||||
|
// port is reported once, as published, rather than again as the runtime's proxy listening for it.
|
||||||
|
func Collect(ctx context.Context, run Runner) ([]Reach, error) {
|
||||||
|
out, err := run(ctx, "ss", "-Hltunp")
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("reading this machine's listening sockets: %w", err)
|
||||||
|
}
|
||||||
|
sockets := Sockets(out)
|
||||||
|
|
||||||
|
var published []Reach
|
||||||
|
if ps, err := run(ctx, "docker", "ps", "--format", "{{.Names}}\t{{.Ports}}"); err == nil {
|
||||||
|
published = Published(ps)
|
||||||
|
}
|
||||||
|
return Merge(sockets, published), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
var process = regexp.MustCompile(`users:\(\("([^"]+)"`)
|
||||||
|
|
||||||
|
// Sockets parses `ss -Hltunp`: each line a netid, a state, two queues, the local address and
|
||||||
|
// port, the peer, and the process when ss may name it.
|
||||||
|
func Sockets(out string) []Reach {
|
||||||
|
var reached []Reach
|
||||||
|
for _, line := range strings.Split(out, "\n") {
|
||||||
|
fields := strings.Fields(line)
|
||||||
|
if len(fields) < 5 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
protocol := fields[0]
|
||||||
|
if protocol != "tcp" && protocol != "udp" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
address, port, ok := splitLocal(fields[4])
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
r := Reach{Protocol: protocol, Address: address, Port: port}
|
||||||
|
if m := process.FindStringSubmatch(line); m != nil {
|
||||||
|
r.By = m[1]
|
||||||
|
}
|
||||||
|
reached = append(reached, r)
|
||||||
|
}
|
||||||
|
return reached
|
||||||
|
}
|
||||||
|
|
||||||
|
// splitLocal reads "127.0.0.1:53", "[::]:22", "*:22" and "[fe80::1]%veth0:123".
|
||||||
|
func splitLocal(local string) (string, int, bool) {
|
||||||
|
i := strings.LastIndex(local, ":")
|
||||||
|
if i < 0 {
|
||||||
|
return "", 0, false
|
||||||
|
}
|
||||||
|
port, err := strconv.Atoi(local[i+1:])
|
||||||
|
if err != nil {
|
||||||
|
return "", 0, false
|
||||||
|
}
|
||||||
|
address := local[:i]
|
||||||
|
if at := strings.Index(address, "%"); at >= 0 {
|
||||||
|
address = address[:at]
|
||||||
|
}
|
||||||
|
address = strings.TrimSuffix(strings.TrimPrefix(address, "["), "]")
|
||||||
|
if address == "*" {
|
||||||
|
address = "0.0.0.0"
|
||||||
|
}
|
||||||
|
return address, port, true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Published parses `docker ps --format '{{.Names}}\t{{.Ports}}'`. Only what is published on the
|
||||||
|
// machine counts; a port a container exposes and nothing publishes is not reachable from outside it.
|
||||||
|
func Published(out string) []Reach {
|
||||||
|
var reached []Reach
|
||||||
|
for _, line := range strings.Split(out, "\n") {
|
||||||
|
name, ports, ok := strings.Cut(strings.TrimSpace(line), "\t")
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, mapping := range strings.Split(ports, ",") {
|
||||||
|
reached = append(reached, mappingOf(name, strings.TrimSpace(mapping))...)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return reached
|
||||||
|
}
|
||||||
|
|
||||||
|
// mappingOf reads "0.0.0.0:9000-9001->9000-9001/tcp" into one reach per port.
|
||||||
|
func mappingOf(name, mapping string) []Reach {
|
||||||
|
outer, inner, ok := strings.Cut(mapping, "->")
|
||||||
|
if !ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
inner, protocol, ok := strings.Cut(inner, "/")
|
||||||
|
if !ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
i := strings.LastIndex(outer, ":")
|
||||||
|
if i < 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
address := strings.TrimSuffix(strings.TrimPrefix(outer[:i], "["), "]")
|
||||||
|
from, to, ok := portRange(outer[i+1:])
|
||||||
|
if !ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
cfrom, _, ok := portRange(inner)
|
||||||
|
if !ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
var reached []Reach
|
||||||
|
for p := from; p <= to; p++ {
|
||||||
|
reached = append(reached, Reach{Protocol: protocol, Address: address, Port: p, By: name,
|
||||||
|
Published: true, ContainerPort: cfrom + (p - from)})
|
||||||
|
}
|
||||||
|
return reached
|
||||||
|
}
|
||||||
|
|
||||||
|
func portRange(s string) (int, int, bool) {
|
||||||
|
a, b, isRange := strings.Cut(s, "-")
|
||||||
|
from, err := strconv.Atoi(a)
|
||||||
|
if err != nil {
|
||||||
|
return 0, 0, false
|
||||||
|
}
|
||||||
|
if !isRange {
|
||||||
|
return from, from, true
|
||||||
|
}
|
||||||
|
to, err := strconv.Atoi(b)
|
||||||
|
if err != nil || to < from {
|
||||||
|
return 0, 0, false
|
||||||
|
}
|
||||||
|
return from, to, true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Merge puts the published ports beside the sockets, dropping the runtime proxy's own socket for a
|
||||||
|
// port that is reported as published already, and sorts the whole by port.
|
||||||
|
func Merge(sockets, published []Reach) []Reach {
|
||||||
|
key := func(r Reach) string { return r.Protocol + " " + r.Address + " " + strconv.Itoa(r.Port) }
|
||||||
|
isPublished := map[string]bool{}
|
||||||
|
for _, p := range published {
|
||||||
|
isPublished[key(p)] = true
|
||||||
|
}
|
||||||
|
var out []Reach
|
||||||
|
for _, s := range sockets {
|
||||||
|
if s.By == "docker-proxy" && isPublished[key(s)] {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
out = append(out, s)
|
||||||
|
}
|
||||||
|
out = append(out, published...)
|
||||||
|
sort.SliceStable(out, func(i, j int) bool {
|
||||||
|
if out[i].Port != out[j].Port {
|
||||||
|
return out[i].Port < out[j].Port
|
||||||
|
}
|
||||||
|
if out[i].Protocol != out[j].Protocol {
|
||||||
|
return out[i].Protocol < out[j].Protocol
|
||||||
|
}
|
||||||
|
return out[i].Address < out[j].Address
|
||||||
|
})
|
||||||
|
return out
|
||||||
|
}
|
||||||
@@ -0,0 +1,100 @@
|
|||||||
|
package reachable
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"os"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Defends novox/hq ADR 0100: converging previews every listening socket and every published
|
||||||
|
// container port. Fixtures are captured from a real machine.
|
||||||
|
|
||||||
|
func fixture(t *testing.T, name string) string {
|
||||||
|
t.Helper()
|
||||||
|
raw, err := os.ReadFile("testdata/" + name)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
return string(raw)
|
||||||
|
}
|
||||||
|
|
||||||
|
func find(rs []Reach, protocol, address string, port int) (Reach, bool) {
|
||||||
|
for _, r := range rs {
|
||||||
|
if r.Protocol == protocol && r.Address == address && r.Port == port {
|
||||||
|
return r, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return Reach{}, false
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSocketsAreReadWithWhatHoldsThem(t *testing.T) {
|
||||||
|
got := Sockets(fixture(t, "ss.txt"))
|
||||||
|
if r, ok := find(got, "tcp", "0.0.0.0", 22); !ok || r.By != "sshd" {
|
||||||
|
t.Errorf("ssh not read: %+v", r)
|
||||||
|
}
|
||||||
|
if r, ok := find(got, "tcp", "::", 445); !ok || r.By != "smbd" {
|
||||||
|
t.Errorf("an IPv6 wildcard listener not read: %+v", r)
|
||||||
|
}
|
||||||
|
if _, ok := find(got, "udp", "fe80::849e:ccff:fea8:24c7", 123); !ok {
|
||||||
|
t.Error("a link-local address with a scope was not read")
|
||||||
|
}
|
||||||
|
if r, ok := find(got, "udp", "127.0.0.1", 53); !ok || r.By != "dnsmasq" {
|
||||||
|
t.Errorf("a loopback udp socket not read: %+v", r)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPublishedPortsNameTheirContainerAndItsPort(t *testing.T) {
|
||||||
|
got := Published(fixture(t, "docker-ps.txt"))
|
||||||
|
if r, ok := find(got, "tcp", "0.0.0.0", 8770); !ok || r.By != "whisper" || r.ContainerPort != 8000 || !r.Published {
|
||||||
|
t.Errorf("a published port: %+v", r)
|
||||||
|
}
|
||||||
|
if r, ok := find(got, "tcp", "0.0.0.0", 9001); !ok || r.ContainerPort != 9001 {
|
||||||
|
t.Errorf("a published range was not expanded: %+v", r)
|
||||||
|
}
|
||||||
|
if r, ok := find(got, "tcp", "127.0.0.1", 15673); !ok || r.ContainerPort != 15672 {
|
||||||
|
t.Errorf("a loopback-published port: %+v", r)
|
||||||
|
}
|
||||||
|
for _, r := range got {
|
||||||
|
if r.By == "umami_db" {
|
||||||
|
t.Errorf("an exposed and unpublished port was reported reachable: %+v", r)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestAPublishedPortIsReportedOnceAsPublished(t *testing.T) {
|
||||||
|
merged := Merge(Sockets(fixture(t, "ss.txt")), Published(fixture(t, "docker-ps.txt")))
|
||||||
|
n := 0
|
||||||
|
for _, r := range merged {
|
||||||
|
if r.Protocol == "tcp" && r.Address == "0.0.0.0" && r.Port == 8770 {
|
||||||
|
n++
|
||||||
|
if !r.Published {
|
||||||
|
t.Errorf("the runtime's proxy was reported instead of the published port: %+v", r)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if n != 1 {
|
||||||
|
t.Errorf("port 8770 reported %d times", n)
|
||||||
|
}
|
||||||
|
if _, ok := find(merged, "tcp", "0.0.0.0", 22); !ok {
|
||||||
|
t.Error("a socket was lost in the merge")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCollectAsksSsAndTheRuntime(t *testing.T) {
|
||||||
|
var asked []string
|
||||||
|
run := func(_ context.Context, name string, args ...string) (string, error) {
|
||||||
|
asked = append(asked, name+" "+strings.Join(args, " "))
|
||||||
|
if name == "ss" {
|
||||||
|
return fixture(t, "ss.txt"), nil
|
||||||
|
}
|
||||||
|
return fixture(t, "docker-ps.txt"), nil
|
||||||
|
}
|
||||||
|
got, err := Collect(context.Background(), run)
|
||||||
|
if err != nil || len(got) == 0 {
|
||||||
|
t.Fatalf("%v %v", got, err)
|
||||||
|
}
|
||||||
|
if len(asked) != 2 {
|
||||||
|
t.Errorf("asked %v", asked)
|
||||||
|
}
|
||||||
|
}
|
||||||
+7
@@ -0,0 +1,7 @@
|
|||||||
|
mesh-controller-check-adoption 127.0.0.1:55541->5432/tcp
|
||||||
|
umami_db 5432/tcp
|
||||||
|
whisper 0.0.0.0:8770->8000/tcp, [::]:8770->8000/tcp
|
||||||
|
keycloak 8443/tcp, 127.0.0.1:28080->8080/tcp
|
||||||
|
minio-lb 0.0.0.0:9000-9001->9000-9001/tcp, [::]:9000-9001->9000-9001/tcp
|
||||||
|
wonderful_mahavira
|
||||||
|
anton-lavinmq 0.0.0.0:5680->5672/tcp, [::]:5680->5672/tcp, 127.0.0.1:15673->15672/tcp
|
||||||
Vendored
+23
@@ -0,0 +1,23 @@
|
|||||||
|
udp UNCONN 0 0 0.0.0.0:55558 0.0.0.0:* users:(("firefox",pid=2283907,fd=288))
|
||||||
|
udp UNCONN 0 0 0.0.0.0:59541 0.0.0.0:* users:(("firefox",pid=2283907,fd=241))
|
||||||
|
udp UNCONN 0 0 127.0.0.1:53 0.0.0.0:* users:(("dnsmasq",pid=1189392,fd=6))
|
||||||
|
udp UNCONN 0 0 0.0.0.0:33525 0.0.0.0:* users:(("firefox",pid=2283907,fd=304))
|
||||||
|
udp UNCONN 0 0 0.0.0.0:41749 0.0.0.0:* users:(("firefox",pid=2283907,fd=351))
|
||||||
|
tcp LISTEN 0 4096 127.0.0.1:55541 0.0.0.0:* users:(("docker-proxy",pid=4108732,fd=7))
|
||||||
|
tcp LISTEN 0 4096 0.0.0.0:9001 0.0.0.0:* users:(("docker-proxy",pid=1849130,fd=7))
|
||||||
|
tcp LISTEN 0 4096 0.0.0.0:8770 0.0.0.0:* users:(("docker-proxy",pid=1920035,fd=7))
|
||||||
|
tcp LISTEN 0 50 0.0.0.0:445 0.0.0.0:* users:(("smbd",pid=1248,fd=29))
|
||||||
|
tcp LISTEN 0 128 0.0.0.0:22 0.0.0.0:* users:(("sshd",pid=1188536,fd=6))
|
||||||
|
tcp LISTEN 0 50 0.0.0.0:139 0.0.0.0:* users:(("smbd",pid=1248,fd=30))
|
||||||
|
tcp LISTEN 0 4096 127.0.0.1:5432 0.0.0.0:* users:(("docker-proxy",pid=1854543,fd=7))
|
||||||
|
tcp LISTEN 0 32 127.0.0.1:53 0.0.0.0:* users:(("dnsmasq",pid=1189392,fd=7))
|
||||||
|
tcp LISTEN 0 4096 127.0.0.1:15673 0.0.0.0:* users:(("docker-proxy",pid=3170,fd=7))
|
||||||
|
tcp LISTEN 0 4096 [::]:9001 [::]:* users:(("docker-proxy",pid=1849138,fd=7))
|
||||||
|
tcp LISTEN 0 4096 [::]:8770 [::]:* users:(("docker-proxy",pid=1920043,fd=7))
|
||||||
|
tcp LISTEN 0 50 [::]:445 [::]:* users:(("smbd",pid=1248,fd=27))
|
||||||
|
tcp LISTEN 0 128 [::]:22 [::]:* users:(("sshd",pid=1188536,fd=7))
|
||||||
|
tcp LISTEN 0 50 [::]:139 [::]:* users:(("smbd",pid=1248,fd=28))
|
||||||
|
udp UNCONN 0 0 [fd42:f8c5:dae:d74c::1]:53 [::]:*
|
||||||
|
udp UNCONN 0 0 [fe80::849e:ccff:fea8:24c7]%veth6b2b7ba:123 [::]:*
|
||||||
|
udp UNCONN 0 0 [fe80::e45a:90ff:feca:148f]%vethb5e5a61:123 [::]:*
|
||||||
|
udp UNCONN 0 0 [fe80::c4ed:ccff:feb1:afd2]%veth005a182:123 [::]:*
|
||||||
Reference in New Issue
Block a user