From 81c0188e2083b099c8026339152f4af7e749ec14 Mon Sep 17 00:00:00 2001 From: Antonio Ojea Date: Sun, 20 Sep 2026 10:49:21 +0200 Subject: [PATCH 1/2] control-plane: publish mesh events for real NewServer installs NewNopMeshAdapter and nothing in the binary ever replaced it, so every MeshEvent the code carefully signs (BANNED, KEY_ROTATION, POLICY_UPDATE) was dropped on the floor. Routers and nodes handle those events; they just never received one outside the tests that wire the adapter themselves. The control plane now has a presence on the mesh, and it is as small as it can be: an ephemeral libp2p identity with no listen address, no DHT, no relay and no stream handlers, holding a gossipsub session with the routers so a publish has somewhere to go. It finds the routers the one way the control plane already knows the mesh's shape, the lease table, and re-dials any it lost every --mesh-reconnect-interval. It is one-way: nothing about the mesh is read back, and every consumer keeps pulling /keys, /info and /policies, so a missed event is a delay and never a divergence. sam-one runs the router in the same process, so its control plane publishes on the embedded router's own topic instead of dialing itself. NewP2PMeshAdapter now takes that topic rather than joining one, and Close only tears down what NewMeshPublisher built. Fixes #317 --- cmd/sam-control-plane/main.go | 15 ++ internal/controlplane/mesh.go | 170 +++++++++++++++--- internal/controlplane/mesh_test.go | 155 +++++++++++++++- .../peer_id_canonicalization_test.go | 6 +- internal/standalone/standalone.go | 9 +- site/content/docs/reference/control-plane.md | 3 +- 6 files changed, 328 insertions(+), 30 deletions(-) diff --git a/cmd/sam-control-plane/main.go b/cmd/sam-control-plane/main.go index f9c7ae7b..ddf353af 100644 --- a/cmd/sam-control-plane/main.go +++ b/cmd/sam-control-plane/main.go @@ -46,6 +46,7 @@ var ( biscuitTTL time.Duration oidcSessionTTL time.Duration nodeRetention time.Duration + meshReconnectInterval time.Duration adminTokenPath string insecureSkipTLSVerify bool logLevel string @@ -142,6 +143,19 @@ func main() { } }() + // Bans, key rotations and policy updates reach the mesh as they + // happen; every consumer also pulls, so this is speed, not truth. + mesh, err := controlplane.NewMeshPublisher(cmd.Context(), store, meshReconnectInterval) + if err != nil { + logger.Fatalf("Failed to start mesh event publisher: %v", err) + } + defer func() { + if err := mesh.Close(); err != nil { + logger.Errorf("Failed to stop mesh event publisher: %v", err) + } + }() + srv.SetMeshAdapter(mesh) + if err := srv.Start(); err != nil { logger.Fatalf("Failed to start control plane: %v", err) } @@ -164,6 +178,7 @@ func main() { rootCmd.Flags().DurationVar(&biscuitTTL, "biscuit-ttl", api.BiscuitTokenTTL, "Lifespan minted into every issued Biscuit's expiration fact. Capped to the OIDC token's own expiry when shorter.") rootCmd.Flags().DurationVar(&oidcSessionTTL, "oidc-session-ttl", api.OIDCSessionTTL, "How long an OIDC enrollment stays refreshable before the identity must re-authenticate with the OIDC provider. Shorter values keep the provider authoritative for offboarding at the cost of more frequent interactive re-enrollment.") rootCmd.Flags().DurationVar(&nodeRetention, "node-retention", controlplane.DefaultNodeRetention, "How long an enrolled node's record is kept after its session expired before it is deleted. Banned nodes are always kept. 0 keeps every record forever.") + rootCmd.Flags().DurationVar(&meshReconnectInterval, "mesh-reconnect-interval", controlplane.DefaultMeshReconnectInterval, "How often the event publisher re-reads the router leases and dials any router it is not connected to.") rootCmd.Flags().StringVar(&adminTokenPath, "admin-token-path", "", "Path to file containing the token for authenticating policy REST API requests (or env SAM_ADMIN_TOKEN)") rootCmd.Flags().BoolVar(&insecureSkipTLSVerify, "insecure-skip-tls-verify", false, "Skip TLS verification for OIDC providers") rootCmd.Flags().StringVar(&logLevel, "log-level", "info", "Log level (debug, info, warn, error)") diff --git a/internal/controlplane/mesh.go b/internal/controlplane/mesh.go index e43e7ca8..152e9f1f 100644 --- a/internal/controlplane/mesh.go +++ b/internal/controlplane/mesh.go @@ -22,10 +22,12 @@ import ( "sync" "time" + "github.com/libp2p/go-libp2p" pubsub "github.com/libp2p/go-libp2p-pubsub" "github.com/libp2p/go-libp2p/core/host" + "github.com/libp2p/go-libp2p/core/network" "github.com/libp2p/go-libp2p/core/peer" - "github.com/libp2p/go-libp2p/core/peerstore" + "github.com/multiformats/go-multiaddr" "google.golang.org/protobuf/proto" "github.com/google/sam/api" @@ -87,42 +89,161 @@ func (n *NopMeshAdapter) Close() error { return nil } -// P2PMeshAdapter implements MeshAdapter using a libp2p Host and GossipSub subscriber/publisher. +// P2PMeshAdapter publishes control plane events on the mesh's gossip topic. +// +// It is the control plane's whole presence on the mesh, and it is one-way: +// events go out so a ban, a key rotation or a policy change reaches routers +// and nodes the moment it happens, and nothing is read back. Every consumer +// still pulls /keys, /info and /policies on its own schedule, so a missed +// event is a delay, never a divergence. Where the adapter runs on a host of +// its own (NewMeshPublisher) that host has no listen address, no DHT, no +// relay and no stream handlers: it can dial routers and nothing can dial it. type P2PMeshAdapter struct { host host.Host - ps *pubsub.PubSub topic *pubsub.Topic store storage.Store mu sync.Mutex + // close tears down what NewMeshPublisher built; nil when the host and + // topic belong to someone else (sam-one's embedded router). + close func() error +} + +// NewP2PMeshAdapter publishes on an existing host's topic. The caller owns +// both and closes them; Close on the adapter is a no-op. +func NewP2PMeshAdapter(h host.Host, topic *pubsub.Topic, store storage.Store) (*P2PMeshAdapter, error) { + if h == nil || topic == nil || store == nil { + return nil, fmt.Errorf("host, topic, and store cannot be nil") + } + if topic.String() != api.GossipEvents { + return nil, fmt.Errorf("topic %q is not the mesh events topic %q", topic.String(), api.GossipEvents) + } + return &P2PMeshAdapter{host: h, topic: topic, store: store}, nil } -func NewP2PMeshAdapter(h host.Host, ps *pubsub.PubSub, store storage.Store) (*P2PMeshAdapter, error) { - if h == nil || ps == nil || store == nil { - return nil, fmt.Errorf("host, pubsub, and store cannot be nil") +// RouterDialTimeout bounds each attempt to reach a leased router. +const RouterDialTimeout = 10 * time.Second + +// DefaultMeshReconnectInterval is how often the publisher re-reads the lease +// table; a router that just enrolled waits at most this long for events. It +// is one query and at most a few dials per tick, so it is kept short. +const DefaultMeshReconnectInterval = 30 * time.Second + +// NewMeshPublisher builds the control plane's own publish-only peer and +// keeps it connected to every router holding a lease, re-checking the lease +// table every reconnect. The lease table is all the control plane needs to +// know about the mesh's shape, and it already has it. Close stops the loop +// and the host. +func NewMeshPublisher(ctx context.Context, store storage.Store, reconnect time.Duration) (*P2PMeshAdapter, error) { + if store == nil { + return nil, fmt.Errorf("store cannot be nil") + } + if reconnect <= 0 { + return nil, fmt.Errorf("reconnect interval must be positive, got %s", reconnect) + } + h, err := libp2p.New(libp2p.NoListenAddrs, libp2p.DisableRelay()) + if err != nil { + return nil, fmt.Errorf("mesh publisher host: %w", err) + } + loopCtx, cancel := context.WithCancel(ctx) + // StrictSign is the default; pinned because routers and nodes key their + // per-author rate limit on the signed sender. + ps, err := pubsub.NewGossipSub(loopCtx, h, pubsub.WithMessageSignaturePolicy(pubsub.StrictSign)) + if err != nil { + cancel() + _ = h.Close() + return nil, fmt.Errorf("mesh publisher gossipsub: %w", err) } topic, err := ps.Join(api.GossipEvents) if err != nil { - return nil, fmt.Errorf("failed to join gossip events topic %s: %w", api.GossipEvents, err) + cancel() + _ = h.Close() + return nil, fmt.Errorf("join %s: %w", api.GossipEvents, err) + } + p := &P2PMeshAdapter{host: h, topic: topic, store: store} + var wg sync.WaitGroup + wg.Add(1) + go func() { + defer wg.Done() + p.keepRoutersConnected(loopCtx, reconnect) + }() + p.close = func() error { + cancel() + wg.Wait() + _ = topic.Close() + return h.Close() } + logger.Infof("[Mesh] Publishing control plane events as %s", h.ID()) + return p, nil +} - return &P2PMeshAdapter{ - host: h, - ps: ps, - topic: topic, - store: store, - }, nil +// keepRoutersConnected dials, now and every interval, each leased router the +// host is not connected to. A publish only reaches peers that have announced +// the topic, so connections are kept warm ahead of the events rather than +// made when one is due. +func (p *P2PMeshAdapter) keepRoutersConnected(ctx context.Context, interval time.Duration) { + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + p.connectRouters(ctx) + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } +} + +func (p *P2PMeshAdapter) connectRouters(ctx context.Context) { + routers, err := p.store.GetActiveRouters(ctx) + if err != nil { + logger.Warnf("[Mesh] Cannot list routers to publish to: %v", err) + return + } + for _, r := range routers { + info, err := routerAddrInfo(r) + if err != nil { + logger.Warnf("[Mesh] Skipping router lease %q: %v", r.PeerID, err) + continue + } + if p.host.Network().Connectedness(info.ID) == network.Connected { + continue + } + dialCtx, cancel := context.WithTimeout(ctx, RouterDialTimeout) + err = p.host.Connect(dialCtx, info) + cancel() + if err != nil { + logger.Warnf("[Mesh] Router %s unreachable for event publishing: %v", info.ID, err) + continue + } + logger.Infof("[Mesh] Connected to router %s", info.ID) + } } -func (p *P2PMeshAdapter) ConnectPeer(ctx context.Context, targetAddr string) error { - info, err := peer.AddrInfoFromString(targetAddr) +// routerAddrInfo is the dial target for a lease. Routers announce their +// addresses with a trailing /p2p/; the dialer wants the id once and the +// addresses bare, and an address naming a different peer is not this +// router's. +func routerAddrInfo(r storage.RouterLease) (peer.AddrInfo, error) { + pid, err := peer.Decode(r.PeerID) if err != nil { - return fmt.Errorf("invalid peer address string %q: %w", targetAddr, err) + return peer.AddrInfo{}, fmt.Errorf("invalid peer ID: %w", err) } - p.host.Peerstore().AddAddrs(info.ID, info.Addrs, peerstore.PermanentAddrTTL) - if err := p.host.Connect(ctx, *info); err != nil { - return fmt.Errorf("failed to connect to peer %s: %w", info.ID, err) + info := peer.AddrInfo{ID: pid} + for _, s := range r.Addresses { + ma, err := multiaddr.NewMultiaddr(s) + if err != nil { + continue + } + addr, id := peer.SplitAddr(ma) + if addr == nil || (id != "" && id != pid) { + continue + } + info.Addrs = append(info.Addrs, addr) } - return nil + if len(info.Addrs) == 0 { + return peer.AddrInfo{}, fmt.Errorf("no dialable address in %v", r.Addresses) + } + return info, nil } func (p *P2PMeshAdapter) PublishEvent(ctx context.Context, eventType api.MeshEvent_Type, peerID string, payload []byte) error { @@ -218,11 +339,8 @@ func (p *P2PMeshAdapter) GetNodeStatus(ctx context.Context, peerID string) (*Nod func (p *P2PMeshAdapter) Close() error { p.mu.Lock() defer p.mu.Unlock() - if p.topic != nil { - _ = p.topic.Close() + if p.close == nil { + return nil } - if p.host != nil { - return p.host.Close() - } - return nil + return p.close() } diff --git a/internal/controlplane/mesh_test.go b/internal/controlplane/mesh_test.go index d6439262..4284fcb7 100644 --- a/internal/controlplane/mesh_test.go +++ b/internal/controlplane/mesh_test.go @@ -23,6 +23,7 @@ import ( "time" pubsub "github.com/libp2p/go-libp2p-pubsub" + "github.com/libp2p/go-libp2p/core/network" "github.com/libp2p/go-libp2p/core/peer" "github.com/libp2p/go-libp2p" @@ -85,8 +86,12 @@ func TestP2PMeshAdapter_PublishAndSubscribe(t *testing.T) { if err != nil { t.Fatalf("failed to create cp pubsub: %v", err) } + cpTopic, err := cpPS.Join(api.GossipEvents) + if err != nil { + t.Fatalf("failed to join cp topic: %v", err) + } - adapter, err := NewP2PMeshAdapter(cpHost, cpPS, store) + adapter, err := NewP2PMeshAdapter(cpHost, cpTopic, store) if err != nil { t.Fatalf("failed to create P2PMeshAdapter: %v", err) } @@ -206,3 +211,151 @@ func TestKeyRotationEventIsSignedByTheRetiringKey(t *testing.T) { t.Error("a BANNED event after rotation must be signed by the current key") } } + +// The shipped control plane's mesh presence (#317): a publish-only peer that +// finds the routers through their leases, dials them, and gets an event to a +// subscriber on the other side. Nothing is listening on the publisher's +// side, so nothing can dial it. +func TestMeshPublisherReachesLeasedRouter(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + store, err := storage.NewSQLStore("sqlite", ":memory:") + if err != nil { + t.Fatalf("failed to create store: %v", err) + } + defer func() { _ = store.Close() }() + pub, priv, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatalf("failed to generate key: %v", err) + } + if err := store.SaveInitialKey(ctx, priv, pub); err != nil { + t.Fatalf("failed to save key: %v", err) + } + + // A router as the publisher sees it: a host subscribed to the events + // topic, known only through the lease it renews. + routerHost, err := libp2p.New(libp2p.ListenAddrStrings("/ip4/127.0.0.1/tcp/0")) + if err != nil { + t.Fatalf("failed to create router host: %v", err) + } + defer func() { _ = routerHost.Close() }() + routerPS, err := pubsub.NewGossipSub(ctx, routerHost, pubsub.WithMessageSignaturePolicy(pubsub.StrictSign)) + if err != nil { + t.Fatalf("failed to create router pubsub: %v", err) + } + routerTopic, err := routerPS.Join(api.GossipEvents) + if err != nil { + t.Fatalf("failed to join topic: %v", err) + } + sub, err := routerTopic.Subscribe() + if err != nil { + t.Fatalf("failed to subscribe: %v", err) + } + now := time.Now() + lease := storage.RouterLease{PeerID: routerHost.ID().String(), LastRenewal: now, ExpiresAt: now.Add(time.Hour)} + for _, a := range routerHost.Addrs() { + lease.Addresses = append(lease.Addresses, a.String()+"/p2p/"+routerHost.ID().String()) + } + // An expired lease must not be dialed: it is a router that is gone. + stale := storage.RouterLease{PeerID: routerHost.ID().String() + "x", Addresses: []string{"/ip4/127.0.0.1/tcp/1"}, LastRenewal: now.Add(-2 * time.Hour), ExpiresAt: now.Add(-time.Hour)} + for _, l := range []*storage.RouterLease{&lease, &stale} { + if err := store.UpsertRouterLease(ctx, l); err != nil { + t.Fatalf("lease %s: %v", l.PeerID, err) + } + } + + publisher, err := NewMeshPublisher(ctx, store, 50*time.Millisecond) + if err != nil { + t.Fatalf("NewMeshPublisher: %v", err) + } + defer func() { + if err := publisher.Close(); err != nil { + t.Errorf("Close: %v", err) + } + }() + if addrs := publisher.host.Network().ListenAddresses(); len(addrs) != 0 { + t.Errorf("publisher must not listen, has %v", addrs) + } + + // The router announces its subscription on connect; the publisher's + // first publish must have somewhere to go. The publisher itself never + // subscribes, so the router's peer list for the topic stays empty. + for routerHost.Network().Connectedness(publisher.host.ID()) != network.Connected { + select { + case <-ctx.Done(): + t.Fatal("publisher never connected to the leased router") + case <-time.After(20 * time.Millisecond): + } + } + for len(publisher.topic.ListPeers()) == 0 { + select { + case <-ctx.Done(): + t.Fatal("publisher never learned that the router is subscribed to the events topic") + case <-time.After(20 * time.Millisecond): + } + } + if peers := routerPS.ListPeers(api.GossipEvents); len(peers) != 0 { + t.Errorf("publisher must not subscribe to the topic, router sees %v", peers) + } + + _, target := newTestKey(t) + if err := publisher.PublishEvent(ctx, api.MeshEvent_BANNED, target.String(), nil); err != nil { + t.Fatalf("PublishEvent: %v", err) + } + msg, err := sub.Next(ctx) + if err != nil { + t.Fatalf("router did not receive the event: %v", err) + } + var event api.MeshEvent + if err := proto.Unmarshal(msg.Data, &event); err != nil { + t.Fatalf("unmarshal event: %v", err) + } + if event.Type != api.MeshEvent_BANNED || event.PeerId != target.String() { + t.Fatalf("got event %v for %q, want BANNED for %q", event.Type, event.PeerId, target) + } +} + +// Lease addresses carry a trailing /p2p/; the dialer gets the id once +// and the addresses bare, and an address naming another peer is dropped. +func TestRouterAddrInfo(t *testing.T) { + _, pid := newTestKey(t) + _, other := newTestKey(t) + + info, err := routerAddrInfo(storage.RouterLease{ + PeerID: pid.String(), + Addresses: []string{ + "/ip4/10.0.0.1/tcp/4501/p2p/" + pid.String(), + "/dnsaddr/bootstrap.example/p2p/" + pid.String(), + "/ip4/10.0.0.2/tcp/4501/p2p/" + other.String(), + "/ip4/10.0.0.3/tcp/4501", + "not a multiaddr", + }, + }) + if err != nil { + t.Fatalf("routerAddrInfo: %v", err) + } + if info.ID != pid { + t.Errorf("peer = %s, want %s", info.ID, pid) + } + var got []string + for _, a := range info.Addrs { + got = append(got, a.String()) + } + want := []string{"/ip4/10.0.0.1/tcp/4501", "/dnsaddr/bootstrap.example", "/ip4/10.0.0.3/tcp/4501"} + if len(got) != len(want) { + t.Fatalf("addrs = %v, want %v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("addrs[%d] = %s, want %s", i, got[i], want[i]) + } + } + + if _, err := routerAddrInfo(storage.RouterLease{PeerID: "not-a-peer", Addresses: []string{"/ip4/10.0.0.1/tcp/4501"}}); err == nil { + t.Error("an undecodable peer ID must be rejected") + } + if _, err := routerAddrInfo(storage.RouterLease{PeerID: pid.String(), Addresses: []string{"/ip4/10.0.0.2/tcp/4501/p2p/" + other.String()}}); err == nil { + t.Error("a lease with no address for its own peer must be rejected") + } +} diff --git a/internal/controlplane/peer_id_canonicalization_test.go b/internal/controlplane/peer_id_canonicalization_test.go index 910e0554..8575fe9d 100644 --- a/internal/controlplane/peer_id_canonicalization_test.go +++ b/internal/controlplane/peer_id_canonicalization_test.go @@ -407,8 +407,12 @@ func TestPublishEventValidatesPeerID(t *testing.T) { if err != nil { t.Fatalf("failed to create pubsub: %v", err) } + topic, err := ps.Join(api.GossipEvents) + if err != nil { + t.Fatalf("failed to join topic: %v", err) + } - adapter, err := NewP2PMeshAdapter(h, ps, store) + adapter, err := NewP2PMeshAdapter(h, topic, store) if err != nil { t.Fatalf("failed to create P2PMeshAdapter: %v", err) } diff --git a/internal/standalone/standalone.go b/internal/standalone/standalone.go index aeb51e04..9bcc8ee6 100644 --- a/internal/standalone/standalone.go +++ b/internal/standalone/standalone.go @@ -356,7 +356,14 @@ func (s *Server) Start(ctx context.Context) error { return fmt.Errorf("failed to start router: %w", err) } s.router = rtr - + // One process, one mesh peer: the control plane publishes its events on + // the embedded router's own topic instead of dialing itself. + mesh, err := controlplane.NewP2PMeshAdapter(rtr.Host, rtr.EventTopic, store) + if err != nil { + _ = rtr.Close() + return fmt.Errorf("failed to attach control plane to the mesh: %w", err) + } + cp.SetMeshAdapter(mesh) if s.publicAddr, err = s.resolvePublicAddr(); err != nil { _ = rtr.Close() return err diff --git a/site/content/docs/reference/control-plane.md b/site/content/docs/reference/control-plane.md index 65f1746e..074ee299 100644 --- a/site/content/docs/reference/control-plane.md +++ b/site/content/docs/reference/control-plane.md @@ -32,9 +32,10 @@ sam-control-plane admin unban --peer lift a ban | `--biscuit-ttl` | `24h` | Lifetime of each credential. If the OIDC token expires sooner, the credential expires with it. | | `--oidc-session-ttl` | `2160h` (90 days) | How long an OIDC enrollment may keep refreshing before the identity must log in again. | | `--key-rotation-interval` | `24h` | How often a new signing key is generated. `0` disables rotation. | -| `--key-grace-period` | `1h` | How long a rotated-out key stays accepted. Credentials signed by a retired key cannot be verified or refreshed. | +| `--key-grace-period` | `1h` | How long a rotated-out key stays accepted. Credentials signed by a retired key cannot be verified or refreshed. Nodes and routers must pull `/keys` well within this window (`sam-node --control-plane-sync-interval`, `sam-router --keys-sync-interval`). | | `--lease-duration` | `15m` | How long a router lease lasts without renewal. | | `--node-retention` | `720h` (30 days) | How long the record of an enrolled node is kept after its session expires. Banned nodes are kept forever. `0` keeps every record. | +| `--mesh-reconnect-interval` | `30s` | How often the event publisher re-reads the router leases and dials any router it is not connected to. | | `--log-level` | `info` | `debug`, `info`, `warn`, `error`. `LOG_FORMAT=json` selects JSON output. | ## HTTP API From c2ea967a682f8e2c3205647cda8a6c89d581a1d5 Mon Sep 17 00:00:00 2001 From: Antonio Ojea Date: Sun, 20 Sep 2026 11:23:07 +0200 Subject: [PATCH 2/2] tests: pass the topic to NewP2PMeshAdapter in the pubsub integration test The constructor takes the joined topic now; this caller was outside the packages built locally and only CI's typecheck caught it. --- tests/integration/controlplane_pubsub_integration_test.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/tests/integration/controlplane_pubsub_integration_test.go b/tests/integration/controlplane_pubsub_integration_test.go index 743eb742..2af22112 100644 --- a/tests/integration/controlplane_pubsub_integration_test.go +++ b/tests/integration/controlplane_pubsub_integration_test.go @@ -77,8 +77,12 @@ func TestControlPlanePubSubEventIntegration(t *testing.T) { if err != nil { t.Fatalf("failed to create cp pubsub: %v", err) } + cpTopic, err := cpPS.Join(api.GossipEvents) + if err != nil { + t.Fatalf("failed to join cp topic: %v", err) + } - meshAdapter, err := controlplane.NewP2PMeshAdapter(cpHost, cpPS, store) + meshAdapter, err := controlplane.NewP2PMeshAdapter(cpHost, cpTopic, store) if err != nil { t.Fatalf("failed to create P2PMeshAdapter: %v", err) }