Skip to main content

atm0s_sdn_network/base/
feature.rs

1use atm0s_sdn_identity::{ConnId, NodeAddr, NodeId};
2use atm0s_sdn_router::{shadow::ShadowRouter, RouteRule};
3use sans_io_runtime::TaskSwitcherChild;
4
5use crate::data_plane::NetPair;
6
7use super::{Buffer, ConnectionCtx, ConnectionEvent, ServiceId, TransportMsgHeader, Ttl};
8
9#[derive(Debug, Default, Clone, PartialEq, Eq)]
10pub struct NetIncomingMeta {
11    pub source: Option<NodeId>,
12    pub ttl: Ttl,
13    pub meta: u8,
14    pub secure: bool,
15}
16
17impl NetIncomingMeta {
18    pub fn new(source: Option<NodeId>, ttl: Ttl, meta: u8, secure: bool) -> Self {
19        Self { source, ttl, meta, secure }
20    }
21}
22
23impl From<&TransportMsgHeader> for NetIncomingMeta {
24    fn from(value: &TransportMsgHeader) -> Self {
25        Self {
26            source: value.from_node,
27            ttl: Ttl(value.ttl),
28            meta: value.meta,
29            secure: value.encrypt,
30        }
31    }
32}
33
34#[derive(Debug, Default, Clone, PartialEq, Eq)]
35pub struct NetOutgoingMeta {
36    pub source: bool,
37    pub ttl: Ttl,
38    pub meta: u8,
39    pub secure: bool,
40}
41
42impl NetOutgoingMeta {
43    pub fn new(source: bool, ttl: Ttl, meta: u8, secure: bool) -> Self {
44        Self { source, ttl, meta, secure }
45    }
46
47    pub fn secure() -> Self {
48        Self {
49            source: false,
50            ttl: Ttl::default(),
51            meta: 0,
52            secure: true,
53        }
54    }
55
56    pub fn to_header(&self, feature: u8, rule: RouteRule, node_id: NodeId) -> TransportMsgHeader {
57        TransportMsgHeader::build(feature, self.meta, rule)
58            .set_ttl(*self.ttl)
59            .set_from_node(if self.source {
60                Some(node_id)
61            } else {
62                None
63            })
64            .set_encrypt(self.secure)
65    }
66
67    pub fn to_incoming(&self, node_id: NodeId) -> NetIncomingMeta {
68        NetIncomingMeta {
69            source: if self.source {
70                Some(node_id)
71            } else {
72                None
73            },
74            ttl: self.ttl,
75            meta: self.meta,
76            secure: self.secure,
77        }
78    }
79}
80
81#[derive(Debug, Clone, Copy, Hash, Eq, PartialEq)]
82pub enum FeatureControlActor<UserData> {
83    Controller(UserData),
84    Worker(u16, UserData),
85    Service(ServiceId),
86}
87
88impl<UserData> FeatureControlActor<UserData> {
89    pub fn into2<UserData2>(self) -> FeatureControlActor<UserData2>
90    where
91        UserData2: From<UserData>,
92    {
93        match self {
94            FeatureControlActor::Controller(u) => FeatureControlActor::Controller(u.into()),
95            FeatureControlActor::Worker(worker, u) => FeatureControlActor::Worker(worker, u.into()),
96            FeatureControlActor::Service(service) => FeatureControlActor::Service(service),
97        }
98    }
99}
100
101#[derive(Debug, Clone)]
102pub enum FeatureSharedInput {
103    Tick(u64),
104    Connection(ConnectionEvent),
105}
106
107#[derive(Debug, Clone)]
108pub enum FeatureInput<'a, UserData, Control, ToController> {
109    FromWorker(ToController),
110    Control(FeatureControlActor<UserData>, Control),
111    Net(&'a ConnectionCtx, NetIncomingMeta, Buffer),
112    Local(NetIncomingMeta, Buffer),
113}
114
115#[derive(Debug, PartialEq, Eq)]
116pub enum FeatureOutput<UserData, Event, ToWorker> {
117    /// First bool is flag for broadcast or not
118    ToWorker(bool, ToWorker),
119    Event(FeatureControlActor<UserData>, Event),
120    SendDirect(ConnId, NetOutgoingMeta, Buffer),
121    SendRoute(RouteRule, NetOutgoingMeta, Buffer),
122    NeighboursConnectTo(NodeAddr),
123    NeighboursDisconnectFrom(NodeId),
124    OnResourceEmpty,
125}
126
127impl<UserData, Event, ToWorker> FeatureOutput<UserData, Event, ToWorker> {
128    pub fn into2<UserData2, Event2, ToWorker2>(self) -> FeatureOutput<UserData2, Event2, ToWorker2>
129    where
130        UserData2: From<UserData>,
131        Event2: From<Event>,
132        ToWorker2: From<ToWorker>,
133    {
134        match self {
135            FeatureOutput::ToWorker(is_broadcast, to) => FeatureOutput::ToWorker(is_broadcast, to.into()),
136            FeatureOutput::Event(actor, event) => FeatureOutput::Event(actor.into2(), event.into()),
137            FeatureOutput::SendDirect(conn, meta, msg) => FeatureOutput::SendDirect(conn, meta, msg),
138            FeatureOutput::SendRoute(rule, ttl, buf) => FeatureOutput::SendRoute(rule, ttl, buf),
139            FeatureOutput::NeighboursConnectTo(addr) => FeatureOutput::NeighboursConnectTo(addr),
140            FeatureOutput::NeighboursDisconnectFrom(id) => FeatureOutput::NeighboursDisconnectFrom(id),
141            FeatureOutput::OnResourceEmpty => FeatureOutput::OnResourceEmpty,
142        }
143    }
144}
145
146pub struct FeatureContext {
147    pub node_id: NodeId,
148    pub session: u64,
149}
150
151pub trait Feature<UserData, Control, Event, ToController, ToWorker>: TaskSwitcherChild<FeatureOutput<UserData, Event, ToWorker>> {
152    fn on_shared_input(&mut self, _ctx: &FeatureContext, _now: u64, _input: FeatureSharedInput);
153    fn on_input(&mut self, _ctx: &FeatureContext, now_ms: u64, input: FeatureInput<'_, UserData, Control, ToController>);
154    fn on_shutdown(&mut self, _ctx: &FeatureContext, _now: u64);
155}
156
157pub enum FeatureWorkerInput<UserData, Control, ToWorker> {
158    /// First bool is flag for broadcast or not
159    FromController(bool, ToWorker),
160    Control(FeatureControlActor<UserData>, Control),
161    Network(ConnId, NetIncomingMeta, Buffer),
162    Local(NetIncomingMeta, Buffer),
163    #[cfg(feature = "vpn")]
164    TunPkt(Buffer),
165}
166
167#[derive(Debug, Clone)]
168pub enum FeatureWorkerOutput<UserData, Control, Event, ToController> {
169    ForwardControlToController(FeatureControlActor<UserData>, Control),
170    ForwardNetworkToController(ConnId, NetIncomingMeta, Buffer),
171    ForwardLocalToController(NetIncomingMeta, Buffer),
172    ToController(ToController),
173    Event(FeatureControlActor<UserData>, Event),
174    SendDirect(ConnId, NetOutgoingMeta, Buffer),
175    SendRoute(RouteRule, NetOutgoingMeta, Buffer),
176    RawDirect(ConnId, Buffer),
177    RawBroadcast(Vec<ConnId>, Buffer),
178    RawDirect2(NetPair, Buffer),
179    RawBroadcast2(Vec<NetPair>, Buffer),
180    #[cfg(feature = "vpn")]
181    TunPkt(Buffer),
182    OnResourceEmpty,
183}
184
185impl<UserData, Control, Event, ToController> FeatureWorkerOutput<UserData, Control, Event, ToController> {
186    pub fn into2<UserData2, Control2, Event2, ToController2>(self) -> FeatureWorkerOutput<UserData2, Control2, Event2, ToController2>
187    where
188        UserData2: From<UserData>,
189        Control2: From<Control>,
190        Event2: From<Event>,
191        ToController2: From<ToController>,
192    {
193        match self {
194            FeatureWorkerOutput::ForwardControlToController(actor, control) => FeatureWorkerOutput::ForwardControlToController(actor.into2(), control.into()),
195            FeatureWorkerOutput::ForwardNetworkToController(conn, header, msg) => FeatureWorkerOutput::ForwardNetworkToController(conn, header, msg),
196            FeatureWorkerOutput::ForwardLocalToController(header, buf) => FeatureWorkerOutput::ForwardLocalToController(header, buf),
197            FeatureWorkerOutput::ToController(to) => FeatureWorkerOutput::ToController(to.into()),
198            FeatureWorkerOutput::Event(actor, event) => FeatureWorkerOutput::Event(actor.into2(), event.into()),
199            FeatureWorkerOutput::SendDirect(conn, meta, buf) => FeatureWorkerOutput::SendDirect(conn, meta, buf),
200            FeatureWorkerOutput::SendRoute(route, meta, buf) => FeatureWorkerOutput::SendRoute(route, meta, buf),
201            FeatureWorkerOutput::RawDirect(conn, buf) => FeatureWorkerOutput::RawDirect(conn, buf),
202            FeatureWorkerOutput::RawBroadcast(conns, buf) => FeatureWorkerOutput::RawBroadcast(conns, buf),
203            FeatureWorkerOutput::RawDirect2(conn, buf) => FeatureWorkerOutput::RawDirect2(conn, buf),
204            FeatureWorkerOutput::RawBroadcast2(conns, buf) => FeatureWorkerOutput::RawBroadcast2(conns, buf),
205            #[cfg(feature = "vpn")]
206            FeatureWorkerOutput::TunPkt(buf) => FeatureWorkerOutput::TunPkt(buf),
207            FeatureWorkerOutput::OnResourceEmpty => FeatureWorkerOutput::OnResourceEmpty,
208        }
209    }
210}
211
212pub struct FeatureWorkerContext {
213    pub node_id: NodeId,
214    pub router: ShadowRouter<ConnId, NetPair>,
215}
216
217pub trait FeatureWorker<UserData, SdkControl, SdkEvent, ToController, ToWorker>: TaskSwitcherChild<FeatureWorkerOutput<UserData, SdkControl, SdkEvent, ToController>> {
218    fn on_tick(&mut self, _ctx: &mut FeatureWorkerContext, _now: u64, _tick_count: u64) {}
219    fn on_network_raw(&mut self, ctx: &mut FeatureWorkerContext, now: u64, conn: ConnId, _pair: NetPair, header: TransportMsgHeader, mut buf: Buffer) {
220        let header_len = header.serialize_size();
221        buf.move_front_right(header_len).expect("Buffer should bigger or equal header");
222        self.on_input(ctx, now, FeatureWorkerInput::Network(conn, (&header).into(), buf));
223    }
224    fn on_input(&mut self, _ctx: &mut FeatureWorkerContext, _now: u64, input: FeatureWorkerInput<UserData, SdkControl, ToWorker>);
225    fn on_shutdown(&mut self, _ctx: &mut FeatureWorkerContext, _now: u64);
226}