1use std::collections::HashMap;
2use std::net::SocketAddr;
3use std::sync::Arc;
4use std::time::{Duration, SystemTime, UNIX_EPOCH};
5
6use foca::{AccumulatingRuntime, Config as FocaConfig, Foca, PostcardCodec, Timer};
7use rand::SeedableRng;
8use rand::rngs::SmallRng;
9use serde::{Deserialize, Serialize};
10use tokio::net::UdpSocket;
11use tokio::sync::{Mutex, RwLock};
12use tracing::{debug, info, warn};
13
14use crate::broadcast::{GossipMessage, HoikeBroadcastHandler};
15use crate::config::GossipConfig;
16use crate::crypto::{self, GossipSigner, GossipVerifier, VerifyPolicy};
17
18fn now_unix() -> u64 {
20 SystemTime::now()
21 .duration_since(UNIX_EPOCH)
22 .map(|d| d.as_secs())
23 .unwrap_or(0)
24}
25
26#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
29#[serde(rename_all = "lowercase")]
30pub enum MemberState {
31 Alive,
32 Suspect,
33 Down,
34}
35
36impl From<foca::State> for MemberState {
37 fn from(s: foca::State) -> Self {
38 match s {
39 foca::State::Alive => MemberState::Alive,
40 foca::State::Suspect => MemberState::Suspect,
41 foca::State::Down => MemberState::Down,
42 }
43 }
44}
45
46#[derive(Clone, Debug, Serialize)]
48pub struct MemberInfo {
49 pub name: String,
50 pub addr: SocketAddr,
51 pub incarnation: u64,
54 pub state: MemberState,
55 pub is_self: bool,
57}
58
59#[derive(Clone, Debug, Serialize)]
62pub struct GenRecord {
63 pub origin_node: String,
65 pub producer_id: String,
66 pub issuer_key_hash: String,
68 pub epoch: u64,
69 pub manifest_digest: String,
70 pub last_seen_unix: u64,
72}
73
74type GenKey = (String, String, String);
76type GenerationTable = Arc<RwLock<HashMap<GenKey, GenRecord>>>;
77
78#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
81pub struct NodeId {
82 pub addr: SocketAddr,
83 pub name: String,
84 pub incarnation: u64,
85}
86
87impl foca::Identity for NodeId {
88 type Addr = SocketAddr;
89
90 fn addr(&self) -> SocketAddr {
91 self.addr
92 }
93
94 fn renew(&self) -> Option<Self> {
95 Some(NodeId {
96 addr: self.addr,
97 name: self.name.clone(),
98 incarnation: self.incarnation + 1,
99 })
100 }
101
102 fn win_addr_conflict(&self, adversary: &Self) -> bool {
103 self.incarnation > adversary.incarnation
104 }
105}
106
107type HoikeFoca = Foca<NodeId, PostcardCodec, SmallRng, HoikeBroadcastHandler>;
108type TimerQueue = Arc<Mutex<Vec<(Duration, Timer<NodeId>)>>>;
109
110pub struct GossipNode {
111 foca: Arc<Mutex<HoikeFoca>>,
112 #[allow(dead_code)]
113 socket: Arc<UdpSocket>,
114 config: GossipConfig,
115 identity: NodeId,
118 generations: GenerationTable,
120 signer: Option<GossipSigner>,
123}
124
125impl GossipNode {
126 pub async fn start(
128 config: GossipConfig,
129 msg_tx: tokio::sync::mpsc::Sender<GossipMessage>,
130 ) -> Result<Self, Box<dyn std::error::Error + Send + Sync>> {
131 let bind_addr: SocketAddr = config
132 .bind
133 .parse()
134 .map_err(|e| format!("invalid gossip bind address '{}': {}", config.bind, e))?;
135
136 let socket = UdpSocket::bind(bind_addr)
137 .await
138 .map_err(|e| format!("failed to bind gossip socket on {}: {}", bind_addr, e))?;
139
140 let local_addr = socket.local_addr()?;
141 info!(addr = %local_addr, name = %config.node_name, "gossip node starting");
142
143 let identity = NodeId {
144 addr: local_addr,
145 name: config.node_name.clone(),
146 incarnation: 0,
147 };
148
149 let foca_config = FocaConfig::simple();
150 let rng = SmallRng::seed_from_u64(
151 std::time::SystemTime::now()
152 .duration_since(std::time::UNIX_EPOCH)
153 .unwrap_or_default()
154 .as_nanos() as u64,
155 );
156 let signer = match &config.identity_key {
161 Some(path) => {
162 let key = crypto::load_signing_key(path)
163 .map_err(|e| format!("gossip identity_key {}: {e}", path.display()))?;
164 info!(key = %path.display(), "gossip message signing enabled");
165 Some(GossipSigner::new(key))
166 }
167 None => None,
168 };
169
170 if !config.peer_keys.is_empty() {
171 return Err("gossip.peer_keys is no longer safe: migrate to gossip.peer_identities (node name -> public key path)".into());
172 }
173 let mut identities = Vec::new();
174 for (name, path) in &config.peer_identities {
175 if name.is_empty() || name.len() > 256 {
176 return Err("gossip peer identity must contain 1..=256 bytes".into());
177 }
178 identities.push((name.clone(), crypto::load_verifying_key(path)?));
179 }
180 if let Some(s) = &signer {
181 identities.push((config.node_name.clone(), s.verifying_key()));
182 }
183 let verifier = if signer.is_some() || !identities.is_empty() {
184 let policy = if config.peer_identities.is_empty() {
185 VerifyPolicy::Permissive
186 } else {
187 VerifyPolicy::Required
188 };
189 Some(GossipVerifier::for_identities(identities, policy))
190 } else {
191 None
192 };
193
194 let broadcast_handler = HoikeBroadcastHandler::new(msg_tx, verifier);
195
196 let self_identity = identity.clone();
200
201 let foca = Foca::with_custom_broadcast(
202 identity,
203 foca_config,
204 rng,
205 PostcardCodec,
206 broadcast_handler,
207 );
208
209 let socket = Arc::new(socket);
210 let foca = Arc::new(Mutex::new(foca));
211 let timer_queue: TimerQueue = Arc::new(Mutex::new(Vec::new()));
212
213 let node = GossipNode {
214 foca: Arc::clone(&foca),
215 socket: Arc::clone(&socket),
216 config: config.clone(),
217 identity: self_identity,
218 generations: Arc::new(RwLock::new(HashMap::new())),
219 signer,
220 };
221
222 {
224 let foca = Arc::clone(&foca);
225 let socket = Arc::clone(&socket);
226 let tq = Arc::clone(&timer_queue);
227 tokio::spawn(async move {
228 receive_loop(foca, socket, tq).await;
229 });
230 }
231
232 {
234 let foca = Arc::clone(&foca);
235 let socket_clone = Arc::clone(&socket);
236 let tq = Arc::clone(&timer_queue);
237 tokio::spawn(async move {
238 timer_loop(foca, socket_clone, tq).await;
239 });
240 }
241
242 for seed in &config.seeds {
244 match seed.parse::<SocketAddr>() {
245 Ok(addr) => {
246 let seed_id = NodeId {
247 addr,
248 name: String::new(),
249 incarnation: 0,
250 };
251 let mut foca_guard = foca.lock().await;
252 let mut runtime = AccumulatingRuntime::new();
253 if let Err(e) = foca_guard.announce(seed_id, &mut runtime) {
254 warn!(seed = %addr, error = %e, "failed to announce to seed");
255 }
256 drain_runtime(&mut runtime, &socket, &timer_queue).await;
257 info!(seed = %addr, "announced to seed");
258 }
259 Err(e) => {
260 warn!(seed = %seed, error = %e, "invalid seed address, skipping");
261 }
262 }
263 }
264
265 Ok(node)
266 }
267
268 fn frame_broadcast(&self, payload: Vec<u8>) -> Vec<u8> {
272 match &self.signer {
273 Some(signer) => signer.frame(&payload),
274 None => payload,
275 }
276 }
277
278 pub async fn announce_generation(
280 &self,
281 producer_id: String,
282 issuer_key_hash: Vec<u8>,
283 epoch: u64,
284 manifest_digest: [u8; 32],
285 bundle_url: Option<String>,
286 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
287 let msg = GossipMessage::GenerationAnnouncement {
288 producer_id,
289 issuer_key_hash,
290 epoch,
291 manifest_digest,
292 bundle_url,
293 origin_node: self.identity.name.clone(),
294 };
295
296 let data = self.frame_broadcast(serde_json::to_vec(&msg)?);
297 let mut foca = self.foca.lock().await;
298 foca.add_broadcast(&data)?;
299
300 info!(msg = %msg, "broadcasting generation announcement");
301 Ok(())
302 }
303
304 pub async fn announce_urgent_revocation(
306 &self,
307 producer_id: String,
308 issuer_key_hash: Vec<u8>,
309 epoch: u64,
310 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
311 let msg = GossipMessage::UrgentRevocation {
312 producer_id,
313 issuer_key_hash,
314 epoch,
315 origin_node: self.identity.name.clone(),
316 };
317
318 let data = self.frame_broadcast(serde_json::to_vec(&msg)?);
319 let mut foca = self.foca.lock().await;
320 foca.add_broadcast(&data)?;
321
322 info!(msg = %msg, "broadcasting urgent revocation");
323 Ok(())
324 }
325
326 pub fn config(&self) -> &GossipConfig {
327 &self.config
328 }
329
330 pub fn identity(&self) -> &NodeId {
332 &self.identity
333 }
334
335 pub async fn members(&self) -> Vec<MemberInfo> {
341 let foca = self.foca.lock().await;
342 let mut out: Vec<MemberInfo> = foca
343 .iter_membership_state()
344 .map(|m| {
345 let id = m.id();
346 MemberInfo {
347 name: id.name.clone(),
348 addr: id.addr,
349 incarnation: id.incarnation,
350 state: m.state().into(),
351 is_self: false,
352 }
353 })
354 .collect();
355 drop(foca);
356
357 out.push(MemberInfo {
358 name: self.identity.name.clone(),
359 addr: self.identity.addr,
360 incarnation: self.identity.incarnation,
361 state: MemberState::Alive,
362 is_self: true,
363 });
364 out
365 }
366
367 pub async fn record_generation(&self, msg: &GossipMessage) {
375 let GossipMessage::GenerationAnnouncement {
376 producer_id,
377 issuer_key_hash,
378 epoch,
379 manifest_digest,
380 origin_node,
381 ..
382 } = msg
383 else {
384 return;
385 };
386
387 if !msg.valid_bounds() {
388 return;
389 }
390 let ikh_hex = hex::encode(issuer_key_hash);
391 let key: GenKey = (origin_node.clone(), producer_id.clone(), ikh_hex.clone());
392 let now = now_unix();
393
394 let mut table = self.generations.write().await;
395 table.retain(|_, row| now.saturating_sub(row.last_seen_unix) < 3600);
396 if !table.contains_key(&key) && table.len() >= 4096 {
397 return;
398 }
399 let entry = table.entry(key).or_insert_with(|| GenRecord {
400 origin_node: origin_node.clone(),
401 producer_id: producer_id.clone(),
402 issuer_key_hash: ikh_hex,
403 epoch: *epoch,
404 manifest_digest: hex::encode(manifest_digest),
405 last_seen_unix: now,
406 });
407 if *epoch >= entry.epoch {
408 entry.epoch = *epoch;
409 entry.manifest_digest = hex::encode(manifest_digest);
410 }
411 entry.last_seen_unix = now;
412 }
413
414 pub async fn generations(&self) -> Vec<GenRecord> {
416 let mut table = self.generations.write().await;
417 let now = now_unix();
418 table.retain(|_, row| now.saturating_sub(row.last_seen_unix) < 3600);
419 table.values().cloned().collect()
420 }
421}
422
423async fn receive_loop(
425 foca: Arc<Mutex<HoikeFoca>>,
426 socket: Arc<UdpSocket>,
427 timer_queue: TimerQueue,
428) {
429 let mut buf = vec![0u8; 2048];
430 loop {
431 match socket.recv_from(&mut buf).await {
432 Ok((len, from)) => {
433 debug!(from = %from, len, "gossip packet received");
434 let mut foca_guard = foca.lock().await;
435 let mut runtime = AccumulatingRuntime::new();
436 if let Err(e) = foca_guard.handle_data(&buf[..len], &mut runtime) {
437 debug!(from = %from, error = %e, "foca handle_data error (expected for arbitrary UDP)");
438 }
439 drop(foca_guard);
440 drain_runtime(&mut runtime, &socket, &timer_queue).await;
441 }
442 Err(e) => {
443 warn!(error = %e, "UDP recv error");
444 tokio::time::sleep(Duration::from_millis(100)).await;
445 }
446 }
447 }
448}
449
450async fn timer_loop(
452 foca: Arc<Mutex<HoikeFoca>>,
453 socket: Arc<UdpSocket>,
454 new_timers_rx: TimerQueue,
455) {
456 let mut pending_timers: Vec<(tokio::time::Instant, Timer<NodeId>)> = Vec::new();
457
458 let tick = Duration::from_millis(200);
459 let mut interval = tokio::time::interval(tick);
460
461 loop {
462 interval.tick().await;
463
464 {
466 let mut incoming = new_timers_rx.lock().await;
467 for (duration, timer) in incoming.drain(..) {
468 pending_timers.push((tokio::time::Instant::now() + duration, timer));
469 }
470 }
471
472 let now = tokio::time::Instant::now();
473
474 let mut due = Vec::new();
475 pending_timers.retain(|(deadline, timer)| {
476 if *deadline <= now {
477 due.push(timer.clone());
478 false
479 } else {
480 true
481 }
482 });
483
484 if due.is_empty() {
485 continue;
486 }
487
488 let mut foca_guard = foca.lock().await;
489 let mut runtime = AccumulatingRuntime::new();
490
491 for timer in due {
492 if let Err(e) = foca_guard.handle_timer(timer, &mut runtime) {
493 debug!(error = %e, "foca timer error");
494 }
495 }
496
497 while let Some(notification) = runtime.to_notify() {
498 handle_notification(¬ification);
499 }
500
501 while let Some((duration, timer)) = runtime.to_schedule() {
502 pending_timers.push((tokio::time::Instant::now() + duration, timer));
503 }
504
505 drop(foca_guard);
506
507 while let Some((to, data)) = runtime.to_send() {
509 if let Err(e) = socket.send_to(&data, to.addr).await {
510 debug!(to = %to.addr, error = %e, "failed to send gossip packet");
511 }
512 }
513 }
514}
515
516fn handle_notification(notification: &foca::OwnedNotification<NodeId>) {
517 match notification {
518 foca::OwnedNotification::MemberUp(id) => {
519 info!(name = %id.name, addr = %id.addr, "member joined");
520 }
521 foca::OwnedNotification::MemberDown(id) => {
522 info!(name = %id.name, addr = %id.addr, "member left");
523 }
524 foca::OwnedNotification::Rename(before, after) => {
525 info!(
526 before_name = %before.name, before_addr = %before.addr,
527 after_name = %after.name, after_addr = %after.addr,
528 "member renamed (rejoin)"
529 );
530 }
531 foca::OwnedNotification::Active => {
532 info!("gossip node is active (known by cluster)");
533 }
534 foca::OwnedNotification::Idle => {
535 info!("gossip node is idle (no known peers)");
536 }
537 foca::OwnedNotification::Defunct => {
538 warn!("gossip node declared defunct — needs manual intervention");
539 }
540 foca::OwnedNotification::Rejoin(id) => {
541 info!(name = %id.name, incarnation = id.incarnation, "auto-rejoined cluster");
542 }
543 }
544}
545
546async fn drain_runtime(
548 runtime: &mut AccumulatingRuntime<NodeId>,
549 socket: &UdpSocket,
550 timer_queue: &TimerQueue,
551) {
552 while let Some(notification) = runtime.to_notify() {
553 handle_notification(¬ification);
554 }
555
556 while let Some((to, data)) = runtime.to_send() {
557 let addr = to.addr;
558 if let Err(e) = socket.send_to(&data, addr).await {
559 debug!(to = %addr, error = %e, "failed to send gossip packet");
560 }
561 }
562
563 let mut timers = timer_queue.lock().await;
565 while let Some((duration, timer)) = runtime.to_schedule() {
566 timers.push((duration, timer));
567 }
568}
569
570#[cfg(test)]
571mod tests {
572 use super::*;
573
574 #[tokio::test]
575 async fn generation_storage_is_bounded_and_expires() {
576 let node = test_node("edge").await;
577 for n in 0..4100 {
578 node.record_generation(&GossipMessage::GenerationAnnouncement {
579 producer_id: format!("producer-{n}"),
580 issuer_key_hash: vec![1; 32],
581 epoch: 1,
582 manifest_digest: [0; 32],
583 bundle_url: None,
584 origin_node: "peer".into(),
585 })
586 .await;
587 }
588 assert_eq!(node.generations().await.len(), 4096);
589 for row in node.generations.write().await.values_mut() {
590 row.last_seen_unix = 0;
591 }
592 assert!(node.generations().await.is_empty());
593 }
594
595 #[test]
596 fn node_id_identity_trait() {
597 let id = NodeId {
598 addr: "127.0.0.1:7946".parse().unwrap(),
599 name: "test-node".into(),
600 incarnation: 0,
601 };
602
603 use foca::Identity;
605 assert_eq!(id.addr(), "127.0.0.1:7946".parse::<SocketAddr>().unwrap());
606
607 let renewed = id.renew().unwrap();
609 assert_eq!(renewed.incarnation, 1);
610 assert_eq!(renewed.addr, id.addr);
611 assert_eq!(renewed.name, id.name);
612
613 assert!(renewed.win_addr_conflict(&id));
615 assert!(!id.win_addr_conflict(&renewed));
616 }
617
618 #[test]
619 fn node_id_serde_round_trip() {
620 let id = NodeId {
621 addr: "192.168.1.1:7946".parse().unwrap(),
622 name: "edge-node-01".into(),
623 incarnation: 5,
624 };
625
626 let json = serde_json::to_string(&id).unwrap();
627 let decoded: NodeId = serde_json::from_str(&json).unwrap();
628 assert_eq!(id, decoded);
629 }
630
631 async fn test_node(name: &str) -> GossipNode {
633 let (tx, _rx) = tokio::sync::mpsc::channel(16);
634 let config = GossipConfig {
635 enabled: true,
636 bind: "127.0.0.1:0".into(),
637 seeds: vec![],
638 node_name: name.into(),
639 identity_key: None,
640 peer_keys: vec![],
641 peer_identities: Default::default(),
642 };
643 GossipNode::start(config, tx).await.expect("node starts")
644 }
645
646 #[tokio::test]
649 async fn members_always_includes_self() {
650 let node = test_node("solo").await;
651 let members = node.members().await;
652 assert_eq!(members.len(), 1, "no peers, only self");
653 assert!(members[0].is_self);
654 assert_eq!(members[0].name, "solo");
655 assert_eq!(members[0].state, MemberState::Alive);
656 }
657
658 #[tokio::test]
659 async fn record_generation_tracks_latest_epoch_per_scope() {
660 let node = test_node("signer").await;
661 let ikh = vec![0xAB; 32];
662
663 let announce = |epoch: u64| GossipMessage::GenerationAnnouncement {
664 producer_id: "prod-1".into(),
665 issuer_key_hash: ikh.clone(),
666 epoch,
667 manifest_digest: [epoch as u8; 32],
668 bundle_url: None,
669 origin_node: "peer-a".into(),
670 };
671
672 node.record_generation(&announce(5)).await;
673 node.record_generation(&announce(7)).await;
674 node.record_generation(&announce(6)).await;
676
677 let gens = node.generations().await;
678 assert_eq!(gens.len(), 1, "one (node, scope) row");
679 assert_eq!(gens[0].epoch, 7);
680 assert_eq!(gens[0].origin_node, "peer-a");
681 assert_eq!(gens[0].producer_id, "prod-1");
682 assert_eq!(gens[0].issuer_key_hash, hex::encode(&ikh));
683 }
684
685 #[tokio::test]
688 async fn record_generation_separates_scopes_and_nodes() {
689 let node = test_node("edge").await;
690 node.record_generation(&GossipMessage::GenerationAnnouncement {
691 producer_id: "prod-1".into(),
692 issuer_key_hash: vec![0x01; 32],
693 epoch: 3,
694 manifest_digest: [0; 32],
695 bundle_url: None,
696 origin_node: "node-a".into(),
697 })
698 .await;
699 node.record_generation(&GossipMessage::GenerationAnnouncement {
700 producer_id: "prod-1".into(),
701 issuer_key_hash: vec![0x01; 32],
702 epoch: 3,
703 manifest_digest: [0; 32],
704 bundle_url: None,
705 origin_node: "node-b".into(),
706 })
707 .await;
708 assert_eq!(node.generations().await.len(), 2, "two distinct nodes");
709 }
710}