1use bytes::Buf;
23use commonware_codec::{EncodeSize, Error as CodecError, RangeCfg, Read, Write};
24use commonware_consensus::types::Epoch;
25use commonware_cryptography::PublicKey;
26use commonware_p2p::{
27 Address, AddressableManager as P2pAddressableManager, AddressableTrackedPeers,
28 Manager as P2pManager, Provider, TrackedPeers,
29};
30use commonware_utils::{
31 ordered::{Map, Set},
32 sequence::Unit,
33};
34use std::{convert::Infallible, fmt, fmt::Debug};
35use thiserror::Error;
36
37pub trait Directory<P: PublicKey>:
50 Clone + Debug + PartialEq + Eq + Send + Sync + 'static + Read + Write + EncodeSize
51{
52 fn codec_config(peers: &Set<P>) -> Self::Cfg;
54
55 fn matches(&self, peers: &Set<P>) -> bool;
57}
58
59impl<P: PublicKey> Directory<P> for Unit {
61 fn codec_config(_: &Set<P>) -> Self::Cfg {}
62
63 fn matches(&self, _: &Set<P>) -> bool {
64 true
65 }
66}
67
68#[derive(Clone, Debug, PartialEq, Eq)]
70pub struct Addresses<P: PublicKey>(Map<P, Address>);
71
72impl<P: PublicKey> Addresses<P> {
73 pub fn get(&self, peer: &P) -> Option<&Address> {
75 self.0.get_value(peer)
76 }
77
78 pub fn into_inner(self) -> Map<P, Address> {
80 self.0
81 }
82}
83
84impl<P: PublicKey> From<Map<P, Address>> for Addresses<P> {
85 fn from(addresses: Map<P, Address>) -> Self {
86 Self(addresses)
87 }
88}
89
90impl<P: PublicKey> FromIterator<(P, Address)> for Addresses<P> {
91 fn from_iter<I: IntoIterator<Item = (P, Address)>>(iter: I) -> Self {
92 Self(Map::from_iter_dedup(iter))
93 }
94}
95
96impl<P: PublicKey> Write for Addresses<P> {
97 fn write(&self, writer: &mut impl bytes::BufMut) {
98 self.0.write(writer);
99 }
100}
101
102impl<P: PublicKey> EncodeSize for Addresses<P> {
103 fn encode_size(&self) -> usize {
104 self.0.encode_size()
105 }
106}
107
108impl<P: PublicKey> Read for Addresses<P> {
109 type Cfg = RangeCfg<usize>;
114
115 fn read_cfg(buf: &mut impl Buf, cfg: &Self::Cfg) -> Result<Self, CodecError> {
116 Ok(Self(Map::read_cfg(buf, &(*cfg, (), ()))?))
117 }
118}
119
120impl<P: PublicKey> Directory<P> for Addresses<P> {
121 fn codec_config(peers: &Set<P>) -> Self::Cfg {
122 RangeCfg::exact(peers.len())
123 }
124
125 fn matches(&self, peers: &Set<P>) -> bool {
126 self.0.keys() == peers
127 }
128}
129
130#[cfg(feature = "arbitrary")]
131impl<P: PublicKey> arbitrary::Arbitrary<'_> for Addresses<P>
132where
133 P: for<'a> arbitrary::Arbitrary<'a>,
134{
135 fn arbitrary(u: &mut arbitrary::Unstructured<'_>) -> arbitrary::Result<Self> {
136 Ok(Self(u.arbitrary()?))
137 }
138}
139
140pub trait Manager: Provider {
142 type Directory: Directory<Self::PublicKey>;
144
145 type Error: std::error::Error + Send + Sync + 'static;
149
150 fn track(
152 &mut self,
153 epoch: Epoch,
154 peers: TrackedPeers<Self::PublicKey>,
155 directory: &Self::Directory,
156 ) -> Result<(), Self::Error>;
157}
158
159impl<M: P2pManager> Manager for M {
160 type Directory = Unit;
161 type Error = Infallible;
162
163 fn track(
164 &mut self,
165 epoch: Epoch,
166 peers: TrackedPeers<Self::PublicKey>,
167 _directory: &Self::Directory,
168 ) -> Result<(), Self::Error> {
169 let _ = P2pManager::track(self, epoch.get(), peers);
170 Ok(())
171 }
172}
173
174#[derive(Clone)]
182pub struct AddressableManager<M> {
183 manager: M,
184}
185
186impl<M> AddressableManager<M> {
187 pub const fn new(manager: M) -> Self {
189 Self { manager }
190 }
191}
192
193impl<M> fmt::Debug for AddressableManager<M> {
194 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
195 f.debug_struct("AddressableManager").finish_non_exhaustive()
196 }
197}
198
199impl<M> Provider for AddressableManager<M>
200where
201 M: P2pAddressableManager,
202{
203 type PublicKey = M::PublicKey;
204
205 async fn peer_set(&mut self, id: u64) -> Option<TrackedPeers<Self::PublicKey>> {
206 self.manager.peer_set(id).await
207 }
208
209 async fn subscribe(&mut self) -> commonware_p2p::PeerSetSubscription<Self::PublicKey> {
210 self.manager.subscribe().await
211 }
212}
213
214#[derive(Clone, Debug, Error, PartialEq, Eq)]
216#[error("epoch directory omitted peer {0:?}")]
217pub struct MissingAddress<P: PublicKey>(pub P);
218
219impl<M> Manager for AddressableManager<M>
220where
221 M: P2pAddressableManager,
222{
223 type Directory = Addresses<M::PublicKey>;
224 type Error = MissingAddress<M::PublicKey>;
225
226 fn track(
227 &mut self,
228 epoch: Epoch,
229 peers: TrackedPeers<Self::PublicKey>,
230 directory: &Self::Directory,
231 ) -> Result<(), Self::Error> {
232 let primary = resolve(&peers.primary, directory)?;
233 let secondary = resolve(&peers.secondary, directory)?;
234 let peers = AddressableTrackedPeers::new(primary, secondary);
235
236 let _ = self.manager.track(epoch.get(), peers);
237 Ok(())
238 }
239}
240
241fn resolve<P: PublicKey>(
242 peers: &Set<P>,
243 directory: &Addresses<P>,
244) -> Result<Map<P, Address>, MissingAddress<P>> {
245 let resolved = peers
246 .iter()
247 .map(|peer| {
248 directory
249 .get(peer)
250 .cloned()
251 .map(|address| (peer.clone(), address))
252 .ok_or_else(|| MissingAddress(peer.clone()))
253 })
254 .collect::<Result<Vec<_>, _>>()?;
255 Ok(Map::from_iter_dedup(resolved))
256}
257
258#[cfg(test)]
259mod tests {
260 use super::*;
261 use commonware_actor::Feedback;
262 use commonware_cryptography::{Signer as _, ed25519};
263 use commonware_macros::test_traced;
264 use commonware_p2p::{
265 PeerSetSubscription, Receiver as _, Recipients, Sender as _, authenticated::lookup,
266 };
267 use commonware_runtime::{
268 Clock as _, Quota, Runner as _, Spawner as _, Supervisor as _, deterministic,
269 };
270 use commonware_utils::{NZU32, NZUsize, channel::mpsc, sync::Mutex};
271 use std::{
272 net::{IpAddr, Ipv4Addr, SocketAddr},
273 sync::Arc,
274 };
275
276 type PublicKey = ed25519::PublicKey;
277 type Tracked = Arc<Mutex<Vec<(u64, TrackedPeers<PublicKey>)>>>;
278 type AddressableTracked = Arc<Mutex<Vec<(u64, AddressableTrackedPeers<PublicKey>)>>>;
279
280 #[derive(Clone, Debug)]
281 struct TestManager {
282 feedback: Feedback,
283 tracked: Tracked,
284 addressable: AddressableTracked,
285 }
286
287 impl TestManager {
288 fn new(feedback: Feedback) -> Self {
289 Self {
290 feedback,
291 tracked: Arc::default(),
292 addressable: Arc::default(),
293 }
294 }
295 }
296
297 impl Provider for TestManager {
298 type PublicKey = PublicKey;
299
300 async fn peer_set(&mut self, _id: u64) -> Option<TrackedPeers<Self::PublicKey>> {
301 None
302 }
303
304 async fn subscribe(&mut self) -> PeerSetSubscription<Self::PublicKey> {
305 let (_, receiver) = mpsc::unbounded_channel();
306 receiver
307 }
308 }
309
310 impl P2pManager for TestManager {
311 fn track<R>(&mut self, id: u64, peers: R) -> Feedback
312 where
313 R: Into<TrackedPeers<Self::PublicKey>> + Send,
314 {
315 self.tracked.lock().push((id, peers.into()));
316 self.feedback
317 }
318 }
319
320 #[derive(Clone, Debug)]
321 struct AddressableTestManager(TestManager);
322
323 impl Provider for AddressableTestManager {
324 type PublicKey = PublicKey;
325
326 async fn peer_set(&mut self, id: u64) -> Option<TrackedPeers<Self::PublicKey>> {
327 self.0.peer_set(id).await
328 }
329
330 async fn subscribe(&mut self) -> PeerSetSubscription<Self::PublicKey> {
331 self.0.subscribe().await
332 }
333 }
334
335 impl P2pAddressableManager for AddressableTestManager {
336 fn track<R>(&mut self, id: u64, peers: R) -> Feedback
337 where
338 R: Into<AddressableTrackedPeers<Self::PublicKey>> + Send,
339 {
340 self.0.addressable.lock().push((id, peers.into()));
341 self.0.feedback
342 }
343
344 fn overwrite(&mut self, _peers: Map<Self::PublicKey, Address>) -> Feedback {
345 self.0.feedback
346 }
347 }
348
349 fn key(seed: u8) -> PublicKey {
350 ed25519::PrivateKey::from_seed(seed.into()).public_key()
351 }
352
353 fn address(port: u16) -> Address {
354 SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port).into()
355 }
356
357 fn peers() -> (TrackedPeers<PublicKey>, [PublicKey; 4]) {
358 let keys = [key(1), key(2), key(3), key(4)];
359 let peers = TrackedPeers::new(
360 Set::from_iter_dedup([keys[0].clone(), keys[1].clone()]),
361 Set::from_iter_dedup([keys[1].clone(), keys[2].clone()]),
362 );
363 (peers, keys)
364 }
365
366 #[test]
367 fn key_only_manager_tracks_peers() {
368 let (peers, _) = peers();
369 let mut manager = TestManager::new(Feedback::Closed);
370 assert_eq!(
371 Manager::track(&mut manager, Epoch::new(7), peers.clone(), &Unit),
372 Ok(())
373 );
374 assert_eq!(manager.tracked.lock()[0], (7, peers));
375 }
376
377 #[test]
378 fn addressable_provider_delegates() {
379 deterministic::Runner::default().start(|_| async move {
380 let inner = TestManager::new(Feedback::Ok);
381 let mut manager = AddressableManager::new(AddressableTestManager(inner));
382
383 assert!(manager.peer_set(7).await.is_none());
384 assert!(manager.subscribe().await.is_closed());
385 });
386 }
387
388 #[test]
389 fn key_only_directory_matches_any_peer_set() {
390 let (peers, _) = peers();
391 assert!(Directory::matches(&Unit, &peers.union()));
392 }
393
394 #[test]
395 fn addresses_requires_exact_peer_set() {
396 let (peers, keys) = peers();
397 let directory =
398 Addresses::from_iter([(keys[0].clone(), address(1)), (keys[1].clone(), address(2))]);
399 assert!(directory.matches(&peers.primary));
400 assert!(!directory.matches(&peers.union()));
401 assert!(!directory.matches(&Set::from_iter_dedup([keys[0].clone()])));
402 }
403
404 #[test]
405 fn addressable_mapping_preserves_roles() {
406 let (peers, keys) = peers();
407 let directory = |offset: u16| {
408 Addresses::from_iter([
409 (keys[0].clone(), address(offset + 1)),
410 (keys[1].clone(), address(offset + 2)),
411 (keys[2].clone(), address(offset + 3)),
412 ])
413 };
414 let inner = TestManager::new(Feedback::Ok);
415 let tracked = inner.addressable.clone();
416 let mut manager = AddressableManager::new(AddressableTestManager(inner));
417
418 Manager::track(&mut manager, Epoch::new(9), peers.clone(), &directory(0)).unwrap();
419 Manager::track(&mut manager, Epoch::new(10), peers.clone(), &directory(10)).unwrap();
420
421 let tracked = tracked.lock();
422 assert_eq!(tracked.len(), 2);
423 assert_eq!(tracked[0].0, 9);
424 assert_eq!(tracked[0].1.primary.keys(), &peers.primary);
425 assert_eq!(tracked[0].1.secondary.keys(), &peers.secondary);
426 assert_eq!(tracked[0].1.primary.get_value(&keys[0]), Some(&address(1)));
427 assert_eq!(tracked[0].1.primary.get_value(&keys[1]), Some(&address(2)));
428 assert_eq!(
429 tracked[0].1.secondary.get_value(&keys[1]),
430 Some(&address(2))
431 );
432 assert_eq!(
433 tracked[0].1.secondary.get_value(&keys[2]),
434 Some(&address(3))
435 );
436 assert_eq!(tracked[1].0, 10);
437 assert_eq!(tracked[1].1.primary.keys(), &peers.primary);
438 assert_eq!(tracked[1].1.secondary.keys(), &peers.secondary);
439 assert_eq!(tracked[1].1.primary.get_value(&keys[0]), Some(&address(11)));
440 assert_eq!(tracked[1].1.primary.get_value(&keys[1]), Some(&address(12)));
441 assert_eq!(
442 tracked[1].1.secondary.get_value(&keys[1]),
443 Some(&address(12))
444 );
445 assert_eq!(
446 tracked[1].1.secondary.get_value(&keys[2]),
447 Some(&address(13))
448 );
449 }
450
451 #[test]
452 fn missing_address_prevents_registration() {
453 let (peers, keys) = peers();
454
455 let inner = TestManager::new(Feedback::Ok);
456 let tracked = inner.addressable.clone();
457 let mut manager = AddressableManager::new(AddressableTestManager(inner));
458 let directory = Addresses::from_iter([(keys[0].clone(), address(1))]);
459 assert_eq!(
460 Manager::track(&mut manager, Epoch::new(1), peers, &directory),
461 Err(MissingAddress(keys[1].clone()))
462 );
463 assert!(tracked.lock().is_empty());
464 }
465
466 #[test_traced]
467 fn lookup_secondary_dials_primary_and_receives_response() {
468 let executor = deterministic::Runner::timed(std::time::Duration::from_secs(10));
469 executor.start(|context| async move {
470 let dealer = ed25519::PrivateKey::from_seed(10);
471 let participant = ed25519::PrivateKey::from_seed(11);
472 let dealer_key = dealer.public_key();
473 let participant_key = participant.public_key();
474 let dealer_socket = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6100);
475 let participant_socket = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6101);
476 let directory = Addresses::from_iter([
477 (dealer_key.clone(), Address::Symmetric(dealer_socket)),
478 (
479 participant_key.clone(),
480 Address::Asymmetric {
481 ingress: participant_socket.into(),
482 egress: participant_socket,
483 },
484 ),
485 ]);
486 let peers = TrackedPeers::new(
487 Set::from_iter_dedup([dealer_key.clone()]),
488 Set::from_iter_dedup([participant_key.clone()]),
489 );
490
491 let (mut dealer_network, dealer_oracle) = lookup::Network::new(
492 context.child("dealer"),
493 lookup::Config::local(
494 dealer,
495 b"_COMMONWARE_GLUE_DKG_LOOKUP_TEST",
496 dealer_socket,
497 NZUsize!(2),
498 1024,
499 ),
500 );
501 let (mut participant_network, participant_oracle) = lookup::Network::new(
502 context.child("participant"),
503 lookup::Config::local(
504 participant,
505 b"_COMMONWARE_GLUE_DKG_LOOKUP_TEST",
506 participant_socket,
507 NZUsize!(2),
508 1024,
509 ),
510 );
511 let (mut dealer_sender, mut dealer_receiver) =
512 dealer_network.register(0, Quota::per_second(NZU32!(100)));
513 let (mut participant_sender, mut participant_receiver) =
514 participant_network.register(0, Quota::per_second(NZU32!(100)));
515
516 let mut dealer_manager = AddressableManager::new(dealer_oracle);
517 let mut participant_manager = AddressableManager::new(participant_oracle);
518 Manager::track(
519 &mut dealer_manager,
520 Epoch::new(3),
521 peers.clone(),
522 &directory,
523 )
524 .unwrap();
525 Manager::track(&mut participant_manager, Epoch::new(3), peers, &directory).unwrap();
526
527 dealer_network.start();
528 participant_network.start();
529
530 let request_sender = context.child("request_sender").spawn({
531 let dealer_key = dealer_key.clone();
532 move |context| async move {
533 loop {
534 participant_sender.send(
535 Recipients::One(dealer_key.clone()),
536 b"request".to_vec(),
537 true,
538 );
539 context.sleep(std::time::Duration::from_millis(100)).await;
540 }
541 }
542 });
543
544 let (sender, request) = dealer_receiver.recv().await.unwrap();
545 request_sender.abort();
546 assert_eq!(sender, participant_key);
547 assert_eq!(request.as_ref(), b"request");
548
549 let response_sender = context.child("response_sender").spawn({
550 let participant_key = participant_key.clone();
551 move |context| async move {
552 loop {
553 dealer_sender.send(
554 Recipients::One(participant_key.clone()),
555 b"response".to_vec(),
556 true,
557 );
558 context.sleep(std::time::Duration::from_millis(100)).await;
559 }
560 }
561 });
562
563 let (sender, response) = participant_receiver.recv().await.unwrap();
564 response_sender.abort();
565 assert_eq!(sender, dealer_key);
566 assert_eq!(response.as_ref(), b"response");
567 });
568 }
569}
570
571#[cfg(all(test, feature = "arbitrary"))]
572mod conformance {
573 use super::*;
574 use commonware_codec::conformance::CodecConformance;
575 use commonware_cryptography::ed25519;
576
577 commonware_conformance::conformance_tests! {
578 CodecConformance<Addresses<ed25519::PublicKey>>,
579 }
580}