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 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) }