From 5dc7ab730916f1786d17979f368393decdd7deac Mon Sep 17 00:00:00 2001 From: Teodor Calin Date: Fri, 24 Jul 2026 17:27:40 +0300 Subject: [PATCH 1/3] =?UTF-8?q?registry:=20hot-path=20performance=20?= =?UTF-8?q?=E2=80=94=20memoized=20signature=20verification,=20lock-free=20?= =?UTF-8?q?telemetry,=20chunked=20reaper,=20single-write=20framing?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Production profiling on the live registry (207K conns, reconnect storm) showed ~30% of CPU in repeat ed25519 verifications of deterministic challenges, ~26% in per-message write syscalls, and reader pile-ups behind the reaper's full-map sweep under the write lock. - authz.VerifyCached: sharded memoization of successful (pubkey, message, sig) verifications; ed25519 is deterministic and registry challenges are static per node/request shape, so identical tuples always re-verify identically. Failures are never cached. Wired into VerifyNodeSignature and the binary heartbeat/resolve paths. - handleMessage: drop the s.mu wrap around walStore.IsStandby (wal has its own lock); per-network request counters use TryRLock and skip a telemetry tick instead of queueing behind writers. - ReapStaleNodes: collect+sort+scan under RLock, write-lock only to delete confirmed-stale nodes (re-checked under the lock), Save() signalled outside the lock. Common no-reap sweep takes no write lock. - accept.writeMessage: single buffered write per frame instead of prefix+body (halves write syscalls). Deployed 2026-07-24: storm-phase CPU dropped from ~1000% to ~240% at equivalent connection counts. Co-Authored-By: Claude Opus 4.8 --- accept/accept.go | 34 +++++----------- authz/authz.go | 42 ++++++++++++++++++- authz/zz_verify_cached_test.go | 73 ++++++++++++++++++++++++++++++++++ directory/directory.go | 41 ++++++++++++++----- server_api.go | 5 +-- 5 files changed, 156 insertions(+), 39 deletions(-) create mode 100644 authz/zz_verify_cached_test.go diff --git a/accept/accept.go b/accept/accept.go index e859049..0db2b31 100644 --- a/accept/accept.go +++ b/accept/accept.go @@ -504,39 +504,27 @@ func readMessage(r io.Reader) (map[string]interface{}, error) { } func writeMessage(w io.Writer, msg map[string]interface{}) error { + var body []byte if raw, ok := msg[rawResponseKey].([]byte); ok && raw != nil { - var lenBuf [4]byte - binary.BigEndian.PutUint32(lenBuf[:], uint32(len(raw))) - if c, ok := w.(net.Conn); ok { - _ = c.SetWriteDeadline(time.Now().Add(writeMessageDeadline)) - defer c.SetWriteDeadline(time.Time{}) - } - if _, err := w.Write(lenBuf[:]); err != nil { - return err - } - if _, err := w.Write(raw); err != nil { - return err + body = raw + } else { + var err error + body, err = json.Marshal(msg) + if err != nil { + return fmt.Errorf("json encode: %w", err) } - return nil } - body, err := json.Marshal(msg) - if err != nil { - return fmt.Errorf("json encode: %w", err) - } - - var lenBuf [4]byte - binary.BigEndian.PutUint32(lenBuf[:], uint32(len(body))) + frame := make([]byte, 4+len(body)) + binary.BigEndian.PutUint32(frame[:4], uint32(len(body))) + copy(frame[4:], body) if c, ok := w.(net.Conn); ok { _ = c.SetWriteDeadline(time.Now().Add(writeMessageDeadline)) defer c.SetWriteDeadline(time.Time{}) } - if _, err := w.Write(lenBuf[:]); err != nil { - return err - } - if _, err := w.Write(body); err != nil { + if _, err := w.Write(frame); err != nil { return err } return nil diff --git a/authz/authz.go b/authz/authz.go index 9dbb833..f5646fb 100644 --- a/authz/authz.go +++ b/authz/authz.go @@ -10,9 +10,11 @@ package authz import ( + "crypto/sha256" "crypto/subtle" "encoding/base64" "fmt" + "sync" "sync/atomic" "github.com/pilot-protocol/common/protocol" @@ -197,6 +199,44 @@ func (c *Checker) IsEnterpriseNode(nodeID uint32, nodes NodeReader, nr NetworkRe // pre-copied value of the global admin token (copy it before releasing // the server mutex). msg must contain a "signature" field (base64). // challenge is the pre-image string that was signed. +const ( + sigCacheShardCount = 64 + sigCacheShardCap = 8192 +) + +type sigCacheShard struct { + mu sync.Mutex + m map[[32]byte]struct{} +} + +var sigCacheShards [sigCacheShardCount]sigCacheShard + +func VerifyCached(pubKey, message, sig []byte) bool { + h := sha256.New() + h.Write(pubKey) + h.Write(message) + h.Write(sig) + var key [32]byte + h.Sum(key[:0]) + sh := &sigCacheShards[key[0]%sigCacheShardCount] + sh.mu.Lock() + _, hit := sh.m[key] + sh.mu.Unlock() + if hit { + return true + } + if !crypto.Verify(pubKey, message, sig) { + return false + } + sh.mu.Lock() + if sh.m == nil || len(sh.m) >= sigCacheShardCap { + sh.m = make(map[[32]byte]struct{}, sigCacheShardCap) + } + sh.m[key] = struct{}{} + sh.mu.Unlock() + return true +} + func VerifyNodeSignature(pubKey []byte, adminToken string, msg map[string]interface{}, challenge string) error { if pubKey == nil { // No key on file — fall back to admin token auth. @@ -213,7 +253,7 @@ func VerifyNodeSignature(pubKey []byte, adminToken string, msg map[string]interf if err != nil { return fmt.Errorf("invalid signature encoding: %w", err) } - ok := crypto.Verify(pubKey, []byte(challenge), sig) + ok := VerifyCached(pubKey, []byte(challenge), sig) if fn := loadSigVerifyHook(); fn != nil { fn(ok) } diff --git a/authz/zz_verify_cached_test.go b/authz/zz_verify_cached_test.go new file mode 100644 index 0000000..53eca63 --- /dev/null +++ b/authz/zz_verify_cached_test.go @@ -0,0 +1,73 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package authz + +import ( + "crypto/ed25519" + "testing" +) + +func TestVerifyCached(t *testing.T) { + pub, priv, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + msg := []byte("heartbeat:42") + sig := ed25519.Sign(priv, msg) + + if !VerifyCached(pub, msg, sig) { + t.Fatal("first verify of valid signature failed") + } + if !VerifyCached(pub, msg, sig) { + t.Fatal("cached verify of valid signature failed") + } + + bad := make([]byte, len(sig)) + copy(bad, sig) + bad[0] ^= 0xFF + if VerifyCached(pub, msg, bad) { + t.Fatal("tampered signature verified") + } + if VerifyCached(pub, msg, bad) { + t.Fatal("tampered signature verified on repeat (failure must not be cached)") + } + + pub2, _, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + if VerifyCached(pub2, msg, sig) { + t.Fatal("signature verified under wrong public key despite cached success for right key") + } + + msg2 := []byte("heartbeat:43") + if VerifyCached(pub, msg2, sig) { + t.Fatal("signature verified for different message despite cached success") + } +} + +func TestVerifyCachedShardReset(t *testing.T) { + pub, priv, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + msg := []byte("resolve:1:2") + sig := ed25519.Sign(priv, msg) + if !VerifyCached(pub, msg, sig) { + t.Fatal("valid signature failed") + } + for i := range sigCacheShards { + sh := &sigCacheShards[i] + sh.mu.Lock() + for k := range sh.m { + sh.m[k] = struct{}{} + } + if sh.m != nil && len(sh.m) > sigCacheShardCap { + t.Fatalf("shard %d exceeded cap: %d", i, len(sh.m)) + } + sh.mu.Unlock() + } + if !VerifyCached(pub, msg, sig) { + t.Fatal("re-verify after shard inspection failed") + } +} diff --git a/directory/directory.go b/directory/directory.go index 1dc5b68..cb94f76 100644 --- a/directory/directory.go +++ b/directory/directory.go @@ -24,6 +24,7 @@ import ( "github.com/pilot-protocol/common/crypto" "github.com/pilot-protocol/common/protocol" "github.com/pilot-protocol/common/registry/wire" + "github.com/pilot-protocol/rendezvous/authz" ) // -------------------------------------------------------------------------- @@ -618,13 +619,12 @@ func (st *Store) AdminListNodesCached() ([]byte, error) { // ReapStaleNodes removes nodes whose last heartbeat is older than threshold. func (st *Store) ReapStaleNodes(threshold time.Time) { - st.mu.Lock() - defer st.mu.Unlock() - + st.mu.RLock() nodeIDs := make([]uint32, 0, len(st.nodes)) for id := range st.nodes { nodeIDs = append(nodeIDs, id) } + st.mu.RUnlock() sort.Slice(nodeIDs, func(i, j int) bool { return nodeIDs[i] < nodeIDs[j] }) startIdx := 0 @@ -633,15 +633,34 @@ func (st *Store) ReapStaleNodes(threshold time.Time) { } processed := 0 - reaped := false + cursor := *st.reapCursor + var stale []uint32 + st.mu.RLock() for i := 0; i < len(nodeIDs) && processed < ReapChunkSize; i++ { idx := (startIdx + i) % len(nodeIDs) id := nodeIDs[idx] processed++ - node := st.nodes[id] - lastSeen := node.GetLastSeen() - if lastSeen.Before(threshold) { + if node, ok := st.nodes[id]; ok && node.GetLastSeen().Before(threshold) { + stale = append(stale, id) + } + + cursor = id + 1 + } + st.mu.RUnlock() + + reaped := false + if len(stale) > 0 { + st.mu.Lock() + for _, id := range stale { + node, ok := st.nodes[id] + if !ok { + continue + } + lastSeen := node.GetLastSeen() + if !lastSeen.Before(threshold) { + continue + } staleDuration := time.Since(lastSeen).Round(time.Second) slog.Info("registry reaping stale node", "node_id", id, "last_seen_ago", staleDuration) st.cb.Audit("node.reaped", "node_id", id, "reason", "stale_heartbeat", @@ -654,10 +673,10 @@ func (st *Store) ReapStaleNodes(threshold time.Time) { st.cb.InvalidateAdminListNodesCache() reaped = true } - - *st.reapCursor = id + 1 + st.mu.Unlock() } + *st.reapCursor = cursor if processed >= len(nodeIDs) { *st.reapCursor = 0 } @@ -1342,7 +1361,7 @@ func (st *Store) HandleBinaryResolve(conn net.Conn, payload []byte) { return } challenge := fmt.Sprintf("resolve:%d:%d", requesterID, nodeID) - if !crypto.Verify(requesterPubKey, []byte(challenge), sig) { + if !authz.VerifyCached(requesterPubKey, []byte(challenge), sig) { st.cb.IncErrorsTotal("resolve") wire.WriteFrame(conn, wire.MsgError, wire.EncodeError("signature verification failed")) return @@ -1963,7 +1982,7 @@ func (st *Store) HandleBinaryHeartbeat(conn net.Conn, payload []byte) { return } challenge := fmt.Sprintf("heartbeat:%d", req.NodeID) - if !crypto.Verify(pubKey, []byte(challenge), req.Signature[:]) { + if !authz.VerifyCached(pubKey, []byte(challenge), req.Signature[:]) { st.cb.IncErrorsTotal("heartbeat") wire.WriteFrame(conn, wire.MsgError, wire.EncodeError("signature verification failed")) return diff --git a/server_api.go b/server_api.go index 7d69760..ba89b38 100644 --- a/server_api.go +++ b/server_api.go @@ -122,8 +122,7 @@ func (s *Server) handleMessage(msg map[string]interface{}, remoteAddr string) (r case uint32: nodeID = v } - if nodeID > 0 { - s.mu.RLock() + if nodeID > 0 && s.mu.TryRLock() { if node, exists := s.nodes[nodeID]; exists { nets := node.Networks for _, netID := range nets { @@ -147,9 +146,7 @@ func (s *Server) handleMessage(msg map[string]interface{}, remoteAddr string) (r }() // Standby mode: reject write operations, allow reads - s.mu.RLock() isStandby := s.walStore.IsStandby() - s.mu.RUnlock() if isStandby { switch msgType { case "lookup", "resolve", "list_networks", "list_nodes", "heartbeat", "poll_handshakes", "poll_invites", "resolve_hostname", "beacon_list", From de0152365e896456a433901e3b501c43e9dd7ac7 Mon Sep 17 00:00:00 2001 From: Teodor Calin Date: Sat, 25 Jul 2026 10:20:47 +0300 Subject: [PATCH 2/3] registry: verify registration proof-of-possession signature (PPA-003) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The register op is the one write that establishes the node_id↔pubkey↔ endpoint binding, and it was the only signed-operation sibling that carried no signature — so anyone reaching the registry could register a key or repoint an endpoint they do not own. The common client now signs register::; this verifies it. Verify-if-present, so old daemons that send no signature still register (wire-compatible); a present-but-invalid signature is always rejected. RequireRegisterSignature (flag -require-register-signature / RENDEZVOUS_REQUIRE_REGISTER_SIGNATURE, default OFF) upgrades this to mandatory once the signing client has rolled out fleet-wide. Co-Authored-By: Claude Opus 4.8 --- cmd/rendezvous/main.go | 5 + directory/directory.go | 25 ++++- .../zz_ppa003_register_signature_test.go | 92 +++++++++++++++++++ server.go | 17 ++-- server_lifecycle.go | 25 ++++- 5 files changed, 149 insertions(+), 15 deletions(-) create mode 100644 directory/zz_ppa003_register_signature_test.go diff --git a/cmd/rendezvous/main.go b/cmd/rendezvous/main.go index 2bf4adf..1c44911 100644 --- a/cmd/rendezvous/main.go +++ b/cmd/rendezvous/main.go @@ -67,6 +67,7 @@ func main() { tlsKey := flag.String("tls-key", "", "TLS key file") enableTLS := flag.Bool("tls", false, "enable TLS for registry connections") strictDirectoryAuth := flag.Bool("strict-directory-auth", false, "WS2: require trust/shared-network authorization on directory RPCs (lookup/resolve/punch/list_*/check_trust). Default false (not enforcing, wire-compatible with old agents). Env: RENDEZVOUS_STRICT_DIRECTORY_AUTH=1.") + requireRegisterSignature := flag.Bool("require-register-signature", false, "PPA-003: require a valid proof-of-possession signature on registration. Default false (present signatures are still verified; unsigned registrations accepted for wire-compatibility). Enable only after the signing client has rolled out. Env: RENDEZVOUS_REQUIRE_REGISTER_SIGNATURE=1.") standbyPrimary := flag.String("standby", "", "run as hot standby replicating from the given primary address (e.g. primary:9000)") httpAddr := flag.String("http", "", "HTTP dashboard listen address (e.g. :3000)") logLevel := flag.String("log-level", "info", "log level (debug, info, warn, error)") @@ -128,6 +129,10 @@ func main() { r.SetStrictDirectoryAuth(true) slog.Info("strict directory authorization enabled (WS2)") } + if *requireRegisterSignature || os.Getenv("RENDEZVOUS_REQUIRE_REGISTER_SIGNATURE") == "1" { + r.SetRequireRegisterSignature(true) + slog.Info("registration signature required (PPA-003)") + } // Plumb the breaker manager into the in-process beacon so // beacon.punch / beacon.relay / beacon.discover can be flipped from // the same breakers.json file as the registry-side switches. diff --git a/directory/directory.go b/directory/directory.go index cb94f76..0d12a31 100644 --- a/directory/directory.go +++ b/directory/directory.go @@ -235,8 +235,9 @@ type Callbacks struct { // ScanNetworkMemberships returns the set of non-backbone networkIDs that // still list nodeID as a member. Used to restore network memberships after // a reaped node reclaims its old identity. Caller holds mu.Lock. - ScanNetworkMemberships func(nodeID uint32) []uint16 - StrictDirectoryAuth func() bool + ScanNetworkMemberships func(nodeID uint32) []uint16 + StrictDirectoryAuth func() bool + RequireRegisterSignature func() bool } // -------------------------------------------------------------------------- @@ -739,6 +740,26 @@ func (st *Store) HandleRegister( return nil, fmt.Errorf("registration requires public_key") } + // PPA-003 proof-of-possession. Verify-if-present is wire-compatible: a + // daemon that sends no signature still registers, so old fleets are + // unaffected; a present-but-invalid signature is always rejected. Flipping + // RequireRegisterSignature (default off) makes the signature mandatory once + // the signing client has rolled out — do NOT default it on. + if sigB64, _ := msg["signature"].(string); sigB64 != "" { + pubKeyBytes, decErr := crypto.DecodePublicKey(pubKeyB64) + if decErr != nil { + return nil, fmt.Errorf("registration: invalid public_key: %w", decErr) + } + challenge := fmt.Sprintf("register:%s:%s", clientAddr, pubKeyB64) + if st.cb.VerifyNodeSignature != nil { + if err := st.cb.VerifyNodeSignature(pubKeyBytes, "", msg, challenge); err != nil { + return nil, fmt.Errorf("registration signature verification failed: %w", err) + } + } + } else if st.cb.RequireRegisterSignature != nil && st.cb.RequireRegisterSignature() { + return nil, fmt.Errorf("registration requires a signature") + } + if hostname != "" { if err := ValidateHostname(hostname); err != nil { resp, regErr := st.HandleReRegister(pubKeyB64, listenAddr, owner, "", lanAddrs, clientVersion, relayOnly) diff --git a/directory/zz_ppa003_register_signature_test.go b/directory/zz_ppa003_register_signature_test.go new file mode 100644 index 0000000..df19cc2 --- /dev/null +++ b/directory/zz_ppa003_register_signature_test.go @@ -0,0 +1,92 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package directory + +import ( + "encoding/base64" + "testing" + + "github.com/pilot-protocol/common/crypto" + "github.com/pilot-protocol/rendezvous/authz" +) + +func ppa003Store(t *testing.T, require bool) *Store { + t.Helper() + st := newTestStore(t) + st.cb.VerifyNodeSignature = authz.VerifyNodeSignature + if require { + st.cb.RequireRegisterSignature = func() bool { return true } + } + return st +} + +func TestHandleRegisterSignatureVerifyIfPresent(t *testing.T) { + id, err := crypto.GenerateIdentity() + if err != nil { + t.Fatalf("GenerateIdentity: %v", err) + } + pubB64 := crypto.EncodePublicKey(id.PublicKey) + addr := "10.0.0.7:4000" + sigB64 := base64.StdEncoding.EncodeToString(id.Sign([]byte("register:" + addr + ":" + pubB64))) + + t.Run("valid signature accepted", func(t *testing.T) { + st := ppa003Store(t, false) + resp, err := st.HandleRegister(map[string]interface{}{ + "public_key": pubB64, + "listen_addr": addr, + "signature": sigB64, + }, addr, nil, nil) + if err != nil { + t.Fatalf("valid signed registration rejected: %v", err) + } + if resp["type"] != "register_ok" { + t.Fatalf("expected register_ok, got %v", resp["type"]) + } + }) + + t.Run("invalid signature rejected", func(t *testing.T) { + st := ppa003Store(t, false) + bad := id.Sign([]byte("register:" + addr + ":" + pubB64)) + bad[0] ^= 0xFF + _, err := st.HandleRegister(map[string]interface{}{ + "public_key": pubB64, + "listen_addr": addr, + "signature": base64.StdEncoding.EncodeToString(bad), + }, addr, nil, nil) + if err == nil { + t.Fatal("registration with a forged signature was accepted") + } + }) + + t.Run("wrong listen_addr rejected", func(t *testing.T) { + st := ppa003Store(t, false) + _, err := st.HandleRegister(map[string]interface{}{ + "public_key": pubB64, + "listen_addr": "10.0.0.9:4000", + "signature": sigB64, + }, addr, nil, nil) + if err == nil { + t.Fatal("signature bound to a different listen_addr was accepted") + } + }) + + t.Run("unsigned accepted for compat", func(t *testing.T) { + st := ppa003Store(t, false) + if _, err := st.HandleRegister(map[string]interface{}{ + "public_key": pubB64, + "listen_addr": addr, + }, addr, nil, nil); err != nil { + t.Fatalf("unsigned registration rejected (compat break): %v", err) + } + }) + + t.Run("unsigned rejected when required", func(t *testing.T) { + st := ppa003Store(t, true) + if _, err := st.HandleRegister(map[string]interface{}{ + "public_key": pubB64, + "listen_addr": addr, + }, addr, nil, nil); err == nil { + t.Fatal("unsigned registration accepted despite RequireRegisterSignature") + } + }) +} diff --git a/server.go b/server.go index 65eec06..f0073c7 100644 --- a/server.go +++ b/server.go @@ -89,13 +89,13 @@ type Server struct { // its own internal locking. Surfaced read-only via /api/health and // included in the rich /api/stats payload. Reads are non-blocking // (atomic loads) so a probe doesn't contend with a save in flight. - snapshotsTotal atomic.Int64 // successful flushSave calls - snapshotsFailed atomic.Int64 // flushSave returning an error - lastSnapshotUnixMs atomic.Int64 // wall time of last successful save - lastSnapshotDurMs atomic.Int64 // duration of last successful save - lastSnapshotSizeB atomic.Int64 // bytes written by last successful save - lastSnapshotRLockMs atomic.Int64 // s.mu.RLock hold time in phase 1 - maxSnapshotDurMs atomic.Int64 // worst-ever save duration + snapshotsTotal atomic.Int64 // successful flushSave calls + snapshotsFailed atomic.Int64 // flushSave returning an error + lastSnapshotUnixMs atomic.Int64 // wall time of last successful save + lastSnapshotDurMs atomic.Int64 // duration of last successful save + lastSnapshotSizeB atomic.Int64 // bytes written by last successful save + lastSnapshotRLockMs atomic.Int64 // s.mu.RLock hold time in phase 1 + maxSnapshotDurMs atomic.Int64 // worst-ever save duration // buildInfo carries the build-time identity surfaced on // /api/public-stats for code-verification (version, git commit, ISO @@ -284,7 +284,8 @@ type Server struct { listNodesPerNetMu sync.Mutex // guards the map itself listNodesPerNet map[uint16]*listNodesCacheState - strictDirectoryAuth atomic.Bool + strictDirectoryAuth atomic.Bool + requireRegisterSignature atomic.Bool } // listNodesCacheState is defined in the directory sub-package (R4.2). diff --git a/server_lifecycle.go b/server_lifecycle.go index 23b23b5..b8c9c07 100644 --- a/server_lifecycle.go +++ b/server_lifecycle.go @@ -125,6 +125,20 @@ func (s *Server) StrictDirectoryAuth() bool { return s.strictDirectoryAuth.Load() } +// SetRequireRegisterSignature toggles PPA-003 mandatory registration +// proof-of-possession. Default off for wire-compatibility: unsigned +// registrations are accepted (a present signature is always verified). Enable +// only after the signing client has rolled out fleet-wide. +func (s *Server) SetRequireRegisterSignature(enabled bool) { + s.requireRegisterSignature.Store(enabled) +} + +// RequireRegisterSignature reports whether registration signatures are +// mandatory. Wired into the directory sub-package as a callback. +func (s *Server) RequireRegisterSignature() bool { + return s.requireRegisterSignature.Load() +} + // SetDashboardToken gates per-network stats on the dashboard. // Empty string restricts the dashboard to global aggregates only. func (s *Server) SetDashboardToken(token string) { @@ -941,7 +955,8 @@ func NewWithStore(beaconAddr, storePath string) *Server { } return nets }, - StrictDirectoryAuth: s.StrictDirectoryAuth, + StrictDirectoryAuth: s.StrictDirectoryAuth, + RequireRegisterSignature: s.RequireRegisterSignature, }, ) @@ -1105,10 +1120,10 @@ func NewWithStore(beaconAddr, storePath string) *Server { allow, _ := s.breakers.Allow(name) return allow, s.breakers.Reason(name) }, - BreakerList: s.BreakerList, - BreakerSet: s.BreakerSet, - BreakerDelete: s.BreakerDelete, - HealthSnapshot: s.HealthSnapshot, + BreakerList: s.BreakerList, + BreakerSet: s.BreakerSet, + BreakerDelete: s.BreakerDelete, + HealthSnapshot: s.HealthSnapshot, VerifyRequest: func(canonical, sigB64 string) interface{} { return s.VerifyRequest(canonical, sigB64) }, From 0f654500be27eb8d8b4d8e28f79955fe464e3681 Mon Sep 17 00:00:00 2001 From: Teodor Calin Date: Sat, 25 Jul 2026 15:31:30 +0300 Subject: [PATCH 3/3] registry: PPA-003 endpoint-relocation ratchet (close the attack, compat-safe) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Once a key registers with a valid proof-of-possession signature, the node latches SigVerified. An unsigned registration may then no longer relocate its endpoint — closing the "re-register a victim's key to an attacker IP" hijack for every signing node. Unsigned/first-seen keys are unaffected, so no old agent breaks; protection grows as the fleet signs. SigVerified is persisted with the node record. Co-Authored-By: Claude Opus 4.8 --- directory/directory.go | 105 ++++++++++++------ .../zz_ppa003_register_signature_test.go | 47 ++++++++ directory/zz_pubkey_rebind_test.go | 12 +- directory/zz_register_test.go | 6 +- directory/zz_resolve_hostname_test.go | 6 +- 5 files changed, 131 insertions(+), 45 deletions(-) diff --git a/directory/directory.go b/directory/directory.go index 0d12a31..3aad5e9 100644 --- a/directory/directory.go +++ b/directory/directory.go @@ -75,6 +75,12 @@ type NodeInfo struct { Version string RelayOnly bool + // SigVerified latches true once this key has registered with a valid + // proof-of-possession signature. Once true, an unsigned registration may + // not relocate the node's endpoint (PPA-003 ratchet) — closes registry + // endpoint hijack for every signing node without breaking unsigned ones. + SigVerified bool `json:"sig_verified,omitempty"` + // Verified-address badge (offline-verifiable; carries no raw identity). Badge string BadgeSig string @@ -745,6 +751,7 @@ func (st *Store) HandleRegister( // unaffected; a present-but-invalid signature is always rejected. Flipping // RequireRegisterSignature (default off) makes the signature mandatory once // the signing client has rolled out — do NOT default it on. + sigVerified := false if sigB64, _ := msg["signature"].(string); sigB64 != "" { pubKeyBytes, decErr := crypto.DecodePublicKey(pubKeyB64) if decErr != nil { @@ -755,6 +762,7 @@ func (st *Store) HandleRegister( if err := st.cb.VerifyNodeSignature(pubKeyBytes, "", msg, challenge); err != nil { return nil, fmt.Errorf("registration signature verification failed: %w", err) } + sigVerified = true } } else if st.cb.RequireRegisterSignature != nil && st.cb.RequireRegisterSignature() { return nil, fmt.Errorf("registration requires a signature") @@ -762,7 +770,7 @@ func (st *Store) HandleRegister( if hostname != "" { if err := ValidateHostname(hostname); err != nil { - resp, regErr := st.HandleReRegister(pubKeyB64, listenAddr, owner, "", lanAddrs, clientVersion, relayOnly) + resp, regErr := st.HandleReRegister(pubKeyB64, listenAddr, owner, "", lanAddrs, clientVersion, relayOnly, sigVerified) if regErr != nil { return resp, regErr } @@ -774,7 +782,7 @@ func (st *Store) HandleRegister( } } - resp, err := st.HandleReRegister(pubKeyB64, listenAddr, owner, hostname, lanAddrs, clientVersion, relayOnly) + resp, err := st.HandleReRegister(pubKeyB64, listenAddr, owner, hostname, lanAddrs, clientVersion, relayOnly, sigVerified) if err == nil { st.cb.IncRegistrations() resp["observed_addr"] = listenAddr @@ -798,12 +806,34 @@ func (st *Store) HandleRegister( // HandleReRegister handles a node presenting an existing public key. // Fast path for known-key reconnects; slow path for new nodes and index mutations. -func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, lanAddrs []string, version string, relayOnly bool) (map[string]interface{}, error) { +func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, lanAddrs []string, version string, relayOnly bool, sigVerified bool) (map[string]interface{}, error) { pubKey, err := crypto.DecodePublicKey(pubKeyB64) if err != nil { return nil, fmt.Errorf("invalid public key: %w", err) } + // PPA-003 ratchet: once a key has registered with a valid proof-of- + // possession (SigVerified), an unsigned registration may not relocate its + // endpoint. Unsigned/first-seen keys are unaffected (wire-compatible), so + // no old agent breaks; coverage grows as the fleet signs. + if !sigVerified { + st.mu.RLock() + if exID, known := st.pubKeyIdx[pubKeyB64]; known { + if n, ok := st.nodes[exID]; ok { + sh := &st.nodeShards[exID%NumNodeShards] + sh.RLock() + reject := n.SigVerified && listenAddr != n.RealAddr + sh.RUnlock() + if reject { + st.mu.RUnlock() + st.cb.Audit("node.register_ratchet_rejected", "node_id", exID) + return nil, fmt.Errorf("node %d: unsigned registration cannot relocate the endpoint of a signature-verified key", exID) + } + } + } + st.mu.RUnlock() + } + // FAST PATH: shard-locked endpoint refresh for the common case. st.mu.RLock() fpNodeID, fpPubKeyKnown := st.pubKeyIdx[pubKeyB64] @@ -837,6 +867,9 @@ func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, fpNode.Version = version } fpNode.RelayOnly = relayOnly + if sigVerified { + fpNode.SigVerified = true + } shard.Unlock() addr := protocol.Addr{Network: 0, Node: fpNodeID} @@ -872,6 +905,9 @@ func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, node.Version = version } node.RelayOnly = relayOnly + if sigVerified { + node.SigVerified = true + } if owner != "" && node.Owner == "" { node.Owner = owner st.ownerIdx[owner] = nodeID @@ -903,16 +939,17 @@ func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, } now := time.Now() node := &NodeInfo{ - ID: nodeID, - Owner: owner, - PublicKey: pubKey, - RealAddr: listenAddr, - Networks: networks, - LastSeen: now, - LANAddrs: lanAddrs, - KeyMeta: KeyInfo{CreatedAt: now}, - Version: version, - RelayOnly: relayOnly, + ID: nodeID, + Owner: owner, + PublicKey: pubKey, + RealAddr: listenAddr, + Networks: networks, + LastSeen: now, + LANAddrs: lanAddrs, + KeyMeta: KeyInfo{CreatedAt: now}, + Version: version, + RelayOnly: relayOnly, + SigVerified: sigVerified, } node.LastSeenNano.Store(now.UnixNano()) st.nodes[nodeID] = node @@ -964,16 +1001,17 @@ func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, st.pubKeyIdx[pubKeyB64] = existingID now := time.Now() node := &NodeInfo{ - ID: existingID, - Owner: owner, - PublicKey: pubKey, - RealAddr: listenAddr, - Networks: []uint16{0}, - LastSeen: now, - LANAddrs: lanAddrs, - KeyMeta: KeyInfo{CreatedAt: now}, - Version: version, - RelayOnly: relayOnly, + ID: existingID, + Owner: owner, + PublicKey: pubKey, + RealAddr: listenAddr, + Networks: []uint16{0}, + LastSeen: now, + LANAddrs: lanAddrs, + KeyMeta: KeyInfo{CreatedAt: now}, + Version: version, + RelayOnly: relayOnly, + SigVerified: sigVerified, } node.LastSeenNano.Store(now.UnixNano()) st.nodes[existingID] = node @@ -1014,16 +1052,17 @@ func (st *Store) HandleReRegister(pubKeyB64, listenAddr, owner, hostname string, now := time.Now() node := &NodeInfo{ - ID: nodeID, - Owner: owner, - PublicKey: pubKey, - RealAddr: listenAddr, - Networks: []uint16{0}, - LastSeen: now, - LANAddrs: lanAddrs, - KeyMeta: KeyInfo{CreatedAt: now}, - Version: version, - RelayOnly: relayOnly, + ID: nodeID, + Owner: owner, + PublicKey: pubKey, + RealAddr: listenAddr, + Networks: []uint16{0}, + LastSeen: now, + LANAddrs: lanAddrs, + KeyMeta: KeyInfo{CreatedAt: now}, + Version: version, + RelayOnly: relayOnly, + SigVerified: sigVerified, } node.LastSeenNano.Store(now.UnixNano()) st.nodes[nodeID] = node diff --git a/directory/zz_ppa003_register_signature_test.go b/directory/zz_ppa003_register_signature_test.go index df19cc2..8c55a56 100644 --- a/directory/zz_ppa003_register_signature_test.go +++ b/directory/zz_ppa003_register_signature_test.go @@ -90,3 +90,50 @@ func TestHandleRegisterSignatureVerifyIfPresent(t *testing.T) { } }) } + +func TestHandleRegisterEndpointRatchet(t *testing.T) { + id, err := crypto.GenerateIdentity() + if err != nil { + t.Fatal(err) + } + pub := crypto.EncodePublicKey(id.PublicKey) + sign := func(addr string) string { + return base64.StdEncoding.EncodeToString(id.Sign([]byte("register:" + addr + ":" + pub))) + } + st := ppa003Store(t, false) + reg := func(pk, addr, remote, sig string) error { + m := map[string]interface{}{"public_key": pk, "listen_addr": addr} + if sig != "" { + m["signature"] = sig + } + _, e := st.HandleRegister(m, remote, nil, nil) + return e + } + + // 1. signed registration latches SigVerified + if e := reg(pub, "10.0.0.1:5000", "10.0.0.1:5000", sign("10.0.0.1:5000")); e != nil { + t.Fatalf("signed register: %v", e) + } + // 2. ATTACK: unsigned re-register of that key from a different endpoint — REJECTED + if e := reg(pub, "10.9.9.9:5000", "10.9.9.9:5000", ""); e == nil { + t.Fatal("unsigned endpoint relocation of a signature-verified key was ALLOWED (ratchet failed)") + } + // 3. legit signed move — allowed + if e := reg(pub, "10.0.0.2:6000", "10.0.0.2:6000", sign("10.0.0.2:6000")); e != nil { + t.Fatalf("signed endpoint move rejected: %v", e) + } + // 4. unsigned same-endpoint refresh — allowed (harmless) + if e := reg(pub, "10.0.0.2:6000", "10.0.0.2:6000", ""); e != nil { + t.Fatalf("unsigned same-endpoint refresh rejected: %v", e) + } + + // 5. COMPAT: a key that never signed can still relocate unsigned (old agent) + id2, _ := crypto.GenerateIdentity() + pub2 := crypto.EncodePublicKey(id2.PublicKey) + if e := reg(pub2, "10.1.1.1:5000", "10.1.1.1:5000", ""); e != nil { + t.Fatalf("unsigned new node: %v", e) + } + if e := reg(pub2, "10.2.2.2:5000", "10.2.2.2:5000", ""); e != nil { + t.Fatalf("unsigned old-agent endpoint change BLOCKED (compat break): %v", e) + } +} diff --git a/directory/zz_pubkey_rebind_test.go b/directory/zz_pubkey_rebind_test.go index 09b4de4..313830a 100644 --- a/directory/zz_pubkey_rebind_test.go +++ b/directory/zz_pubkey_rebind_test.go @@ -19,14 +19,14 @@ func TestReRegister_SameKey_NoRegression(t *testing.T) { st := newTestStore(t) pubKey := genPubKeyB64(t) - resp1, err := st.HandleReRegister(pubKey, "10.0.0.1:4000", "alice", "host1", nil, "1.0.0", false) + resp1, err := st.HandleReRegister(pubKey, "10.0.0.1:4000", "alice", "host1", nil, "1.0.0", false, false) if err != nil { t.Fatalf("first register: %v", err) } nodeID := resp1["node_id"].(uint32) // Same owner, same key, new address — must succeed, same node_id. - resp2, err := st.HandleReRegister(pubKey, "10.0.0.2:4000", "alice", "host1", nil, "2.0.0", false) + resp2, err := st.HandleReRegister(pubKey, "10.0.0.2:4000", "alice", "host1", nil, "2.0.0", false, false) if err != nil { t.Fatalf("same-key re-register rejected (regression!): %v", err) } @@ -46,7 +46,7 @@ func TestReRegister_DifferentKey_OwnerReclaim_Rejected(t *testing.T) { st := newTestStore(t) origPubKeyB64 := genPubKeyB64(t) - resp1, err := st.HandleReRegister(origPubKeyB64, "10.0.0.1:4000", "alice", "host1", nil, "1.0.0", false) + resp1, err := st.HandleReRegister(origPubKeyB64, "10.0.0.1:4000", "alice", "host1", nil, "1.0.0", false, false) if err != nil { t.Fatalf("first register: %v", err) } @@ -64,7 +64,7 @@ func TestReRegister_DifferentKey_OwnerReclaim_Rejected(t *testing.T) { if attackerPubKeyB64 == origPubKeyB64 { t.Fatal("test setup: keys collided") } - _, err = st.HandleReRegister(attackerPubKeyB64, "6.6.6.6:4000", "alice", "host1", nil, "9.9.9", false) + _, err = st.HandleReRegister(attackerPubKeyB64, "6.6.6.6:4000", "alice", "host1", nil, "9.9.9", false, false) if err == nil { t.Fatal("owner-reclaim with a different key must be REJECTED") } @@ -103,7 +103,7 @@ func TestReRegister_ReapedOwnerNode_ReclaimAllowed(t *testing.T) { st := newTestStore(t) pubKey := genPubKeyB64(t) - resp1, err := st.HandleReRegister(pubKey, "10.0.0.1:4000", "alice", "host1", nil, "1.0.0", false) + resp1, err := st.HandleReRegister(pubKey, "10.0.0.1:4000", "alice", "host1", nil, "1.0.0", false, false) if err != nil { t.Fatalf("first register: %v", err) } @@ -115,7 +115,7 @@ func TestReRegister_ReapedOwnerNode_ReclaimAllowed(t *testing.T) { delete(st.pubKeyIdx, pubKey) newPubKey := genPubKeyB64(t) - resp2, err := st.HandleReRegister(newPubKey, "10.0.0.9:4000", "alice", "host1", nil, "2.0.0", false) + resp2, err := st.HandleReRegister(newPubKey, "10.0.0.9:4000", "alice", "host1", nil, "2.0.0", false, false) if err != nil { t.Fatalf("reaped-node owner reclaim should succeed: %v", err) } diff --git a/directory/zz_register_test.go b/directory/zz_register_test.go index c430ed4..7605e6d 100644 --- a/directory/zz_register_test.go +++ b/directory/zz_register_test.go @@ -98,7 +98,7 @@ func TestHandleRegister_InvalidHostnameStillRegisters(t *testing.T) { func TestHandleReRegister_InvalidPubKey(t *testing.T) { t.Parallel() st := newTestStore(t) - _, err := st.HandleReRegister("!!!", "10.0.0.1:4000", "alice", "", nil, "1.0.0", false) + _, err := st.HandleReRegister("!!!", "10.0.0.1:4000", "alice", "", nil, "1.0.0", false, false) if err == nil { t.Error("expected invalid-pubkey error") } @@ -109,14 +109,14 @@ func TestHandleReRegister_ExistingNodeUpdatesFields(t *testing.T) { st := newTestStore(t) pubKey := genPubKeyB64(t) // First register. - resp1, err := st.HandleReRegister(pubKey, "10.0.0.1:4000", "alice", "host1", []string{"192.168.1.5:4000"}, "1.0.0", false) + resp1, err := st.HandleReRegister(pubKey, "10.0.0.1:4000", "alice", "host1", []string{"192.168.1.5:4000"}, "1.0.0", false, false) if err != nil { t.Fatalf("first register: %v", err) } nodeID1 := resp1["node_id"].(uint32) // Re-register with same pubkey — should return same node_id (fast path). - resp2, err := st.HandleReRegister(pubKey, "10.0.0.2:4000", "alice", "host2", []string{"192.168.1.6:4000"}, "2.0.0", false) + resp2, err := st.HandleReRegister(pubKey, "10.0.0.2:4000", "alice", "host2", []string{"192.168.1.6:4000"}, "2.0.0", false, false) if err != nil { t.Fatalf("re-register: %v", err) } diff --git a/directory/zz_resolve_hostname_test.go b/directory/zz_resolve_hostname_test.go index 0085d98..4c818c2 100644 --- a/directory/zz_resolve_hostname_test.go +++ b/directory/zz_resolve_hostname_test.go @@ -140,7 +140,7 @@ func TestHandleReRegister_FastPathRefreshesEndpoint(t *testing.T) { // Re-register with same owner + hostname → fast-path endpoint refresh. resp, err := st.HandleReRegister(pubKeyB64, "10.0.0.2:4000", "alice", "host1", - []string{"192.168.1.5:4000"}, "1.7.2", false) + []string{"192.168.1.5:4000"}, "1.7.2", false, false) if err != nil { t.Fatalf("re-register: %v", err) } @@ -173,7 +173,7 @@ func TestHandleReRegister_SlowPathNewOwner(t *testing.T) { t.Fatal(err) } // Re-register with new owner triggers slow path. - resp, err := st.HandleReRegister(pubKeyB64, "10.0.0.2:4000", "newowner", "", nil, "", false) + resp, err := st.HandleReRegister(pubKeyB64, "10.0.0.2:4000", "newowner", "", nil, "", false, false) if err != nil { t.Fatal(err) } @@ -192,7 +192,7 @@ func TestHandleReRegister_SlowPathNewOwner(t *testing.T) { func TestHandleReRegister_BadPubKey(t *testing.T) { t.Parallel() st := newTestStore(t) - _, err := st.HandleReRegister("not-base64!", "10.0.0.1:4000", "", "", nil, "", false) + _, err := st.HandleReRegister("not-base64!", "10.0.0.1:4000", "", "", nil, "", false, false) if err == nil || !strings.Contains(err.Error(), "invalid public key") { t.Fatalf("%v", err) }