Skip to main content

wowlab_node/core/
connection.rs

1use tokio::sync::mpsc;
2use tokio_util::sync::CancellationToken;
3
4use super::{NodeCore, NodeCoreEvent, RegisterResult, VerifyError};
5use crate::{
6    NodeState, claim,
7    realtime::{NodeRealtime, RealtimeConfig},
8};
9
10impl NodeCore {
11    pub(super) fn start_verification(&mut self) {
12        let (tx, rx) = mpsc::channel(1);
13
14        self.verify_rx = Some(rx);
15        let sentinel = self.sentinel.clone();
16
17        tracing::info!("Verifying node...");
18
19        self.runtime.spawn(async move {
20            let result: Result<wowlab_types::sensitive::Sensitive<String>, VerifyError> =
21                match sentinel.refresh_token().await {
22                    Ok(token) => Ok(token),
23                    Err(e) => {
24                        let err_str = e.to_string().to_lowercase();
25
26                        tracing::debug!(error = %err_str, "Node verification failed");
27
28                        if err_str.contains("not found") || err_str.contains("not claimed") {
29                            Err(VerifyError::NotFound)
30                        } else {
31                            Err(VerifyError::Unavailable)
32                        }
33                    }
34                };
35            let _ = tx.send(result).await;
36        });
37    }
38
39    pub(super) fn check_verification(&mut self) {
40        let Some(ref mut rx) = self.verify_rx else {
41            return;
42        };
43
44        match rx.try_recv() {
45            Ok(Ok(_token)) => {
46                tracing::info!("Node verified, starting...");
47                self.verify_rx = None;
48                self.backoff.reset();
49                self.registered = true;
50                self.set_state(NodeState::Running);
51                self.start_realtime();
52            }
53            Ok(Err(VerifyError::NotFound)) => {
54                tracing::warn!("Node not found, needs re-registration");
55                self.verify_rx = None;
56                self.registered = false;
57                self.set_state(NodeState::NotFound);
58            }
59            Ok(Err(VerifyError::Unavailable)) => {
60                // #t(rust_hardcoded_url) static status page URL in user-facing log message
61                tracing::error!("Server unavailable, please check https://wowlab.gg/status");
62                self.verify_rx = None;
63                self.set_unavailable();
64            }
65            Err(mpsc::error::TryRecvError::Empty) => {}
66            Err(mpsc::error::TryRecvError::Disconnected) => {
67                self.verify_rx = None;
68            }
69        }
70    }
71
72    pub(super) fn check_retry(&mut self) {
73        if !matches!(self.state, NodeState::Unavailable) {
74            return;
75        }
76
77        if !self.backoff.should_retry() {
78            return;
79        }
80
81        self.backoff.on_retry();
82
83        tracing::info!("Retrying connection...");
84
85        if self.registered {
86            self.set_state(NodeState::Verifying);
87            self.start_verification();
88        } else {
89            self.set_state(NodeState::Registering);
90            self.start_registration();
91        }
92    }
93
94    pub(super) fn start_registration(&mut self) {
95        let (tx, rx) = mpsc::channel(1);
96
97        self.register_rx = Some(rx);
98        let sentinel = self.sentinel.clone();
99        let config = self.config.clone();
100
101        self.runtime.spawn(async move {
102            let Some(token_claim) = config.token_claim.as_ref() else {
103                let _ = tx
104                    .send(RegisterResult::Failed("Missing claim token".to_string()))
105                    .await;
106
107                return;
108            };
109
110            match claim::register(&sentinel, &config, token_claim.expose()).await {
111                Ok(()) => {
112                    let _ = tx.send(RegisterResult::Success).await;
113                }
114                Err(error) => {
115                    tracing::error!(%error, "Registration failed");
116                    let _ = tx.send(RegisterResult::Failed(error.to_string())).await;
117                }
118            }
119        });
120    }
121
122    pub(super) fn start_realtime(&mut self) {
123        if let Some(token) = self.realtime_shutdown.take() {
124            token.cancel();
125        }
126
127        let shutdown = CancellationToken::new();
128        let sentinel = self.sentinel.clone();
129        let beacon_url = self.config.beacon_url.clone();
130        let app_name = self.app_name.clone();
131        let app_version = self.app_version.clone();
132        let public_key = self.config.public_key.clone();
133        let handle = self.runtime.handle().clone();
134        let (tx, rx) = mpsc::channel(super::EVENT_CHANNEL_SIZE);
135
136        self.realtime_rx = Some(rx);
137        self.realtime_shutdown = Some(shutdown.clone());
138
139        self.runtime.spawn(async move {
140            let beacon_token = match sentinel.refresh_token().await {
141                Ok(token) => token,
142                Err(error) => {
143                    tracing::error!(%error, "Failed to get beacon token");
144                    let _ = tx
145                        .send(crate::realtime::RealtimeEvent::Error(error.to_string()))
146                        .await;
147
148                    return;
149                }
150            };
151
152            let realtime = NodeRealtime::new(RealtimeConfig {
153                url: beacon_url,
154                token: beacon_token,
155                name: app_name,
156                version: app_version,
157                sentinel,
158            });
159
160            let mut inner_rx = realtime.subscribe(&public_key, &handle, shutdown);
161
162            while let Some(event) = inner_rx.recv().await {
163                if tx.send(event).await.is_err() {
164                    break;
165                }
166            }
167        });
168    }
169
170    pub(super) fn check_registration(&mut self) {
171        let Some(ref mut rx) = self.register_rx else {
172            return;
173        };
174
175        match rx.try_recv() {
176            Ok(RegisterResult::Success) => {
177                self.registered = true;
178                self.backoff.reset();
179                tracing::info!("Registered node, now running");
180                self.register_rx = None;
181                self.set_state(NodeState::Running);
182                self.start_realtime();
183            }
184            Ok(RegisterResult::Failed(ref err)) => {
185                tracing::error!(error = %err, "Registration failed");
186                self.register_rx = None;
187                let _ = self.event_tx.try_send(NodeCoreEvent::Error(err.clone()));
188
189                if err.contains("Invalid token") || err.contains("401") {
190                    self.config.token_claim = None;
191                    self.set_state(NodeState::Setup);
192                } else {
193                    self.set_unavailable();
194                }
195            }
196            Err(mpsc::error::TryRecvError::Empty) => {}
197            Err(mpsc::error::TryRecvError::Disconnected) => {
198                self.register_rx = None;
199            }
200        }
201    }
202
203    fn set_unavailable(&mut self) {
204        let _ = self.backoff.next_retry_at();
205
206        tracing::info!(
207            retry_seconds = self.backoff.current_backoff().as_secs(),
208            "Scheduled connection retry"
209        );
210        self.set_state(NodeState::Unavailable);
211    }
212}