atm0s_sdn_network/features/
data.rs1use 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 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}