//go:build p2p package client import ( "context" "encoding/json" "fmt" "log" "sync" "github.com/libp2p/go-libp2p" "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/p2p/discovery/mdns" ) const MeshProtocol = "/aetherforge/mesh/1.0.0" const DiscoveryTag = "aetherforge-mesh-discovery" // MeshNode represents a libp2p peer on the local network. // // Relay limitation: mesh uplink is one-way. Orphaned peers may forward messages // to the Hub through a connected relay node, but Hub responses are not sent back // over the mesh. Treat mesh as a best-effort share/stats uplink, not full C2. type MeshNode struct { mu sync.Mutex host host.Host mdns mdns.Service client *AgentClient } // NewMeshNode creates a new P2P fallback node. func NewMeshNode(c *AgentClient) *MeshNode { return &MeshNode{client: c} } // Start initializes the libp2p host and mDNS discovery. func (m *MeshNode) Start() error { m.mu.Lock() defer m.mu.Unlock() if m.host != nil { return nil } // Bind to any available local port automatically h, err := libp2p.New(libp2p.ListenAddrStrings("/ip4/0.0.0.0/tcp/0")) if err != nil { return err } // Register the protocol handler for incoming mesh streams h.SetStreamHandler(MeshProtocol, m.handleStream) // Start mDNS discovery to find other agents on the LAN ser := mdns.NewMdnsService(h, DiscoveryTag, m) if err := ser.Start(); err != nil { _ = h.Close() return err } m.host = h m.mdns = ser log.Printf("[Mesh] P2P Node started. PeerID: %s", h.ID().String()) return nil } // Stop tears down mDNS discovery and the libp2p host. Safe to call multiple times. func (m *MeshNode) Stop() { m.mu.Lock() defer m.mu.Unlock() if m.mdns != nil { _ = m.mdns.Close() m.mdns = nil } if m.host != nil { _ = m.host.Close() m.host = nil } } // HandlePeerFound is a callback for mDNS discovery. func (m *MeshNode) HandlePeerFound(pi peer.AddrInfo) { m.mu.Lock() h := m.host m.mu.Unlock() if h == nil || pi.ID == h.ID() { return } log.Printf("[Mesh] Discovered peer on LAN: %s", pi.ID.String()) if err := h.Connect(context.Background(), pi); err != nil { log.Printf("[Mesh] Failed to connect to peer %s: %v", pi.ID, err) } } // handleStream processes incoming messages from orphaned peers. func (m *MeshNode) handleStream(s network.Stream) { defer s.Close() var msg Message if err := json.NewDecoder(s).Decode(&msg); err != nil { return } if err := m.relayToHub(msg); err == nil { log.Printf("[Mesh] Relaying %s message from orphaned peer to Hub", msg.Type) } } // relayToHub forwards a mesh message to the active Hub WebSocket session. // write() acquires AgentClient.mu — never read client.conn directly (data race). func (m *MeshNode) relayToHub(msg Message) error { if m.client == nil { return fmt.Errorf("mesh: no agent client") } return m.client.write(msg) } // BroadcastToMesh sends a message to all connected P2P peers. func (m *MeshNode) BroadcastToMesh(msg Message) { m.mu.Lock() h := m.host m.mu.Unlock() if h == nil { return } for _, p := range h.Network().Peers() { s, err := h.NewStream(context.Background(), p, MeshProtocol) if err != nil { continue } _ = json.NewEncoder(s).Encode(msg) s.Close() } } // PeerCount returns the number of connected mesh peers. func (m *MeshNode) PeerCount() int { m.mu.Lock() h := m.host m.mu.Unlock() if h == nil { return 0 } return len(h.Network().Peers()) }