package main import ( "encoding/json" "fmt" "log" "os" "strings" "sync/atomic" "time" "github.com/nats-io/nats.go" "github.com/novox/mesh-controller/internal/broker" ) // What the mesh issued this proxy, read on the bus (novox/hq ADR 0160, ADR 0167). // // **The proxy is told, not left to work it out.** Its membership carries the routes it is given — // the same contributions its file is written from — and every machine's address on the private // network, which is who may be served an internal name. Read once at connect and followed, so a // route added or a machine joining reaches a running proxy without a restart. // credential is the bus account the mesh delivered as this module's own secret named broker. type credential struct { URL string `json:"url"` Fingerprint string `json:"fingerprint"` Node string `json:"node"` Module string `json:"module"` User string `json:"user"` Password string `json:"password"` } // followMembership connects with the credential in path and applies every membership the mesh // issues this proxy. It retries the first connection for as long as it takes: a proxy that started // before the bus keeps serving the file, and takes the bus when it answers. func followMembership(path string, held *table, fromBus *atomic.Bool) { for { err := followOnce(path, held, fromBus) if err == nil { return } log.Printf("cannot follow this proxy's membership, serving the file meanwhile: %v", err) time.Sleep(30 * time.Second) } } func followOnce(path string, held *table, fromBus *atomic.Bool) error { raw, err := os.ReadFile(path) if err != nil { return err } var cred credential if err := json.Unmarshal(raw, &cred); err != nil { return fmt.Errorf("the broker credential is not one: %w", err) } if cred.Node == "" || cred.Module == "" { return fmt.Errorf("the broker credential names no node or module, so it has no membership") } opts := []nats.Option{ nats.Name(cred.Node + "." + cred.Module), nats.UserInfo(cred.User, cred.Password), // Its own inbox, and nothing wider: every principal is granted `_INBOX..>` alone. nats.CustomInboxPrefix("_INBOX." + cred.User), // The bus being restarted is an upgrade, not a reason to stop following. nats.MaxReconnects(-1), } if strings.TrimSpace(cred.Fingerprint) != "" { opts = append(opts, nats.Secure(broker.PinnedToFingerprint(cred.Fingerprint))) } conn, err := nats.Connect(cred.URL, opts...) if err != nil { return fmt.Errorf("connecting to the bus at %s: %w", broker.BareAddress(cred.URL), err) } subject := broker.MembershipSubject(cred.Node, cred.Module) apply := func(body []byte) { var issued broker.Membership if err := json.Unmarshal(body, &issued); err != nil { log.Printf("a membership arrived that is not one: %v", err) return } if took := applyMembership(issued, held); took && !fromBus.Swap(true) { log.Printf("routes now come from this proxy's membership on %s", subject) } } // Followed first, read second: an issue landing between the two is applied, not missed. if _, err := conn.Subscribe(subject, func(m *nats.Msg) { apply(m.Data) }); err != nil { conn.Close() return fmt.Errorf("cannot follow %s: %w", subject, err) } // The subject-addressed direct get: the one request this account may make of the stream. got, err := conn.Request("$JS.API.DIRECT.GET."+broker.AssignmentsStream+"."+subject, nil, 5*time.Second) switch { case err != nil: log.Printf("cannot read the membership issued on %s yet (%v); following it", subject, err) case got.Header.Get("Status") != "" || len(got.Data) == 0: log.Printf("no membership issued on %s yet; serving the file until one is", subject) default: apply(got.Data) } return nil } // applyMembership serves what a membership says, and says whether it said anything about routes. // // A membership with no routes in it is one from a controller older than ADR 0167, and the file stays // the source rather than every route being withdrawn because a field was absent. func applyMembership(issued broker.Membership, held *table) bool { raw, carries := issued.Receives["route"] if !carries { return false } var contributions []contribution if err := json.Unmarshal(raw, &contributions); err != nil { log.Printf("the routes in this proxy's membership are not contributions, keeping what is served: %v", err) return false } inside, err := sourcesOf(issued.Mesh) if err != nil { log.Printf("the mesh in this proxy's membership is unreadable, keeping what is served: %v", err) return false } routes, public := routesOf(contributions) held.set(routes, public) held.setInside(inside) log.Printf("serving %d route(s) from the membership, internal names to %d machine(s): %s", len(routes), len(inside), strings.Join(held.names(), ", ")) return true }