Skip to main content

hoike_gossip/
broadcast.rs

1use serde::{Deserialize, Serialize};
2use std::fmt;
3
4/// Gossip message payload carried as a foca custom broadcast.
5///
6/// These messages announce the *existence* of signed artifacts — they never
7/// carry certificate status themselves. A node acts on a gossip message only
8/// by fetching and validating a sealed bundle.
9#[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        /// Name of the node that produced this announcement, so receivers can
18        /// attribute the epoch to a specific fleet member and compute
19        /// per-node staleness. Added after the initial wire format; `#[serde(default)]`
20        /// keeps it backward-compatible — messages from older nodes decode with an
21        /// empty origin (the JSON payload simply omits the field).
22        #[serde(default)]
23        origin_node: String,
24    },
25    UrgentRevocation {
26        producer_id: String,
27        issuer_key_hash: Vec<u8>,
28        epoch: u64,
29        /// See `GenerationAnnouncement::origin_node`.
30        #[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    /// The announcing node's name, or `""` if the message came from an older
74    /// node that predates the `origin_node` field.
75    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/// Broadcast key for foca's deduplication. Two broadcasts with the same
115/// (producer_id, issuer_key_hash) scope where the newer one has a higher
116/// epoch invalidate the older one.
117#[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
132/// BroadcastHandler that processes GossipMessage broadcasts.
133pub struct HoikeBroadcastHandler {
134    tx: tokio::sync::mpsc::Sender<GossipMessage>,
135    /// Inbound message authentication (FPT_ITT.1). `None` disables verification
136    /// entirely — today's unauthenticated behavior — used when no gossip
137    /// `identity_key`/`peer_keys` are configured.
138    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        // Authenticate before decoding. A dropped message returns `Ok(None)` so
179        // foca neither stores nor re-broadcasts it — the forgery dies at this hop.
180        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            // No verifier configured: this node has opted out of authentication.
189            // It must still strip a signing peer's frame to reach the inner JSON,
190            // or a signed broadcast would fail to decode and die at this hop —
191            // breaking mixed-fleet interop during a rolling upgrade. Accepting
192            // the unverified payload is no weaker than the unsigned messages this
193            // node already accepts.
194            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    // A node running the pre-`origin_node` wire format emits JSON without that
283    // field. `#[serde(default)]` must let a new node decode it (empty origin)
284    // rather than rejecting the message — otherwise a mixed fleet partitions
285    // during a rolling upgrade.
286    #[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    // The authentication integration point: a verifier-equipped handler must
303    // admit a validly signed broadcast (onto the channel + re-broadcast) and
304    // silently drop a forged one — foca sees `Ok(None)` so nothing propagates.
305    #[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        // Valid signature → accepted (returns a broadcast key, lands on channel).
330        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        // Forged by an untrusted key → dropped (Ok(None), nothing on channel).
338        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        // Unsigned legacy under Required policy → dropped too.
350        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    // A node with no gossip keys (verifier: None) has opted out of auth, but a
361    // *signing* peer still frames its broadcasts. The keyless node must strip the
362    // frame and decode the inner payload — otherwise every signed announcement
363    // dies at this hop and a mid-upgrade mixed fleet partitions.
364    #[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        // Signed frame from a peer → keyless node strips the frame and decodes.
384        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        // Plain unsigned payload still decodes on the same node.
393        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)); // same epoch doesn't invalidate
423    }
424}