Skip to main content

atm0s_sdn_network/features/
router_sync.rs

1use std::collections::{HashMap, VecDeque};
2
3use atm0s_sdn_identity::{ConnId, NodeId};
4use atm0s_sdn_router::{
5    core::{DestDelta, Metric, RegistryDelta, RegistryRemoteDestDelta, Router, RouterDelta, RouterDump, RouterSync, TableDelta},
6    shadow::ShadowRouterDelta,
7};
8use derivative::Derivative;
9use sans_io_runtime::{collections::DynamicDeque, TaskSwitcherChild};
10
11use crate::{
12    base::{ConnectionEvent, Feature, FeatureContext, FeatureInput, FeatureOutput, FeatureSharedInput, FeatureWorker, FeatureWorkerContext, FeatureWorkerInput, FeatureWorkerOutput, NetOutgoingMeta},
13    data_plane::NetPair,
14};
15
16pub const FEATURE_ID: u8 = 2;
17pub const FEATURE_NAME: &str = "router_sync";
18
19const INIT_RTT_MS: u16 = 1000;
20const INIT_BW: u32 = 100_000_000;
21
22#[derive(Debug, Clone, PartialEq, Eq)]
23pub enum Control {
24    DumpRouter,
25}
26
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub enum Event {
29    DumpRouter(Box<RouterDump>),
30}
31
32pub type ToWorker = ShadowRouterDelta<ConnId, NetPair>;
33pub type ToController = ();
34
35pub type Output<UserData> = FeatureOutput<UserData, Event, ToWorker>;
36pub type WorkerOutput<UserData> = FeatureWorkerOutput<UserData, Control, Event, ToController>;
37
38pub struct RouterSyncFeature<UserData> {
39    router: Router,
40    conns: HashMap<ConnId, (NodeId, NetPair, Metric)>,
41    queue: VecDeque<Output<UserData>>,
42    services: Vec<u8>,
43    shutdown: bool,
44}
45
46impl<UserData> RouterSyncFeature<UserData> {
47    pub fn new(node: NodeId, services: Vec<u8>) -> Self {
48        log::info!("[RouterSync] started node {} with public services {:?}", node, services);
49
50        Self {
51            router: Router::new(node),
52            services,
53            conns: HashMap::new(),
54            queue: VecDeque::new(),
55            shutdown: false,
56        }
57    }
58
59    fn send_sync_to(router: &Router, queue: &mut VecDeque<Output<UserData>>, conn: ConnId, node: NodeId) {
60        let sync = router.create_sync(node);
61        log::debug!("[RouterSync] send sync to {node} content {sync:?}");
62        queue.push_back(FeatureOutput::SendDirect(
63            conn,
64            NetOutgoingMeta::new(false, 1.into(), 0, true),
65            bincode::serialize(&sync).expect("").into(),
66        ));
67    }
68}
69
70impl<UserData> Feature<UserData, Control, Event, ToController, ToWorker> for RouterSyncFeature<UserData> {
71    fn on_shared_input(&mut self, _ctx: &FeatureContext, _now: u64, input: FeatureSharedInput) {
72        match input {
73            FeatureSharedInput::Tick(tick_count) => {
74                if tick_count < 1 {
75                    //we need to wait all workers to be ready
76                    return;
77                }
78
79                while let Some(service) = self.services.pop() {
80                    log::info!("[RouterSync] register local service {}", service);
81                    self.router.register_service(service);
82                }
83
84                for (conn, (node, _, _)) in self.conns.iter() {
85                    Self::send_sync_to(&self.router, &mut self.queue, *conn, *node);
86                }
87            }
88            FeatureSharedInput::Connection(event) => match event {
89                ConnectionEvent::Connecting(_conn_ctx) => {}
90                ConnectionEvent::ConnectError(_conn_ctx, _err) => {}
91                ConnectionEvent::Connected(conn_ctx, _) => {
92                    log::info!("[RouterSync] Connection {} connected", conn_ctx.pair);
93                    let metric = Metric::direct(INIT_RTT_MS, conn_ctx.node, INIT_BW);
94                    self.conns.insert(conn_ctx.conn, (conn_ctx.node, conn_ctx.pair, metric.clone()));
95                    self.router.set_direct(conn_ctx.conn, conn_ctx.node, metric);
96                    Self::send_sync_to(&self.router, &mut self.queue, conn_ctx.conn, conn_ctx.node);
97                }
98                ConnectionEvent::Stats(conn_ctx, stats) => {
99                    log::debug!("[RouterSync] Connection {} stats rtt_ms {}", conn_ctx.pair, stats.rtt_ms);
100                    let metric = Metric::direct(stats.rtt_ms as u16, conn_ctx.node, INIT_BW);
101                    self.conns.insert(conn_ctx.conn, (conn_ctx.node, conn_ctx.pair, metric.clone()));
102                    self.router.set_direct(conn_ctx.conn, conn_ctx.node, metric);
103                }
104                ConnectionEvent::Disconnected(conn_ctx) => {
105                    log::info!("[RouterSync] Connection {} disconnected", conn_ctx.pair);
106                    self.conns.remove(&conn_ctx.conn);
107                    self.router.del_direct(conn_ctx.conn);
108                }
109            },
110        }
111    }
112
113    fn on_input(&mut self, _ctx: &FeatureContext, _now_ms: u64, input: FeatureInput<'_, UserData, Control, ToController>) {
114        match input {
115            FeatureInput::FromWorker(_) => {}
116            FeatureInput::Control(actor, control) => match control {
117                Control::DumpRouter => {
118                    self.queue.push_back(FeatureOutput::Event(actor, Event::DumpRouter(Box::new(self.router.dump()))));
119                }
120            },
121            FeatureInput::Net(conn_ctx, meta, buf) => {
122                if !meta.secure {
123                    log::warn!("[RouterSync] reject unsecure message");
124                    return;
125                }
126                if let Some((node, remote, metric)) = self.conns.get(&conn_ctx.conn) {
127                    if let Ok(sync) = bincode::deserialize::<RouterSync>(&buf) {
128                        log::debug!("[RouterSync] Receive sync from {node} {remote:?}");
129                        self.router.apply_sync(conn_ctx.conn, conn_ctx.node, metric.clone(), sync);
130                    } else {
131                        log::warn!("[RouterSync] Receive invalid sync from {}", conn_ctx.pair);
132                    }
133                } else {
134                    log::warn!("[RouterSync] Receive sync from unknown connection {}", conn_ctx.pair);
135                }
136            }
137            FeatureInput::Local(..) => {}
138        }
139    }
140
141    fn on_shutdown(&mut self, _ctx: &FeatureContext, _now: u64) {
142        log::info!("[RouterSync] Shutdown");
143        self.shutdown = true;
144    }
145}
146
147impl<UserData> TaskSwitcherChild<Output<UserData>> for RouterSyncFeature<UserData> {
148    type Time = u64;
149
150    fn is_empty(&self) -> bool {
151        self.shutdown && self.queue.is_empty()
152    }
153
154    fn empty_event(&self) -> Output<UserData> {
155        Output::OnResourceEmpty
156    }
157
158    fn pop_output(&mut self, _now: u64) -> Option<Output<UserData>> {
159        if let Some(rule) = self.router.pop_delta() {
160            log::debug!("[RouterSync] broadcast to all workers {:?}", rule);
161            let rule = match rule {
162                RouterDelta::Table(layer, TableDelta(index, DestDelta::SetBestPath(conn))) => ShadowRouterDelta::SetTable {
163                    layer,
164                    index,
165                    conn,
166                    remote: self.conns.get(&conn)?.1,
167                },
168                RouterDelta::Table(layer, TableDelta(index, DestDelta::DelBestPath)) => ShadowRouterDelta::DelTable { layer, index },
169                RouterDelta::Registry(RegistryDelta::SetServiceLocal(service)) => ShadowRouterDelta::SetServiceLocal { service },
170                RouterDelta::Registry(RegistryDelta::DelServiceLocal(service)) => ShadowRouterDelta::DelServiceLocal { service },
171                RouterDelta::Registry(RegistryDelta::ServiceRemote(service, RegistryRemoteDestDelta::SetServicePath(conn, dest, score))) => {
172                    let conn_info = self.conns.get(&conn)?;
173                    ShadowRouterDelta::SetServiceRemote {
174                        service,
175                        conn,
176                        remote: conn_info.1,
177                        next: conn_info.0,
178                        dest,
179                        score,
180                    }
181                }
182                RouterDelta::Registry(RegistryDelta::ServiceRemote(service, RegistryRemoteDestDelta::DelServicePath(conn))) => ShadowRouterDelta::DelServiceRemote { service, conn },
183            };
184            return Some(FeatureOutput::ToWorker(true, rule));
185        }
186        self.queue.pop_front()
187    }
188}
189
190#[derive(Derivative)]
191#[derivative(Default(bound = ""))]
192pub struct RouterSyncFeatureWorker<UserData> {
193    queue: DynamicDeque<WorkerOutput<UserData>, 1>,
194    shutdown: bool,
195}
196
197impl<UserData> FeatureWorker<UserData, Control, Event, ToController, ToWorker> for RouterSyncFeatureWorker<UserData> {
198    fn on_input(&mut self, ctx: &mut FeatureWorkerContext, _now: u64, input: FeatureWorkerInput<UserData, Control, ToWorker>) {
199        match input {
200            FeatureWorkerInput::Control(service, control) => self.queue.push_back(FeatureWorkerOutput::ForwardControlToController(service, control)),
201            FeatureWorkerInput::Network(conn, header, msg) => self.queue.push_back(FeatureWorkerOutput::ForwardNetworkToController(conn, header, msg)),
202            FeatureWorkerInput::FromController(_, delta) => {
203                log::debug!("[RouterSyncWorker] apply router delta {:?}", delta);
204                ctx.router.apply_delta(delta);
205            }
206            FeatureWorkerInput::Local(_header, _msg) => {
207                log::warn!("No handler for local message in {}", FEATURE_NAME);
208            }
209            #[cfg(feature = "vpn")]
210            FeatureWorkerInput::TunPkt(_buf) => {
211                log::warn!("No handler for tun packet in {}", FEATURE_NAME);
212            }
213        }
214    }
215
216    fn on_shutdown(&mut self, _ctx: &mut FeatureWorkerContext, _now: u64) {
217        log::info!("[RouterSyncFeatureWorker] Shutdown");
218        self.shutdown = true;
219    }
220}
221
222impl<UserData> TaskSwitcherChild<WorkerOutput<UserData>> for RouterSyncFeatureWorker<UserData> {
223    type Time = u64;
224
225    fn is_empty(&self) -> bool {
226        self.shutdown && self.queue.is_empty()
227    }
228
229    fn empty_event(&self) -> WorkerOutput<UserData> {
230        WorkerOutput::OnResourceEmpty
231    }
232
233    fn pop_output(&mut self, _now: u64) -> Option<WorkerOutput<UserData>> {
234        self.queue.pop_front()
235    }
236}
237
238#[cfg(test)]
239mod tests {
240    use atm0s_sdn_identity::ConnId;
241    use atm0s_sdn_router::core::{Metric, RegistrySync, RouterSync, TableSync};
242    use sans_io_runtime::TaskSwitcherChild;
243
244    use crate::{
245        base::{ConnectionCtx, ConnectionEvent, Feature, FeatureContext, FeatureInput, FeatureSharedInput, MockDecryptor, MockEncryptor, NetIncomingMeta, SecureContext, Ttl},
246        data_plane::NetPair,
247    };
248
249    use super::{Output, RouterSyncFeature, ShadowRouterDelta};
250
251    #[test]
252    fn should_send_registry_delta_after_disconnect() {
253        let service_id = 0;
254        let remote_node = 1;
255        let ctx = FeatureContext { node_id: 0, session: 0 };
256        let conn = ConnectionCtx {
257            conn: ConnId::from_in(0, 0),
258            node: remote_node,
259            pair: NetPair::new("127.0.0.1:1000".parse().unwrap(), "127.0.0.1:1001".parse().unwrap()),
260        };
261        let mut service = RouterSyncFeature::<()>::new(0, vec![]);
262        let sync = RouterSync(RegistrySync(vec![(service_id, Metric::local())]), [None, None, None, None]);
263        let sync_msg = bincode::serialize(&sync).expect("");
264
265        // fake connected
266        service.on_shared_input(
267            &ctx,
268            0,
269            FeatureSharedInput::Connection(ConnectionEvent::Connected(
270                conn.clone(),
271                SecureContext {
272                    encryptor: Box::new(MockEncryptor::new()),
273                    decryptor: Box::new(MockDecryptor::new()),
274                },
275            )),
276        );
277        assert_eq!(
278            service.pop_output(0),
279            Some(Output::ToWorker(
280                true,
281                ShadowRouterDelta::SetTable {
282                    layer: 0,
283                    index: 1,
284                    conn: conn.conn,
285                    remote: conn.pair
286                }
287            ))
288        );
289
290        // fake sync service
291        service.on_input(&ctx, 0, FeatureInput::Net(&conn, NetIncomingMeta::new(Some(1), Ttl(1), 0, true), sync_msg.into()));
292
293        assert_eq!(
294            service.pop_output(0),
295            Some(Output::ToWorker(
296                true,
297                ShadowRouterDelta::SetServiceRemote {
298                    service: service_id,
299                    conn: conn.conn,
300                    remote: conn.pair,
301                    next: remote_node,
302                    dest: remote_node,
303                    score: 1010,
304                }
305            ))
306        );
307
308        // fake disconnected
309        service.on_shared_input(&ctx, 0, FeatureSharedInput::Connection(ConnectionEvent::Disconnected(conn.clone())));
310
311        assert_eq!(
312            service.pop_output(0),
313            Some(Output::ToWorker(true, ShadowRouterDelta::DelServiceRemote { service: service_id, conn: conn.conn }))
314        );
315        assert_eq!(service.pop_output(0), Some(Output::ToWorker(true, ShadowRouterDelta::DelTable { layer: 0, index: 1 })));
316    }
317
318    #[test]
319    fn router_sync_should_fit_udp() {
320        const MAX_SIZE: usize = 1200;
321        const NUMBER_SERVICES: usize = 2;
322        const NUMBER_NEIGHBORS: usize = 10;
323        const NUMBER_NODE_PATH: u32 = 3;
324
325        let mut service_sync = RegistrySync(vec![]);
326        let mut table_sync = [None, None, None, None];
327
328        for _ in 0..NUMBER_SERVICES {
329            service_sync.0.push((rand::random(), Metric::new(0, (0..NUMBER_NODE_PATH).collect::<Vec<_>>(), 0)));
330        }
331
332        for i in &mut table_sync {
333            let mut table = TableSync(vec![]);
334            for _ in 0..NUMBER_NEIGHBORS {
335                table.0.push((rand::random(), Metric::new(0, (0..NUMBER_NODE_PATH).collect::<Vec<_>>(), 0)));
336            }
337            *i = Some(table);
338        }
339
340        let sync = RouterSync(service_sync, table_sync);
341        let sync_msg_len = bincode::serialize(&sync).expect("").len();
342        assert!(sync_msg_len <= MAX_SIZE, "SYNC msg not fit in UDP {} vs {}", sync_msg_len, MAX_SIZE);
343    }
344}