wowlab_node/core/
connection.rs1use 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 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}