Skip to main content

codex_sync/
discovery.rs

1use crate::config::Config;
2use anyhow::Result;
3use serde::{Deserialize, Serialize};
4use std::{
5    collections::HashMap,
6    net::{IpAddr, Ipv4Addr, SocketAddr},
7    time::Duration,
8};
9use tokio::net::UdpSocket;
10
11#[derive(Debug, Clone, Serialize, Deserialize)]
12pub struct Peer {
13    pub node_id: String,
14    pub name: String,
15    pub address: String,
16}
17
18pub async fn responder(cfg: Config) -> Result<()> {
19    let socket = UdpSocket::bind((Ipv4Addr::UNSPECIFIED, cfg.discovery_port)).await?;
20    socket.set_broadcast(true)?;
21    let mut buf = [0u8; 512];
22    loop {
23        let (n, from) = socket.recv_from(&mut buf).await?;
24        if &buf[..n] == b"CODEX_SYNC_DISCOVER_V1" {
25            let port = cfg.listen.rsplit(':').next().unwrap_or("8787");
26            let peer = Peer {
27                node_id: cfg.node_id(),
28                name: cfg.node_name.clone(),
29                address: format!("{}:{port}", from.ip()),
30            };
31            socket.send_to(&serde_json::to_vec(&peer)?, from).await?;
32        }
33    }
34}
35
36pub async fn discover(cfg: &Config, seconds: u64) -> Result<Vec<Peer>> {
37    let socket = UdpSocket::bind((Ipv4Addr::UNSPECIFIED, 0)).await?;
38    socket.set_broadcast(true)?;
39    socket
40        .send_to(
41            b"CODEX_SYNC_DISCOVER_V1",
42            SocketAddr::new(IpAddr::V4(Ipv4Addr::BROADCAST), cfg.discovery_port),
43        )
44        .await?;
45    let deadline = tokio::time::Instant::now() + Duration::from_secs(seconds);
46    let mut peers = HashMap::new();
47    let mut buf = [0u8; 1024];
48    loop {
49        let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
50        if remaining.is_zero() {
51            break;
52        }
53        match tokio::time::timeout(remaining, socket.recv_from(&mut buf)).await {
54            Ok(Ok((n, _))) => {
55                if let Ok(peer) = serde_json::from_slice::<Peer>(&buf[..n])
56                    && peer.node_id != cfg.node_id()
57                {
58                    peers.insert(peer.node_id.clone(), peer);
59                }
60            }
61            _ => break,
62        }
63    }
64    Ok(peers.into_values().collect())
65}