Skip to main content

atm0s_sdn_network/features/
alias.rs

1use std::{
2    collections::{HashMap, VecDeque},
3    fmt::Debug,
4};
5
6use atm0s_sdn_identity::NodeId;
7use atm0s_sdn_router::{RouteRule, ServiceBroadcastLevel};
8use derivative::Derivative;
9use sans_io_runtime::{collections::DynamicDeque, TaskSwitcherChild};
10use serde::{Deserialize, Serialize};
11
12use crate::base::{Feature, FeatureContext, FeatureControlActor, FeatureInput, FeatureOutput, FeatureSharedInput, FeatureWorker, FeatureWorkerInput, FeatureWorkerOutput, NetOutgoingMeta, Ttl};
13
14pub const FEATURE_ID: u8 = 6;
15pub const FEATURE_NAME: &str = "alias";
16pub const HINT_TIMEOUT_MS: u64 = 2000;
17pub const SCAN_TIMEOUT_MS: u64 = 5000;
18
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub enum Control {
21    Register { alias: u64, service: u8, level: ServiceBroadcastLevel },
22    Query { alias: u64, service: u8, level: ServiceBroadcastLevel },
23    Unregister { alias: u64 },
24}
25
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub enum FoundLocation {
28    Local,
29    Notify(NodeId),
30    CachedHint(NodeId),
31    RemoteHint(NodeId),
32    RemoteScan(NodeId),
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub enum Event {
37    QueryResult(u64, Option<FoundLocation>),
38}
39
40#[derive(Debug, PartialEq, Eq, Clone)]
41pub struct ToWorker;
42
43#[derive(Debug, PartialEq, Eq, Clone)]
44pub struct ToController;
45
46#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
47pub enum Message {
48    Notify(u64),
49    Scan(u64),
50    Check(u64),
51    Found(u64, bool),
52}
53
54#[derive(Debug)]
55enum QueryState {
56    CheckHint(NodeId, u64),
57    Scan(u64),
58}
59
60#[derive(Debug)]
61struct QuerySlot<UserData> {
62    waiters: Vec<FeatureControlActor<UserData>>,
63    state: QueryState,
64    service: u8,
65    level: ServiceBroadcastLevel,
66}
67
68#[derive(Debug, PartialEq, Eq)]
69struct HintSlot {
70    node: NodeId,
71    ts: u64,
72}
73
74pub type Output<UserData> = FeatureOutput<UserData, Event, ToWorker>;
75pub type WorkerOutput<UserData> = FeatureWorkerOutput<UserData, Control, Event, ToController>;
76
77#[derive(Debug, Derivative)]
78#[derivative(Default(bound = ""))]
79pub struct AliasFeature<UserData> {
80    queries: HashMap<u64, QuerySlot<UserData>>,
81    hint_slots: HashMap<u64, HintSlot>,
82    local_slots: HashMap<u64, u64>,
83    queue: VecDeque<Output<UserData>>,
84    scan_seq: u16,
85    shutdown: bool,
86}
87
88impl<UserData: Debug + Copy> AliasFeature<UserData> {
89    fn process_control(&mut self, now_ms: u64, actor: FeatureControlActor<UserData>, control: Control) {
90        match control {
91            Control::Register { alias, service, level } => {
92                log::info!("[AliasFeature] Register local alias {} and broadcast hint", alias);
93                self.local_slots.insert(alias, now_ms);
94                let seq = Self::gen_seq(&mut self.scan_seq);
95                Self::send_to(&mut self.queue, RouteRule::ToServices(service, level, seq), Message::Notify(alias));
96            }
97            Control::Query { alias, service, level } => {
98                if self.local_slots.contains_key(&alias) {
99                    log::debug!("[AliasFeature] Found alias {} at local", alias);
100                    self.queue.push_back(FeatureOutput::Event(actor, Event::QueryResult(alias, Some(FoundLocation::Local))));
101                } else if let Some(slot) = self.queries.get_mut(&alias) {
102                    log::debug!("[AliasFeature] Alias {} is already in query state => push to wait queue", alias);
103                    slot.waiters.push(actor);
104                } else if let Some(slot) = self.hint_slots.get(&alias) {
105                    if slot.ts + HINT_TIMEOUT_MS >= now_ms {
106                        log::debug!("[AliasFeature] Alias {alias} is very newly added ({} vs now {}) to hint {} => reuse", slot.ts, now_ms, slot.node);
107                        self.queue.push_back(FeatureOutput::Event(actor, Event::QueryResult(alias, Some(FoundLocation::CachedHint(slot.node)))));
108                    } else {
109                        log::debug!("[AliasFeature] Alias {alias} is not in query state but has hint {} => check hint", slot.node);
110                        self.queries.insert(
111                            alias,
112                            QuerySlot {
113                                waiters: vec![actor],
114                                state: QueryState::CheckHint(slot.node, now_ms),
115                                service,
116                                level,
117                            },
118                        );
119                        Self::send_to(&mut self.queue, RouteRule::ToNode(slot.node), Message::Check(alias));
120                    }
121                } else {
122                    log::debug!("[AliasFeature] Alias {alias} is not in query state and has no hint => scan");
123                    self.queries.insert(
124                        alias,
125                        QuerySlot {
126                            waiters: vec![actor],
127                            state: QueryState::Scan(now_ms),
128                            service,
129                            level,
130                        },
131                    );
132                    let seq = Self::gen_seq(&mut self.scan_seq);
133                    Self::send_to(&mut self.queue, RouteRule::ToServices(service, level, seq), Message::Scan(alias));
134                }
135            }
136            Control::Unregister { alias } => {
137                log::info!("[AliasFeature] Unregister alias {}", alias);
138                self.local_slots.remove(&alias);
139            }
140        }
141    }
142
143    fn process_remote(&mut self, now_ms: u64, from: NodeId, msg: Message) {
144        log::debug!("[AliasFeature] Received message from {from}: {:?}", msg);
145        match msg {
146            Message::Notify(alias) => {
147                self.hint_slots.insert(alias, HintSlot { node: from, ts: now_ms });
148                if let Some(slot) = self.queries.remove(&alias) {
149                    for actor in &slot.waiters {
150                        self.queue.push_back(FeatureOutput::Event(*actor, Event::QueryResult(alias, Some(FoundLocation::Notify(from)))));
151                    }
152                }
153            }
154            Message::Scan(alias) => {
155                if self.local_slots.contains_key(&alias) {
156                    log::debug!("[AliasFeature] Received Scan alias {alias}, found at local");
157                    Self::send_to(&mut self.queue, RouteRule::ToNode(from), Message::Found(alias, true));
158                } else {
159                    log::debug!("[AliasFeature] Received Scan alias {alias}, not found at local");
160                }
161            }
162            Message::Check(alias) => {
163                let found = self.local_slots.contains_key(&alias);
164                log::debug!("[AliasFeature] Received Check alias {alias}, found at local: {found}");
165                Self::send_to(&mut self.queue, RouteRule::ToNode(from), Message::Found(alias, found));
166            }
167            Message::Found(alias, found) => {
168                if found {
169                    self.hint_slots.insert(alias, HintSlot { node: from, ts: now_ms });
170                }
171                if let Some(slot) = self.queries.get_mut(&alias) {
172                    match slot.state {
173                        QueryState::CheckHint(node, _) => {
174                            if node != from {
175                                log::warn!("[AliasFeature] Reject Found message from wrong hint {node} vs {from}");
176                                return;
177                            }
178                            if found {
179                                log::debug!("[AliasFeature] Found alias {alias} at {node} => notify waiters {:?}", slot.waiters);
180                                for actor in &slot.waiters {
181                                    self.queue.push_back(FeatureOutput::Event(*actor, Event::QueryResult(alias, Some(FoundLocation::RemoteHint(from)))));
182                                }
183                                self.queries.remove(&alias);
184                            } else {
185                                log::debug!("[AliasFeature] Not found alias {alias} at hint {node} => switch to Scan");
186                                let seq = self.scan_seq;
187                                self.scan_seq = self.scan_seq.wrapping_add(1);
188                                slot.state = QueryState::Scan(now_ms);
189                                Self::send_to(&mut self.queue, RouteRule::ToServices(slot.service, slot.level, seq), Message::Scan(alias));
190                            }
191                        }
192                        QueryState::Scan(_) => {
193                            if !found {
194                                log::warn!("[AliasFeature] Remote should not reply with Found=false for Scan");
195                                return;
196                            }
197                            log::debug!("[AliasFeature] Found alias {alias} at {from} with Scan => notify waiters {:?}", slot.waiters);
198                            for actor in &slot.waiters {
199                                self.queue.push_back(FeatureOutput::Event(*actor, Event::QueryResult(alias, Some(FoundLocation::RemoteScan(from)))));
200                            }
201                            self.queries.remove(&alias);
202                        }
203                    }
204                }
205            }
206        }
207    }
208
209    fn send_to(queue: &mut VecDeque<FeatureOutput<UserData, Event, ToWorker>>, rule: RouteRule, msg: Message) {
210        let msg = bincode::serialize(&msg).expect("Should to bytes");
211        queue.push_back(FeatureOutput::SendRoute(rule, NetOutgoingMeta::new(true, Ttl::default(), 0, true), msg.into()));
212    }
213
214    fn gen_seq(scan_seq: &mut u16) -> u16 {
215        let seq = *scan_seq;
216        *scan_seq = scan_seq.wrapping_add(1);
217        seq
218    }
219}
220
221impl<UserData: Debug + Copy> Feature<UserData, Control, Event, ToController, ToWorker> for AliasFeature<UserData> {
222    fn on_shared_input(&mut self, _ctx: &FeatureContext, now: u64, input: FeatureSharedInput) {
223        if let FeatureSharedInput::Tick(_) = input {
224            let mut timeout = vec![];
225            for (alias, slot) in &mut self.queries {
226                match &slot.state {
227                    QueryState::CheckHint(hint, started_at) => {
228                        if now >= *started_at + HINT_TIMEOUT_MS {
229                            log::debug!("[AliasFeature] check {alias} hint node {hint} timeout => switch to Scan");
230
231                            let seq = self.scan_seq;
232                            self.scan_seq = self.scan_seq.wrapping_add(1);
233                            slot.state = QueryState::Scan(now);
234                            Self::send_to(&mut self.queue, RouteRule::ToServices(slot.service, slot.level, seq), Message::Scan(*alias));
235                        }
236                    }
237                    QueryState::Scan(started_at) => {
238                        if now >= *started_at + SCAN_TIMEOUT_MS {
239                            timeout.push(*alias);
240                        }
241                    }
242                }
243            }
244
245            for alias in timeout {
246                let slot = self.queries.remove(&alias).expect("Should have slot");
247                log::debug!("[AliasFeature] scan {alias} timeout => notify waiters {:?}", slot.waiters);
248                for actor in slot.waiters {
249                    self.queue.push_back(FeatureOutput::Event(actor, Event::QueryResult(alias, None)));
250                }
251            }
252        }
253    }
254
255    fn on_input(&mut self, _ctx: &FeatureContext, now_ms: u64, input: FeatureInput<'_, UserData, Control, ToController>) {
256        match input {
257            FeatureInput::Control(actor, control) => self.process_control(now_ms, actor, control),
258            FeatureInput::Local(meta, msg) | FeatureInput::Net(_, meta, msg) => {
259                if !meta.secure {
260                    log::warn!("[AliasFeature] reject unsecure message");
261                    return;
262                }
263                if let (Some(from), Ok(msg)) = (meta.source, bincode::deserialize::<Message>(&msg)) {
264                    self.process_remote(now_ms, from, msg)
265                }
266            }
267            _ => {}
268        }
269    }
270
271    fn on_shutdown(&mut self, _ctx: &FeatureContext, _now: u64) {
272        self.shutdown = true;
273    }
274}
275
276impl<UserData> TaskSwitcherChild<Output<UserData>> for AliasFeature<UserData> {
277    type Time = u64;
278
279    fn is_empty(&self) -> bool {
280        self.shutdown && self.queue.is_empty()
281    }
282
283    fn empty_event(&self) -> Output<UserData> {
284        Output::OnResourceEmpty
285    }
286
287    fn pop_output(&mut self, _now: u64) -> Option<Output<UserData>> {
288        self.queue.pop_front()
289    }
290}
291
292#[derive(Derivative)]
293#[derivative(Default(bound = ""))]
294pub struct AliasFeatureWorker<UserData> {
295    queue: DynamicDeque<WorkerOutput<UserData>, 1>,
296    shutdown: bool,
297}
298
299impl<UserData> FeatureWorker<UserData, Control, Event, ToController, ToWorker> for AliasFeatureWorker<UserData> {
300    fn on_input(&mut self, _ctx: &mut crate::base::FeatureWorkerContext, _now: u64, input: crate::base::FeatureWorkerInput<UserData, Control, ToWorker>) {
301        match input {
302            FeatureWorkerInput::Control(actor, control) => self.queue.push_back(FeatureWorkerOutput::ForwardControlToController(actor, control)),
303            FeatureWorkerInput::Network(conn, header, buf) => self.queue.push_back(FeatureWorkerOutput::ForwardNetworkToController(conn, header, buf)),
304            #[cfg(feature = "vpn")]
305            FeatureWorkerInput::TunPkt(..) => {}
306            FeatureWorkerInput::FromController(..) => {
307                log::warn!("No handler for FromController");
308            }
309            FeatureWorkerInput::Local(header, buf) => self.queue.push_back(FeatureWorkerOutput::ForwardLocalToController(header, buf)),
310        }
311    }
312
313    fn on_shutdown(&mut self, _ctx: &mut crate::base::FeatureWorkerContext, _now: u64) {
314        log::info!("[AliasFeatureWorker] Shutdown");
315        self.shutdown = true;
316    }
317}
318
319impl<UserData> TaskSwitcherChild<WorkerOutput<UserData>> for AliasFeatureWorker<UserData> {
320    type Time = u64;
321
322    fn is_empty(&self) -> bool {
323        self.shutdown && self.queue.is_empty()
324    }
325
326    fn empty_event(&self) -> WorkerOutput<UserData> {
327        WorkerOutput::OnResourceEmpty
328    }
329
330    fn pop_output(&mut self, _now: u64) -> Option<WorkerOutput<UserData>> {
331        self.queue.pop_front()
332    }
333}
334
335#[cfg(test)]
336mod tests {
337    use atm0s_sdn_router::{RouteRule, ServiceBroadcastLevel};
338    use sans_io_runtime::TaskSwitcherChild;
339
340    use crate::{
341        base::{Feature, FeatureContext, FeatureControlActor, FeatureInput, FeatureOutput, FeatureSharedInput},
342        features::alias::{HintSlot, HINT_TIMEOUT_MS, SCAN_TIMEOUT_MS},
343    };
344
345    use super::{AliasFeature, Control, Event, FoundLocation, Message, ToWorker};
346
347    fn decode_msg(msg: Option<FeatureOutput<(), Event, ToWorker>>) -> Option<(RouteRule, Message)> {
348        match msg? {
349            FeatureOutput::SendRoute(rule, _, msg) => Some((rule, bincode::deserialize(&msg).expect("Should decode"))),
350            _ => panic!("Should be SendRoute"),
351        }
352    }
353
354    #[test]
355    fn local_alias_simple() {
356        let mut alias = AliasFeature::default();
357        let ctx = FeatureContext { node_id: 0, session: 0 };
358        let service = 1;
359        let level = ServiceBroadcastLevel::Global;
360        alias.on_input(&ctx, 0, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Register { alias: 1000, service, level }));
361        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToServices(service, level, 0), Message::Notify(1000))));
362        assert_eq!(alias.pop_output(0), None);
363
364        alias.on_input(&ctx, 0, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }));
365        assert_eq!(
366            alias.pop_output(0),
367            Some(FeatureOutput::Event(FeatureControlActor::Controller(()), Event::QueryResult(1000, Some(FoundLocation::Local))))
368        );
369        assert_eq!(alias.pop_output(0), None);
370    }
371
372    #[test]
373    fn local_alias_handle_check() {
374        let mut alias = AliasFeature::default();
375        let ctx = FeatureContext { node_id: 0, session: 0 };
376        let service = 1;
377        let level = ServiceBroadcastLevel::Global;
378        alias.on_input(&ctx, 0, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Register { alias: 1000, service, level }));
379        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToServices(service, level, 0), Message::Notify(1000))));
380        assert_eq!(alias.pop_output(0), None);
381
382        alias.process_remote(0, 123, Message::Check(1000));
383        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToNode(123), Message::Found(1000, true))));
384        assert_eq!(alias.pop_output(0), None);
385
386        alias.process_remote(0, 123, Message::Check(1001));
387        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToNode(123), Message::Found(1001, false))));
388        assert_eq!(alias.pop_output(0), None);
389    }
390
391    #[test]
392    fn local_alias_handle_scan() {
393        let mut alias = AliasFeature::default();
394        let ctx = FeatureContext { node_id: 0, session: 0 };
395        let service = 1;
396        let level = ServiceBroadcastLevel::Global;
397        alias.on_input(&ctx, 0, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Register { alias: 1000, service, level }));
398        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToServices(service, level, 0), Message::Notify(1000))));
399        assert_eq!(alias.pop_output(0), None);
400
401        alias.process_remote(0, 123, Message::Scan(1000));
402        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToNode(123), Message::Found(1000, true))));
403        assert_eq!(alias.pop_output(0), None);
404
405        alias.process_remote(0, 123, Message::Scan(1001));
406        assert_eq!(alias.pop_output(0), None);
407    }
408
409    #[test]
410    fn found_cached_hint() {
411        let mut alias = AliasFeature::default();
412        let ctx = FeatureContext { node_id: 0, session: 0 };
413        let service = 1;
414        let level = ServiceBroadcastLevel::Global;
415
416        alias.hint_slots.insert(1000, HintSlot { node: 123, ts: 0 });
417
418        alias.on_input(
419            &ctx,
420            HINT_TIMEOUT_MS,
421            FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }),
422        );
423
424        assert_eq!(
425            alias.pop_output(HINT_TIMEOUT_MS),
426            Some(FeatureOutput::Event(
427                FeatureControlActor::Controller(()),
428                Event::QueryResult(1000, Some(FoundLocation::CachedHint(123)))
429            ))
430        );
431        assert_eq!(alias.pop_output(HINT_TIMEOUT_MS), None);
432    }
433
434    #[test]
435    fn found_remote_with_hint() {
436        let mut alias = AliasFeature::default();
437        let ctx = FeatureContext { node_id: 0, session: 0 };
438        let service = 1;
439        let level = ServiceBroadcastLevel::Global;
440
441        alias.hint_slots.insert(1000, HintSlot { node: 123, ts: 0 });
442
443        alias.on_input(&ctx, 10000, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }));
444        assert_eq!(decode_msg(alias.pop_output(10000)), Some((RouteRule::ToNode(123), Message::Check(1000))));
445        assert_eq!(alias.pop_output(10000), None);
446
447        //simulate remote found
448        alias.process_remote(10100, 123, Message::Found(1000, true));
449
450        assert_eq!(
451            alias.pop_output(10100),
452            Some(FeatureOutput::Event(
453                FeatureControlActor::Controller(()),
454                Event::QueryResult(1000, Some(FoundLocation::RemoteHint(123)))
455            ))
456        );
457        assert_eq!(alias.pop_output(10100), None);
458    }
459
460    #[test]
461    fn found_remote_with_scan() {
462        let mut alias = AliasFeature::default();
463        let ctx = FeatureContext { node_id: 0, session: 0 };
464        let service = 1;
465        let level = ServiceBroadcastLevel::Global;
466
467        alias.on_input(&ctx, 0, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }));
468        assert_eq!(decode_msg(alias.pop_output(0)), Some((RouteRule::ToServices(service, level, 0), Message::Scan(1000))));
469        assert_eq!(alias.pop_output(0), None);
470
471        //simulate scan found
472        alias.process_remote(100, 123, Message::Found(1000, true));
473
474        assert_eq!(
475            alias.pop_output(100),
476            Some(FeatureOutput::Event(
477                FeatureControlActor::Controller(()),
478                Event::QueryResult(1000, Some(FoundLocation::RemoteScan(123)))
479            ))
480        );
481        assert_eq!(alias.pop_output(100), None);
482    }
483
484    #[test]
485    fn found_remote_with_hint_then_scan_fallback() {
486        let mut alias = AliasFeature::default();
487        let ctx = FeatureContext { node_id: 0, session: 0 };
488        let service = 1;
489        let level = ServiceBroadcastLevel::Global;
490
491        alias.hint_slots.insert(1000, HintSlot { node: 122, ts: 0 });
492
493        alias.on_input(&ctx, 10000, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }));
494        assert_eq!(decode_msg(alias.pop_output(10000)), Some((RouteRule::ToNode(122), Message::Check(1000))));
495        assert_eq!(alias.pop_output(10000), None);
496
497        //simulate remote not found
498        alias.process_remote(10100, 122, Message::Found(1000, false));
499
500        // will fallback to scan
501        assert_eq!(decode_msg(alias.pop_output(10100)), Some((RouteRule::ToServices(service, level, 0), Message::Scan(1000))));
502        assert_eq!(alias.pop_output(10100), None);
503
504        //simulate scan found
505        alias.process_remote(10100, 123, Message::Found(1000, true));
506
507        assert_eq!(
508            alias.pop_output(10100),
509            Some(FeatureOutput::Event(
510                FeatureControlActor::Controller(()),
511                Event::QueryResult(1000, Some(FoundLocation::RemoteScan(123)))
512            ))
513        );
514        assert_eq!(alias.pop_output(10100), None);
515    }
516
517    #[test]
518    fn found_remote_with_hint_timeout_then_scan_fallback() {
519        let mut alias = AliasFeature::default();
520        let ctx = FeatureContext { node_id: 0, session: 0 };
521        let service = 1;
522        let level = ServiceBroadcastLevel::Global;
523
524        alias.hint_slots.insert(1000, HintSlot { node: 122, ts: 0 });
525
526        alias.on_input(&ctx, 10000, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }));
527        assert_eq!(decode_msg(alias.pop_output(10000)), Some((RouteRule::ToNode(122), Message::Check(1000))));
528        assert_eq!(alias.pop_output(10000), None);
529
530        //simulate remote not found
531        alias.on_shared_input(&ctx, 10000 + HINT_TIMEOUT_MS, FeatureSharedInput::Tick(0));
532
533        // will fallback to scan
534        assert_eq!(
535            decode_msg(alias.pop_output(10000 + HINT_TIMEOUT_MS)),
536            Some((RouteRule::ToServices(service, level, 0), Message::Scan(1000)))
537        );
538        assert_eq!(alias.pop_output(10000 + HINT_TIMEOUT_MS), None);
539
540        //simulate scan found
541        alias.process_remote(10100 + HINT_TIMEOUT_MS, 123, Message::Found(1000, true));
542
543        assert_eq!(
544            alias.pop_output(10100 + HINT_TIMEOUT_MS),
545            Some(FeatureOutput::Event(
546                FeatureControlActor::Controller(()),
547                Event::QueryResult(1000, Some(FoundLocation::RemoteScan(123)))
548            ))
549        );
550        assert_eq!(alias.pop_output(10100 + HINT_TIMEOUT_MS), None);
551
552        //after that hint should be saved
553        assert_eq!(
554            alias.hint_slots.get(&1000),
555            Some(&HintSlot {
556                node: 123,
557                ts: 10100 + HINT_TIMEOUT_MS
558            })
559        );
560    }
561
562    #[test]
563    fn timeout_both_hint_and_scan() {
564        let mut alias = AliasFeature::default();
565        let ctx = FeatureContext { node_id: 0, session: 0 };
566        let service = 1;
567        let level = ServiceBroadcastLevel::Global;
568
569        alias.hint_slots.insert(1000, HintSlot { node: 122, ts: 0 });
570
571        alias.on_input(&ctx, 10000, FeatureInput::Control(FeatureControlActor::Controller(()), Control::Query { alias: 1000, service, level }));
572        assert_eq!(decode_msg(alias.pop_output(10000)), Some((RouteRule::ToNode(122), Message::Check(1000))));
573        assert_eq!(alias.pop_output(10000), None);
574
575        //simulate remote not found
576        alias.on_shared_input(&ctx, 10000 + HINT_TIMEOUT_MS, FeatureSharedInput::Tick(0));
577
578        // will fallback to scan
579        assert_eq!(
580            decode_msg(alias.pop_output(10000 + HINT_TIMEOUT_MS)),
581            Some((RouteRule::ToServices(service, level, 0), Message::Scan(1000)))
582        );
583        assert_eq!(alias.pop_output(10000 + HINT_TIMEOUT_MS), None);
584
585        //simulate scan found
586        alias.on_shared_input(&ctx, 10000 + HINT_TIMEOUT_MS + SCAN_TIMEOUT_MS, FeatureSharedInput::Tick(1));
587
588        assert_eq!(
589            alias.pop_output(10000 + HINT_TIMEOUT_MS + SCAN_TIMEOUT_MS),
590            Some(FeatureOutput::Event(FeatureControlActor::Controller(()), Event::QueryResult(1000, None)))
591        );
592        assert_eq!(alias.pop_output(10000 + HINT_TIMEOUT_MS + SCAN_TIMEOUT_MS), None);
593    }
594
595    #[test]
596    fn handle_notify_from_remote() {
597        let mut alias = AliasFeature::<()>::default();
598        alias.process_remote(100, 123, Message::Notify(1000));
599        assert_eq!(alias.hint_slots.get(&1000), Some(&HintSlot { node: 123, ts: 100 }));
600    }
601}