Skip to main content

atm0s_sdn_network/features/
data.rs

1use std::{
2    collections::{HashMap, VecDeque},
3    fmt::Debug,
4};
5
6use atm0s_sdn_identity::NodeId;
7use atm0s_sdn_router::RouteRule;
8use derivative::Derivative;
9use sans_io_runtime::{collections::DynamicDeque, TaskSwitcherChild};
10use serde::{Deserialize, Serialize};
11
12use crate::base::{
13    Feature, FeatureContext, FeatureControlActor, FeatureInput, FeatureOutput, FeatureSharedInput, FeatureWorker, FeatureWorkerInput, FeatureWorkerOutput, NetIncomingMeta, NetOutgoingMeta,
14};
15
16pub const FEATURE_ID: u8 = 1;
17pub const FEATURE_NAME: &str = "data_transfer";
18
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub enum Control {
21    Ping(NodeId),
22    DataListen(u16),
23    DataUnlisten(u16),
24    DataSendRule(u16, RouteRule, NetOutgoingMeta, Vec<u8>),
25}
26
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub enum Event {
29    Pong(NodeId, Option<u16>),
30    Recv(u16, NetIncomingMeta, Vec<u8>),
31}
32
33#[derive(Debug, Clone)]
34pub struct ToWorker;
35
36#[derive(Debug, Clone)]
37pub struct ToController;
38
39#[derive(Debug, Serialize, Deserialize)]
40enum DataMsg {
41    Ping { id: u64, ts: u64, from: NodeId },
42    Pong { id: u64, ts: u64 },
43    Data(u16, Vec<u8>),
44}
45
46pub type Output<UserData> = FeatureOutput<UserData, Event, ToWorker>;
47pub type WorkerOutput<UserData> = FeatureWorkerOutput<UserData, Control, Event, ToController>;
48
49pub struct DataFeature<UserData> {
50    waits: HashMap<u64, (u64, FeatureControlActor<UserData>, NodeId)>,
51    ping_seq: u64,
52    queue: VecDeque<Output<UserData>>,
53    data_dest: HashMap<u16, FeatureControlActor<UserData>>,
54    shutdown: bool,
55}
56
57impl<UserData> Default for DataFeature<UserData> {
58    fn default() -> Self {
59        Self {
60            waits: HashMap::new(),
61            ping_seq: 0,
62            queue: VecDeque::new(),
63            data_dest: HashMap::new(),
64            shutdown: false,
65        }
66    }
67}
68
69impl<UserData: Copy> Feature<UserData, Control, Event, ToController, ToWorker> for DataFeature<UserData> {
70    fn on_shared_input(&mut self, _ctx: &FeatureContext, now: u64, input: FeatureSharedInput) {
71        if let FeatureSharedInput::Tick(_) = input {
72            //clean timeout ping
73            let mut timeout_list = Vec::new();
74            for (id, (sent_ms, _, _)) in self.waits.iter() {
75                if now >= sent_ms + 2000 {
76                    timeout_list.push(*id);
77                }
78            }
79
80            for id in timeout_list {
81                let (_, actor, dest) = self.waits.remove(&id).expect("Should have");
82                self.queue.push_back(FeatureOutput::Event(actor, Event::Pong(dest, None)));
83            }
84        }
85    }
86
87    fn on_input(&mut self, ctx: &FeatureContext, now_ms: u64, input: FeatureInput<'_, UserData, Control, ToController>) {
88        match input {
89            FeatureInput::Control(actor, control) => match control {
90                Control::Ping(dest) => {
91                    log::info!("[DataFeature] send ping to: {}", dest);
92                    let seq = self.ping_seq;
93                    self.ping_seq += 1;
94                    self.waits.insert(seq, (now_ms, actor, dest));
95                    let msg = bincode::serialize(&DataMsg::Ping {
96                        id: seq,
97                        ts: now_ms,
98                        from: ctx.node_id,
99                    })
100                    .expect("should work");
101                    let rule = RouteRule::ToNode(dest);
102                    self.queue.push_back(FeatureOutput::SendRoute(rule, NetOutgoingMeta::default(), msg.into()));
103                }
104                Control::DataListen(port) => {
105                    self.data_dest.insert(port, actor);
106                }
107                Control::DataUnlisten(port) => {
108                    self.data_dest.remove(&port);
109                }
110                Control::DataSendRule(port, rule, ttl, data) => {
111                    let data = DataMsg::Data(port, data);
112                    let msg = bincode::serialize(&data).expect("should work");
113                    self.queue.push_back(FeatureOutput::SendRoute(rule, ttl, msg.into()));
114                }
115            },
116            FeatureInput::Net(_, meta, buf) | FeatureInput::Local(meta, buf) => {
117                log::debug!("[DataFeature] on message from {:?} len {}", meta.source, buf.len());
118                if let Ok(msg) = bincode::deserialize::<DataMsg>(&buf) {
119                    match msg {
120                        DataMsg::Pong { id, ts } => {
121                            if let Some((_, actor, dest)) = self.waits.remove(&id) {
122                                self.queue.push_back(FeatureOutput::Event(actor, Event::Pong(dest, Some((now_ms - ts) as u16))));
123                            } else {
124                                log::warn!("[DataFeature] pong with unknown id: {}", id);
125                            }
126                        }
127                        DataMsg::Ping { id, ts, from } => {
128                            log::info!("[DataFeature] got ping from: {}", from);
129                            let msg = bincode::serialize(&DataMsg::Pong { id, ts }).expect("should work");
130                            let rule = RouteRule::ToNode(from);
131                            self.queue.push_back(FeatureOutput::SendRoute(rule, NetOutgoingMeta::default(), msg.into()));
132                        }
133                        DataMsg::Data(port, data) => {
134                            if let Some(actor) = self.data_dest.get(&port) {
135                                self.queue.push_back(FeatureOutput::Event(*actor, Event::Recv(port, meta, data)));
136                            }
137                        }
138                    }
139                }
140            }
141            _ => {}
142        }
143    }
144
145    fn on_shutdown(&mut self, _ctx: &FeatureContext, _now: u64) {
146        self.shutdown = true;
147    }
148}
149
150impl<UserData> TaskSwitcherChild<Output<UserData>> for DataFeature<UserData> {
151    type Time = u64;
152
153    fn is_empty(&self) -> bool {
154        self.shutdown && self.queue.is_empty()
155    }
156
157    fn empty_event(&self) -> Output<UserData> {
158        Output::OnResourceEmpty
159    }
160
161    fn pop_output(&mut self, _now: u64) -> Option<Output<UserData>> {
162        self.queue.pop_front()
163    }
164}
165
166#[derive(Derivative)]
167#[derivative(Default(bound = ""))]
168pub struct DataFeatureWorker<UserData> {
169    queue: DynamicDeque<WorkerOutput<UserData>, 1>,
170    shutdown: bool,
171}
172
173impl<UserData: Debug> FeatureWorker<UserData, Control, Event, ToController, ToWorker> for DataFeatureWorker<UserData> {
174    fn on_input(&mut self, _ctx: &mut crate::base::FeatureWorkerContext, _now: u64, input: crate::base::FeatureWorkerInput<UserData, Control, ToWorker>) {
175        match input {
176            FeatureWorkerInput::Control(actor, control) => self.queue.push_back(FeatureWorkerOutput::ForwardControlToController(actor, control)),
177            FeatureWorkerInput::Network(conn, header, buf) => self.queue.push_back(FeatureWorkerOutput::ForwardNetworkToController(conn, header, buf)),
178            #[cfg(feature = "vpn")]
179            FeatureWorkerInput::TunPkt(..) => {}
180            FeatureWorkerInput::FromController(..) => {
181                log::warn!("No handler for FromController");
182            }
183            FeatureWorkerInput::Local(header, buf) => self.queue.push_back(FeatureWorkerOutput::ForwardLocalToController(header, buf)),
184        }
185    }
186
187    fn on_shutdown(&mut self, _ctx: &mut crate::base::FeatureWorkerContext, _now: u64) {
188        log::info!("[DataFeatureWorker] Shutdown");
189        self.shutdown = true;
190    }
191}
192
193impl<UserData> TaskSwitcherChild<WorkerOutput<UserData>> for DataFeatureWorker<UserData> {
194    type Time = u64;
195
196    fn is_empty(&self) -> bool {
197        self.shutdown && self.queue.is_empty()
198    }
199
200    fn empty_event(&self) -> WorkerOutput<UserData> {
201        WorkerOutput::OnResourceEmpty
202    }
203
204    fn pop_output(&mut self, _now: u64) -> Option<WorkerOutput<UserData>> {
205        self.queue.pop_front()
206    }
207}