Skip to main content

myko_server/
peer_registry.rs

1//! Peer registry for the cell-based server.
2//!
3//! Convergence model:
4//! - Desired peers are derived from `GetAllServers` (minus self).
5//! - Actual peers are active RAII handles keyed by `server_id`.
6//! - Reconcile only performs set-difference operations:
7//!   - create missing handles
8//!   - drop extra handles
9//! - Handles self-terminate on disconnect/id-mismatch and are re-created only
10//!   by future snapshots.
11
12use 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/// Status of a peer connection.
28#[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
37/// Peer registry for managing connections to other servers.
38pub 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 another row points at our local endpoint but has a different id,
110            // it is a stale incarnation and should be tombstoned.
111            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}