1use std::sync::Arc;
14
15use dashmap::DashMap;
16use tokio::io::{AsyncBufReadExt, AsyncWrite, AsyncWriteExt, BufReader};
17use tokio::net::TcpListener;
18use tokio::sync::watch;
19
20use crate::conn::ClientConn;
21use crate::listener_status::{
22 ListenerBindEventSender, ListenerSurface, notify_listener_failed, notify_listener_opened,
23};
24use crate::metrics::GatewayMetrics;
25use crate::subscription::SubscriptionManager;
26
27pub type ReloginFn = Arc<dyn Fn() + Send + Sync>;
40
41pub struct TelnetServer {
43 listen_addr: String,
44 connections: Arc<DashMap<u64, ClientConn>>,
45 subscriptions: Option<Arc<SubscriptionManager>>,
46 metrics: Option<Arc<GatewayMetrics>>,
47 shutdown_tx: watch::Sender<bool>,
48 relogin_fn: Option<ReloginFn>,
51}
52
53const MAX_TELNET_CONNECTIONS: usize = 5;
54const VERSION: &str = env!("CARGO_PKG_VERSION");
55
56async fn write_telnet<W>(
57 writer: &mut W,
58 peer_addr: std::net::SocketAddr,
59 context: &'static str,
60 bytes: &[u8],
61) -> bool
62where
63 W: AsyncWrite + Unpin,
64{
65 match writer.write_all(bytes).await {
66 Ok(()) => true,
67 Err(err) => {
68 tracing::debug!(
69 peer = %peer_addr,
70 context,
71 error = %err,
72 "telnet write failed; client likely disconnected"
73 );
74 false
75 }
76 }
77}
78
79impl TelnetServer {
80 pub fn new(
81 listen_addr: String,
82 connections: Arc<DashMap<u64, ClientConn>>,
83 subscriptions: Option<Arc<SubscriptionManager>>,
84 metrics: Option<Arc<GatewayMetrics>>,
85 shutdown_tx: watch::Sender<bool>,
86 ) -> Self {
87 Self {
88 listen_addr,
89 connections,
90 subscriptions,
91 metrics,
92 shutdown_tx,
93 relogin_fn: None,
94 }
95 }
96
97 pub fn with_relogin_fn(mut self, relogin_fn: ReloginFn) -> Self {
99 self.relogin_fn = Some(relogin_fn);
100 self
101 }
102
103 pub async fn run(&self) -> anyhow::Result<()> {
104 let (_shutdown_tx, shutdown_rx) = watch::channel(false);
105 self.run_until_shutdown(shutdown_rx).await
106 }
107
108 pub async fn run_until_shutdown(
109 &self,
110 shutdown_rx: watch::Receiver<bool>,
111 ) -> anyhow::Result<()> {
112 self.run_until_shutdown_with_listener_events(shutdown_rx, None)
113 .await
114 }
115
116 pub async fn run_until_shutdown_with_listener_events(
118 &self,
119 mut shutdown_rx: watch::Receiver<bool>,
120 listener_events: Option<ListenerBindEventSender>,
121 ) -> anyhow::Result<()> {
122 let listener = TcpListener::bind(&self.listen_addr)
123 .await
124 .map_err(|error| {
125 notify_listener_failed(&listener_events, ListenerSurface::Telnet);
126 anyhow::Error::new(crate::bind_hint::io_bind_error(
127 "Telnet",
128 "--telnet-port",
129 &self.listen_addr,
130 error,
131 ))
132 })?;
133 let _serving =
134 notify_listener_opened(&listener_events, ListenerSurface::Telnet, &shutdown_rx).await?;
135 drop(listener_events);
136 tracing::info!(addr = %self.listen_addr, "Telnet server listening");
137
138 let active = Arc::new(std::sync::atomic::AtomicUsize::new(0));
139
140 loop {
141 let (stream, peer_addr) = tokio::select! {
142 _ = crate::listener::shutdown_requested(&mut shutdown_rx) => {
143 tracing::info!("Telnet server accept loop stopped by shutdown signal");
144 break;
145 }
146 accepted = listener.accept() => accepted?,
147 };
148 let count = active.load(std::sync::atomic::Ordering::Relaxed);
149 if count >= MAX_TELNET_CONNECTIONS {
150 tracing::warn!(peer = %peer_addr, "telnet max connections reached");
151 drop(stream);
152 continue;
153 }
154
155 active.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
156 let active_clone = Arc::clone(&active);
157 let connections = Arc::clone(&self.connections);
158 let subscriptions = self.subscriptions.clone();
159 let metrics = self.metrics.clone();
160 let shutdown_tx = self.shutdown_tx.clone();
161 let relogin_fn = self.relogin_fn.clone();
163
164 tokio::spawn(async move {
165 tracing::info!(peer = %peer_addr, "telnet client connected");
166 let (reader, mut writer) = stream.into_split();
167 let mut reader = BufReader::new(reader);
168
169 if !write_telnet(
170 &mut writer,
171 peer_addr,
172 "greeting",
173 b"FutuOpenD Rust Gateway Telnet Console\r\nType 'help' for commands.\r\n> ",
174 )
175 .await
176 {
177 active_clone.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
178 tracing::info!(peer = %peer_addr, "telnet client disconnected");
179 return;
180 }
181
182 let mut line = String::new();
183 loop {
184 line.clear();
185 match reader.read_line(&mut line).await {
186 Ok(0) => break, Ok(_) => {}
188 Err(_) => break,
189 }
190
191 let cmd = line.trim().to_lowercase();
192 let response = match cmd.as_str() {
193 "help" => "Available commands:\r\n\
194 help — show this help\r\n\
195 ping — connectivity test\r\n\
196 version — show version\r\n\
197 show_subinfo — show subscription info\r\n\
198 show_conn — show active connections\r\n\
199 show_metrics — show gateway metrics\r\n\
200 set_loglevel — change log level (debug/info/warn/error)\r\n\
201 relogin — clear login_cache; next P1-D tick triggers relogin (v1.4.97)\r\n\
202 exit — shutdown gateway\r\n"
203 .to_string(),
204 "ping" => "pong\r\n".to_string(),
205 "version" => format!("futu-opend-rs v{VERSION}\r\n"),
206 "show_subinfo" => {
207 if let Some(ref subs) = subscriptions {
208 let total_quota = subs.get_total_used_quota();
209 format!("Total subscription quota used: {total_quota}/4000\r\n")
210 } else {
211 "Subscription manager not available\r\n".to_string()
212 }
213 }
214 "show_conn" => {
215 let count = connections.len();
216 let mut info = format!("Active connections: {count}\r\n");
217 for entry in connections.iter() {
218 let c = entry.value();
219 info += &format!(" conn_id={} state={:?}\r\n", c.conn_id, c.state);
220 }
221 info
222 }
223 "show_metrics" => {
224 if let Some(ref m) = metrics {
225 m.report()
226 } else {
227 "Metrics not available\r\n".to_string()
228 }
229 }
230 s if s.starts_with("set_loglevel") => {
231 let parts: Vec<&str> = s.split_whitespace().collect();
232 if parts.len() < 2 {
233 "Usage: set_loglevel <debug|info|warn|error>\r\n".to_string()
234 } else {
235 let level = parts[1];
238 tracing::info!(level, "log level change requested via telnet");
239 format!(
240 "Log level change to '{level}' noted (restart for full effect)\r\n"
241 )
242 }
243 }
244 "relogin" => {
245 if let Some(ref f) = relogin_fn {
254 tracing::info!(peer = %peer_addr,
255 "v1.4.97 P1-D-F: relogin requested via telnet");
256 f();
257 "Relogin triggered. Watch daemon log for \
258 \"P1-D ladder fired\" + \"relogin succeeded\" \
259 in next 30s tick.\r\n"
260 .to_string()
261 } else {
262 "Relogin callback not wired (daemon was started \
263 without bridge handle). Relogin requires daemon \
264 restart in this configuration.\r\n"
265 .to_string()
266 }
267 }
268 "exit" => {
269 write_telnet(
270 &mut writer,
271 peer_addr,
272 "shutdown_notice",
273 b"Shutting down...\r\n",
274 )
275 .await;
276 if shutdown_tx.send(true).is_err() {
277 tracing::warn!(
278 peer = %peer_addr,
279 "telnet exit requested but shutdown receiver dropped"
280 );
281 }
282 break;
283 }
284 "" => String::new(),
285 _ => format!(
286 "Unknown command: '{cmd}'. Type 'help' for available commands.\r\n"
287 ),
288 };
289
290 if !response.is_empty()
291 && !write_telnet(&mut writer, peer_addr, "response", response.as_bytes())
292 .await
293 {
294 break;
295 }
296 if !write_telnet(&mut writer, peer_addr, "prompt", b"> ").await {
297 break;
298 }
299 }
300
301 active_clone.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
302 tracing::info!(peer = %peer_addr, "telnet client disconnected");
303 });
304 }
305
306 Ok(())
307 }
308}