1use std::sync::Arc;
13
14use dashmap::DashMap;
15use hyphae::{
16 Cell, CellImmutable, MaterializeDefinite, Signal, SubscriptionGuard, TapExt, Watchable,
17};
18use myko::{
19 entities::server::{GetAllServers, GetPeerServers, Server, ServerId},
20 server::CellServerCtx,
21};
22use tracing::info;
23
24mod peer_connection_handle;
25use peer_connection_handle::{PeerConnectionHandle, PeerState};
26
27#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
29pub struct PeerStatus {
30 pub peer_id: String,
31 pub is_connected: bool,
32 pub is_alive: bool,
33 pub latency_ms: Option<u64>,
34 pub last_seen: Option<String>,
35}
36
37pub struct PeerRegistry {
39 _peers_guard: SubscriptionGuard,
40 _self_advertise_guard: SubscriptionGuard,
41 _connections: Arc<DashMap<ServerId, PeerConnectionHandle>>,
42 _remove_guards: Arc<DashMap<ServerId, SubscriptionGuard>>,
43}
44
45#[derive(Debug, Clone)]
46pub struct PeerRegistryConfig {
47 pub address: String,
48 pub port: u16,
49 pub version: String,
50}
51
52impl PeerRegistry {
53 fn build_local_server(config: &PeerRegistryConfig, host_id: &ServerId) -> Server {
54 Server {
55 address: config.address.clone(),
56 id: host_id.clone(),
57 port: config.port,
58 started_at: chrono::Utc::now().to_rfc3339(),
59 version: config.version.clone(),
60 }
61 }
62
63 fn spawn_self_advertise_guard(
64 ctx: CellServerCtx,
65 local_server: Server,
66 self_host_id: ServerId,
67 ) -> SubscriptionGuard {
68 let all_servers = ctx
69 .query_map(GetAllServers {}, ctx.new_server_transaction())
70 .items();
71 all_servers.subscribe(move |signal| {
72 if let Signal::Value(servers) = signal {
73 let has_self = servers.iter().any(|s| s.id == self_host_id);
74 if !has_self {
75 tracing::info!(
76 "Local server {} missing from GetAllServers; re-advertising",
77 self_host_id
78 );
79 if let Err(e) = ctx.set(&local_server) {
80 tracing::error!("Failed to re-advertise local server: {e}");
81 }
82 }
83 }
84 })
85 }
86
87 fn reconcile_peer_snapshot<T>(
88 peers: &T,
89 host_id: &ServerId,
90 local_address: &str,
91 local_port: u16,
92 ctx: &CellServerCtx,
93 connections: &Arc<DashMap<ServerId, PeerConnectionHandle>>,
94 remove_guards: &Arc<DashMap<ServerId, SubscriptionGuard>>,
95 ) where
96 T: AsRef<[Arc<Server>]>,
97 {
98 tracing::info!(
99 "Current Peers: {}",
100 peers
101 .as_ref()
102 .iter()
103 .map(|s| format!("{}/{}:{}", s.id, s.address, s.port))
104 .collect::<Vec<String>>()
105 .join(", ")
106 );
107
108 for server in peers.as_ref() {
109 if server.id != *host_id && server.address == local_address && server.port == local_port
112 {
113 tracing::warn!(
114 "Deleting stale entries: {}:{}:{}",
115 server.id,
116 server.address,
117 server.port
118 );
119 ctx.unregister_peer_client(server.id.as_ref());
120 if let Err(e) = ctx.del(server.as_ref()) {
121 tracing::error!("Failed to delete stale server entry: {e}");
122 }
123 continue;
124 }
125
126 if connections.contains_key(&server.id) {
127 continue;
128 }
129
130 let handle = PeerConnectionHandle::new(server.clone());
131
132 let remove_connection_handles = connections.clone();
133 let remove_state_guards = remove_guards.clone();
134
135 let remove_server = server.clone();
136
137 let remove_ctx = ctx.clone();
138
139 let state_guard = handle.signal_state.subscribe(move |state| {
140 if let Signal::Value(v) = state
141 && v.as_ref() == &PeerState::Delete
142 {
143 tracing::warn!(
144 "Deleting: {}:{}:{}",
145 remove_server.id,
146 remove_server.address,
147 remove_server.port
148 );
149 remove_connection_handles.remove(&remove_server.id);
150 remove_state_guards.remove(&remove_server.id);
151 remove_ctx.unregister_peer_client(remove_server.id.as_ref());
152 if let Err(e) = remove_ctx.del(remove_server.as_ref()) {
153 tracing::error!("Failed to delete peer server: {e}");
154 }
155 }
156 });
157
158 let server_id = server.id.clone();
159 ctx.register_peer_client(server_id.clone(), handle.client());
160 connections.insert(server_id.clone(), handle);
161 remove_guards.insert(server_id, state_guard);
162 }
163 }
164
165 fn spawn_peer_reconcile_guard(
166 peer_servers: Cell<Vec<Arc<Server>>, CellImmutable>,
167 host_id: ServerId,
168 local_address: String,
169 local_port: u16,
170 ctx: CellServerCtx,
171 connections: Arc<DashMap<ServerId, PeerConnectionHandle>>,
172 remove_guards: Arc<DashMap<ServerId, SubscriptionGuard>>,
173 ) -> SubscriptionGuard {
174 peer_servers
175 .tap(move |peers| {
176 Self::reconcile_peer_snapshot(
177 peers,
178 &host_id,
179 &local_address,
180 local_port,
181 &ctx,
182 &connections,
183 &remove_guards,
184 );
185 })
186 .materialize()
187 .subscribe(|_| {})
188 }
189
190 pub fn new(ctx: CellServerCtx, config: PeerRegistryConfig) -> Self {
191 let server_req = ctx.new_server_transaction();
192
193 let connections = Arc::new(DashMap::new());
194 let remove_guards = Arc::new(DashMap::new());
195 let peer_servers = ctx.query_map(GetPeerServers {}, server_req).items();
196 let host_id = ServerId(ctx.host_id.to_string().into());
197 let server = Self::build_local_server(&config, &host_id);
198
199 let self_advertise_guard =
200 Self::spawn_self_advertise_guard(ctx.clone(), server.clone(), host_id.clone());
201
202 let peer_sub = Self::spawn_peer_reconcile_guard(
203 peer_servers,
204 host_id.clone(),
205 config.address.clone(),
206 config.port,
207 ctx.clone(),
208 connections.clone(),
209 remove_guards.clone(),
210 );
211
212 tracing::info!(
213 "Publishing local server bootstrap advert: {}:{}:{}",
214 server.id,
215 server.address,
216 server.port
217 );
218 if let Err(e) = ctx.set(&server) {
219 tracing::error!("Failed to publish local server bootstrap advert: {e}");
220 }
221
222 Self {
223 _peers_guard: peer_sub,
224 _self_advertise_guard: self_advertise_guard,
225 _connections: connections,
226 _remove_guards: remove_guards,
227 }
228 }
229
230 pub fn shutdown(&self) {
231 info!("PeerRegistry shutting down");
232 }
233}
234
235impl Drop for PeerRegistry {
236 fn drop(&mut self) {
237 tracing::warn!("Dropping Peer Registry");
238 self.shutdown();
239 }
240}