1use serde::{Deserialize, Serialize};
2use std::fmt;
3
4#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
10pub enum GossipMessage {
11 GenerationAnnouncement {
12 producer_id: String,
13 issuer_key_hash: Vec<u8>,
14 epoch: u64,
15 manifest_digest: [u8; 32],
16 bundle_url: Option<String>,
17 #[serde(default)]
23 origin_node: String,
24 },
25 UrgentRevocation {
26 producer_id: String,
27 issuer_key_hash: Vec<u8>,
28 epoch: u64,
29 #[serde(default)]
31 origin_node: String,
32 },
33}
34
35impl GossipMessage {
36 pub(crate) fn valid_bounds(&self) -> bool {
37 let (producer, hash) = self.scope_key();
38 !producer.is_empty()
39 && producer.len() <= 256
40 && self.origin_node().len() <= 256
41 && matches!(hash.len(), 20 | 32)
42 && match self {
43 Self::GenerationAnnouncement {
44 bundle_url: Some(url),
45 ..
46 } => url.len() <= 1024,
47 _ => true,
48 }
49 }
50
51 pub fn scope_key(&self) -> (&str, &[u8]) {
52 match self {
53 GossipMessage::GenerationAnnouncement {
54 producer_id,
55 issuer_key_hash,
56 ..
57 } => (producer_id, issuer_key_hash),
58 GossipMessage::UrgentRevocation {
59 producer_id,
60 issuer_key_hash,
61 ..
62 } => (producer_id, issuer_key_hash),
63 }
64 }
65
66 pub fn epoch(&self) -> u64 {
67 match self {
68 GossipMessage::GenerationAnnouncement { epoch, .. } => *epoch,
69 GossipMessage::UrgentRevocation { epoch, .. } => *epoch,
70 }
71 }
72
73 pub fn origin_node(&self) -> &str {
76 match self {
77 GossipMessage::GenerationAnnouncement { origin_node, .. } => origin_node,
78 GossipMessage::UrgentRevocation { origin_node, .. } => origin_node,
79 }
80 }
81}
82
83impl fmt::Display for GossipMessage {
84 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
85 match self {
86 GossipMessage::GenerationAnnouncement {
87 producer_id,
88 issuer_key_hash,
89 epoch,
90 ..
91 } => write!(
92 f,
93 "GenerationAnnouncement(producer={}, ikh={}, epoch={})",
94 producer_id,
95 hex::encode(&issuer_key_hash[..8.min(issuer_key_hash.len())]),
96 epoch
97 ),
98 GossipMessage::UrgentRevocation {
99 producer_id,
100 issuer_key_hash,
101 epoch,
102 ..
103 } => write!(
104 f,
105 "UrgentRevocation(producer={}, ikh={}, epoch={})",
106 producer_id,
107 hex::encode(&issuer_key_hash[..8.min(issuer_key_hash.len())]),
108 epoch
109 ),
110 }
111 }
112}
113
114#[derive(Clone, Debug, PartialEq, Eq)]
118pub struct BroadcastKey {
119 pub producer_id: String,
120 pub issuer_key_hash: Vec<u8>,
121 pub epoch: u64,
122}
123
124impl foca::Invalidates for BroadcastKey {
125 fn invalidates(&self, other: &Self) -> bool {
126 self.producer_id == other.producer_id
127 && self.issuer_key_hash == other.issuer_key_hash
128 && self.epoch > other.epoch
129 }
130}
131
132pub struct HoikeBroadcastHandler {
134 tx: tokio::sync::mpsc::Sender<GossipMessage>,
135 verifier: Option<crate::crypto::GossipVerifier>,
139 admitted: std::collections::HashMap<(String, Vec<u8>), std::time::Instant>,
140}
141
142impl HoikeBroadcastHandler {
143 pub fn new(
144 tx: tokio::sync::mpsc::Sender<GossipMessage>,
145 verifier: Option<crate::crypto::GossipVerifier>,
146 ) -> Self {
147 Self {
148 tx,
149 verifier,
150 admitted: Default::default(),
151 }
152 }
153}
154
155#[derive(Debug)]
156pub struct BroadcastError(String);
157
158impl fmt::Display for BroadcastError {
159 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
160 write!(f, "broadcast error: {}", self.0)
161 }
162}
163
164impl std::error::Error for BroadcastError {}
165
166impl<T> foca::BroadcastHandler<T> for HoikeBroadcastHandler {
167 type Key = BroadcastKey;
168 type Error = BroadcastError;
169
170 fn receive_item(
171 &mut self,
172 data: &[u8],
173 _sender: Option<&T>,
174 ) -> Result<Option<Self::Key>, Self::Error> {
175 if data.len() > 2048 {
176 return Ok(None);
177 }
178 let payload: &[u8] = match &self.verifier {
181 Some(v) => match v.check(data) {
182 crate::crypto::VerifyOutcome::Accept(p) => p,
183 crate::crypto::VerifyOutcome::Reject(reason) => {
184 tracing::warn!(reason, "gossip message rejected (authentication)");
185 return Ok(None);
186 }
187 },
188 None => crate::crypto::unwrap_frame(data),
195 };
196
197 let msg: GossipMessage =
198 serde_json::from_slice(payload).map_err(|e| BroadcastError(format!("decode: {e}")))?;
199
200 if !msg.valid_bounds() {
201 return Ok(None);
202 }
203 self.admitted
204 .retain(|_, seen| seen.elapsed() < std::time::Duration::from_secs(3600));
205 let scope = (msg.scope_key().0.to_string(), msg.scope_key().1.to_vec());
206 if !self.admitted.contains_key(&scope) && self.admitted.len() >= 4096 {
207 return Ok(None);
208 }
209 self.admitted.insert(scope, std::time::Instant::now());
210 let key = BroadcastKey {
211 producer_id: msg.scope_key().0.to_string(),
212 issuer_key_hash: msg.scope_key().1.to_vec(),
213 epoch: msg.epoch(),
214 };
215
216 if let Err(e) = self.tx.try_send(msg) {
217 tracing::warn!(error = %e, "gossip message channel full — announcement may be lost");
218 }
219
220 Ok(Some(key))
221 }
222}
223
224#[cfg(test)]
225mod tests {
226 use super::*;
227
228 #[test]
229 fn handler_bounds_new_scopes_and_payloads() {
230 use foca::BroadcastHandler;
231 let (tx, _rx) = tokio::sync::mpsc::channel(8192);
232 let mut handler = HoikeBroadcastHandler::new(tx, None);
233 for n in 0..4100 {
234 let msg = GossipMessage::UrgentRevocation {
235 producer_id: format!("p{n}"),
236 issuer_key_hash: vec![1; 32],
237 epoch: 1,
238 origin_node: "peer".into(),
239 };
240 let accepted = handler
241 .receive_item(&serde_json::to_vec(&msg).unwrap(), None::<&()>)
242 .unwrap()
243 .is_some();
244 assert_eq!(accepted, n < 4096);
245 }
246 assert!(
247 handler
248 .receive_item(&vec![0; 2049], None::<&()>)
249 .unwrap()
250 .is_none()
251 );
252 assert_eq!(handler.admitted.len(), 4096);
253 }
254
255 #[test]
256 fn gossip_message_serialization_round_trip() {
257 let msg = GossipMessage::GenerationAnnouncement {
258 producer_id: "signer-a".into(),
259 issuer_key_hash: vec![0xAA; 32],
260 epoch: 42,
261 manifest_digest: [0xBB; 32],
262 bundle_url: Some("https://signer-a.example/ahu/latest.ahu".into()),
263 origin_node: "edge-1".into(),
264 };
265
266 let json = serde_json::to_vec(&msg).unwrap();
267 let decoded: GossipMessage = serde_json::from_slice(&json).unwrap();
268 assert_eq!(msg, decoded);
269 assert_eq!(decoded.origin_node(), "edge-1");
270
271 let msg2 = GossipMessage::UrgentRevocation {
272 producer_id: "signer-b".into(),
273 issuer_key_hash: vec![0xCC; 32],
274 epoch: 100,
275 origin_node: "signer-b".into(),
276 };
277 let json2 = serde_json::to_vec(&msg2).unwrap();
278 let decoded2: GossipMessage = serde_json::from_slice(&json2).unwrap();
279 assert_eq!(msg2, decoded2);
280 }
281
282 #[test]
287 fn legacy_announcement_without_origin_node_decodes() {
288 let legacy = serde_json::json!({
289 "GenerationAnnouncement": {
290 "producer_id": "signer-a",
291 "issuer_key_hash": [1, 2, 3],
292 "epoch": 7,
293 "manifest_digest": vec![0u8; 32],
294 "bundle_url": null
295 }
296 });
297 let decoded: GossipMessage = serde_json::from_value(legacy).unwrap();
298 assert_eq!(decoded.epoch(), 7);
299 assert_eq!(decoded.origin_node(), "", "missing origin decodes to empty");
300 }
301
302 #[test]
306 fn handler_drops_forged_and_admits_signed() {
307 use crate::crypto::{GossipSigner, GossipVerifier, VerifyPolicy};
308 use ed25519_dalek::SigningKey;
309 use foca::BroadcastHandler;
310
311 let trusted_key = SigningKey::from_bytes(&[1u8; 32]);
312 let signer = GossipSigner::new(trusted_key.clone());
313 let verifier =
314 GossipVerifier::new(vec![trusted_key.verifying_key()], VerifyPolicy::Required);
315
316 let (tx, mut rx) = tokio::sync::mpsc::channel::<GossipMessage>(4);
317 let mut handler = HoikeBroadcastHandler::new(tx, Some(verifier));
318
319 let msg = GossipMessage::GenerationAnnouncement {
320 producer_id: "signer-a".into(),
321 issuer_key_hash: vec![0xAA; 32],
322 epoch: 9,
323 manifest_digest: [0u8; 32],
324 bundle_url: None,
325 origin_node: "node-a".into(),
326 };
327 let payload = serde_json::to_vec(&msg).unwrap();
328
329 let good = signer.frame(&payload);
331 let key = handler
332 .receive_item(&good, None::<&()>)
333 .expect("no decode error");
334 assert!(key.is_some(), "valid signed message must be accepted");
335 assert_eq!(rx.try_recv().unwrap().epoch(), 9);
336
337 let attacker = GossipSigner::new(SigningKey::from_bytes(&[9u8; 32]));
339 let forged = attacker.frame(&payload);
340 let key = handler
341 .receive_item(&forged, None::<&()>)
342 .expect("drop is not a decode error");
343 assert!(key.is_none(), "forged message must be dropped");
344 assert!(
345 rx.try_recv().is_err(),
346 "nothing forwarded for forged message"
347 );
348
349 let key = handler
351 .receive_item(&payload, None::<&()>)
352 .expect("drop is not a decode error");
353 assert!(
354 key.is_none(),
355 "unsigned message dropped under Required policy"
356 );
357 assert!(rx.try_recv().is_err());
358 }
359
360 #[test]
365 fn none_verifier_decodes_signed_frame_from_peer() {
366 use crate::crypto::GossipSigner;
367 use ed25519_dalek::SigningKey;
368 use foca::BroadcastHandler;
369
370 let (tx, mut rx) = tokio::sync::mpsc::channel::<GossipMessage>(4);
371 let mut handler = HoikeBroadcastHandler::new(tx, None);
372
373 let msg = GossipMessage::GenerationAnnouncement {
374 producer_id: "signer-a".into(),
375 issuer_key_hash: vec![0xAA; 32],
376 epoch: 11,
377 manifest_digest: [0u8; 32],
378 bundle_url: None,
379 origin_node: "node-a".into(),
380 };
381 let payload = serde_json::to_vec(&msg).unwrap();
382
383 let signer = GossipSigner::new(SigningKey::from_bytes(&[7u8; 32]));
385 let framed = signer.frame(&payload);
386 let key = handler
387 .receive_item(&framed, None::<&()>)
388 .expect("keyless node must decode a signed frame, not error");
389 assert!(key.is_some(), "signed frame accepted by keyless node");
390 assert_eq!(rx.try_recv().unwrap().epoch(), 11);
391
392 let key = handler
394 .receive_item(&payload, None::<&()>)
395 .expect("keyless node decodes unsigned too");
396 assert!(key.is_some());
397 assert_eq!(rx.try_recv().unwrap().epoch(), 11);
398 }
399
400 #[test]
401 fn broadcast_key_invalidation() {
402 let old = BroadcastKey {
403 producer_id: "prod-1".into(),
404 issuer_key_hash: vec![0x01; 32],
405 epoch: 5,
406 };
407 let new = BroadcastKey {
408 producer_id: "prod-1".into(),
409 issuer_key_hash: vec![0x01; 32],
410 epoch: 10,
411 };
412 let other_producer = BroadcastKey {
413 producer_id: "prod-2".into(),
414 issuer_key_hash: vec![0x01; 32],
415 epoch: 10,
416 };
417
418 use foca::Invalidates;
419 assert!(new.invalidates(&old));
420 assert!(!old.invalidates(&new));
421 assert!(!new.invalidates(&other_producer));
422 assert!(!old.invalidates(&old)); }
424}