@@ -24,8 +24,78 @@ use peat_mesh::transport::{
2424} ;
2525use std:: net:: SocketAddr ;
2626use std:: sync:: Arc ;
27+ use std:: time:: Instant ;
28+ use tokio:: sync:: broadcast;
2729use tracing:: { error, info, warn} ;
2830
31+ struct NodeBrokerState {
32+ node_id : String ,
33+ version : String ,
34+ started_at : Instant ,
35+ transport : Arc < MeshSyncTransport > ,
36+ event_tx : broadcast:: Sender < peat_mesh:: broker:: state:: MeshEvent > ,
37+ }
38+
39+ #[ async_trait:: async_trait]
40+ impl peat_mesh:: broker:: state:: MeshBrokerState for NodeBrokerState {
41+ fn node_info ( & self ) -> peat_mesh:: broker:: state:: MeshNodeInfo {
42+ peat_mesh:: broker:: state:: MeshNodeInfo {
43+ node_id : self . node_id . clone ( ) ,
44+ uptime_secs : self . started_at . elapsed ( ) . as_secs ( ) ,
45+ version : self . version . clone ( ) ,
46+ }
47+ }
48+
49+ async fn list_peers ( & self ) -> Vec < peat_mesh:: broker:: state:: PeerSummary > {
50+ self . transport
51+ . connected_peers ( )
52+ . into_iter ( )
53+ . map ( |peer_id| peat_mesh:: broker:: state:: PeerSummary {
54+ id : peer_id. to_string ( ) ,
55+ connected : true ,
56+ state : "active" . to_string ( ) ,
57+ rtt_ms : self
58+ . transport
59+ . peer_rtt ( & peer_id)
60+ . map ( |rtt| rtt. as_millis ( ) as u64 ) ,
61+ } )
62+ . collect ( )
63+ }
64+
65+ async fn get_peer ( & self , id : & str ) -> Option < peat_mesh:: broker:: state:: PeerSummary > {
66+ self . transport
67+ . connected_peers ( )
68+ . into_iter ( )
69+ . find ( |existing| existing. to_string ( ) == id)
70+ . map ( |peer_id| peat_mesh:: broker:: state:: PeerSummary {
71+ id : peer_id. to_string ( ) ,
72+ connected : true ,
73+ state : "active" . to_string ( ) ,
74+ rtt_ms : self
75+ . transport
76+ . peer_rtt ( & peer_id)
77+ . map ( |rtt| rtt. as_millis ( ) as u64 ) ,
78+ } )
79+ }
80+
81+ fn topology ( & self ) -> peat_mesh:: broker:: state:: TopologySummary {
82+ let peer_count = self . transport . connected_peers ( ) . len ( ) ;
83+ peat_mesh:: broker:: state:: TopologySummary {
84+ peer_count,
85+ role : if peer_count > 0 {
86+ "connected" . to_string ( )
87+ } else {
88+ "standalone" . to_string ( )
89+ } ,
90+ hierarchy_level : 0 ,
91+ }
92+ }
93+
94+ fn subscribe_events ( & self ) -> broadcast:: Receiver < peat_mesh:: broker:: state:: MeshEvent > {
95+ self . event_tx . subscribe ( )
96+ }
97+ }
98+
2999fn main ( ) -> anyhow:: Result < ( ) > {
30100 // Install rustls crypto provider (required by kube's rustls-tls)
31101 rustls:: crypto:: ring:: default_provider ( )
@@ -150,10 +220,30 @@ async fn run() -> anyhow::Result<()> {
150220 let mut discovery: Box < dyn peat_mesh:: discovery:: DiscoveryStrategy > =
151221 match discovery_mode. as_str ( ) {
152222 "kubernetes" | "k8s" => {
153- info ! ( "Using Kubernetes EndpointSlice discovery" ) ;
154- Box :: new ( KubernetesDiscovery :: new (
155- KubernetesDiscoveryConfig :: default ( ) ,
156- ) )
223+ let k8s_namespace = std:: env:: var ( "PEAT_K8S_NAMESPACE" ) . ok ( ) ;
224+ let k8s_label_selector = std:: env:: var ( "PEAT_K8S_LABEL_SELECTOR" )
225+ . unwrap_or_else ( |_| KubernetesDiscoveryConfig :: default ( ) . label_selector ) ;
226+ let k8s_annotation_prefix = std:: env:: var ( "PEAT_K8S_ANNOTATION_PREFIX" )
227+ . unwrap_or_else ( |_| KubernetesDiscoveryConfig :: default ( ) . annotation_prefix ) ;
228+ let k8s_poll_interval = std:: env:: var ( "PEAT_K8S_POLL_INTERVAL_SECS" )
229+ . ok ( )
230+ . and_then ( |v| v. parse :: < u64 > ( ) . ok ( ) )
231+ . map ( std:: time:: Duration :: from_secs)
232+ . unwrap_or_else ( || KubernetesDiscoveryConfig :: default ( ) . poll_interval ) ;
233+
234+ info ! (
235+ namespace = ?k8s_namespace,
236+ label_selector = %k8s_label_selector,
237+ annotation_prefix = %k8s_annotation_prefix,
238+ poll_interval_secs = k8s_poll_interval. as_secs( ) ,
239+ "Using Kubernetes EndpointSlice discovery"
240+ ) ;
241+ Box :: new ( KubernetesDiscovery :: new ( KubernetesDiscoveryConfig {
242+ namespace : k8s_namespace,
243+ label_selector : k8s_label_selector,
244+ annotation_prefix : k8s_annotation_prefix,
245+ poll_interval : k8s_poll_interval,
246+ } ) )
157247 }
158248 "mdns" => {
159249 info ! ( "Using mDNS discovery" ) ;
@@ -465,12 +555,14 @@ async fn run() -> anyhow::Result<()> {
465555 ) ;
466556
467557 // ── Build mesh ───────────────────────────────────────────────
468- let mesh = PeatMeshBuilder :: new ( mesh_config)
469- . with_device_keypair_from_seed ( & seed, & hostname)
470- . map_err ( |e| anyhow:: anyhow!( "Keypair derivation failed: {}" , e) ) ?
471- . with_formation_key ( formation_key)
472- . with_discovery ( discovery)
473- . build ( ) ;
558+ let mesh = Arc :: new (
559+ PeatMeshBuilder :: new ( mesh_config)
560+ . with_device_keypair_from_seed ( & seed, & hostname)
561+ . map_err ( |e| anyhow:: anyhow!( "Keypair derivation failed: {}" , e) ) ?
562+ . with_formation_key ( formation_key)
563+ . with_discovery ( discovery)
564+ . build ( ) ,
565+ ) ;
474566
475567 mesh. start ( )
476568 . map_err ( |e| anyhow:: anyhow!( "Failed to start mesh: {}" , e) ) ?;
@@ -506,6 +598,7 @@ async fn run() -> anyhow::Result<()> {
506598 let coordinator = coordinator. clone ( ) ;
507599 let transport = sync_transport. clone ( ) ;
508600 let ttl_for_sync = ttl_manager. clone ( ) ;
601+ let formation_peers_for_sync = formation_peers. clone ( ) ;
509602 tokio:: spawn ( async move {
510603 let mut interval = tokio:: time:: interval ( std:: time:: Duration :: from_secs ( 5 ) ) ;
511604 interval. set_missed_tick_behavior ( tokio:: time:: MissedTickBehavior :: Skip ) ;
@@ -517,6 +610,32 @@ async fn run() -> anyhow::Result<()> {
517610 break ;
518611 }
519612 }
613+ // Bootstrap sync connections from discovered formation members.
614+ for peer_id in formation_peers_for_sync. snapshot ( ) {
615+ if peer_id == transport. endpoint ( ) . id ( ) {
616+ continue ;
617+ }
618+ if transport. get_connection ( & peer_id) . is_some ( ) {
619+ continue ;
620+ }
621+ match transport. connect_and_authenticate ( peer_id) . await {
622+ Ok ( conn) => {
623+ transport. start_sync_connection ( conn, coordinator. clone ( ) ) ;
624+ info ! (
625+ peer = %peer_id. fmt_short( ) ,
626+ "Bootstrapped sync connection to discovered formation peer"
627+ ) ;
628+ }
629+ Err ( e) => {
630+ warn ! (
631+ peer = %peer_id. fmt_short( ) ,
632+ error = %e,
633+ "Failed to bootstrap sync connection to discovered formation peer"
634+ ) ;
635+ }
636+ }
637+ }
638+
520639 let peers = transport. connected_peers ( ) ;
521640 // When offline (no peers), extend TTLs to prevent premature eviction
522641 if peers. is_empty ( ) {
@@ -628,14 +747,21 @@ async fn run() -> anyhow::Result<()> {
628747 } ;
629748
630749 // ── Broker HTTP server ───────────────────────────────────────
631- let mesh = Arc :: new ( mesh) ;
632750 let broker_config = BrokerConfig {
633751 bind_addr : SocketAddr :: from ( ( [ 0 , 0 , 0 , 0 ] , broker_port) ) ,
634752 ..Default :: default ( )
635753 } ;
636754 let store_adapter = peat_mesh:: broker:: StoreBrokerAdapter :: new ( automerge_store. clone ( ) ) ;
755+ let ( broker_event_tx, _) = broadcast:: channel ( 256 ) ;
756+ let node_state = Arc :: new ( NodeBrokerState {
757+ node_id : hostname. clone ( ) ,
758+ version : env ! ( "CARGO_PKG_VERSION" ) . to_string ( ) ,
759+ started_at : Instant :: now ( ) ,
760+ transport : sync_transport. clone ( ) ,
761+ event_tx : broker_event_tx,
762+ } ) ;
637763 let composite_state = Arc :: new ( peat_mesh:: broker:: CompositeBrokerState :: new (
638- mesh . clone ( ) as Arc < dyn peat_mesh:: broker:: state:: MeshBrokerState > ,
764+ node_state as Arc < dyn peat_mesh:: broker:: state:: MeshBrokerState > ,
639765 store_adapter,
640766 ) ) ;
641767
0 commit comments