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 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 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 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 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}