1mod combo;
4mod indicator;
5mod market;
6mod option_analytics;
7mod qot_3401_plus;
8mod reference;
9mod reference_corporate_short_brokers;
10mod reference_f10;
11mod reference_price_reminder;
12mod reference_screen_unusual;
13mod reference_shareholders_insider;
14mod search;
15mod subscription;
16mod system;
17mod trade_read;
18mod trade_unlock;
19mod trade_write;
20
21use crate::state::ServerState;
22use std::borrow::Cow;
23use std::sync::{
24 Arc,
25 atomic::{AtomicU8, Ordering},
26};
27
28#[derive(Clone)]
31pub struct FutuServer {
32 pub state: ServerState,
33 legacy_log_level: Arc<AtomicU8>,
34 legacy_push_service_lease: Arc<crate::state::LegacyPushServiceLease>,
35}
36
37impl FutuServer {
38 #[allow(deprecated)]
39 pub fn new(state: ServerState) -> Self {
40 Self {
41 legacy_push_service_lease: Arc::new(crate::state::LegacyPushServiceLease::new(
42 state.clone(),
43 )),
44 state,
45 legacy_log_level: Arc::new(AtomicU8::new(crate::state::LEGACY_PUSH_INFO_LEVEL_RANK)),
46 }
47 }
48
49 pub(crate) fn legacy_push_delivery(
50 &self,
51 peer: rmcp::service::Peer<rmcp::RoleServer>,
52 ) -> crate::state::PushDeliveryTarget {
53 crate::state::PushDeliveryTarget::LegacyPeer(peer, Arc::clone(&self.legacy_log_level))
54 }
55
56 pub(crate) async fn register_legacy_push_subscriber(
57 &self,
58 delivery: crate::state::PushDeliveryTarget,
59 acc_ids: std::collections::HashSet<u64>,
60 owner_key_id: Option<String>,
61 allowed_acc_ids_snapshot: Option<std::collections::HashSet<u64>>,
62 allowed_markets_snapshot: Option<std::collections::HashSet<String>>,
63 ) -> Result<String, String> {
64 self.legacy_push_service_lease
65 .register_push_subscriber(
66 delivery,
67 acc_ids,
68 owner_key_id,
69 allowed_acc_ids_snapshot,
70 allowed_markets_snapshot,
71 )
72 .await
73 }
74}
75
76#[allow(deprecated)]
77fn legacy_logging_level_rank(level: rmcp::model::LoggingLevel) -> u8 {
78 match level {
79 rmcp::model::LoggingLevel::Debug => 0,
80 rmcp::model::LoggingLevel::Info => 1,
81 rmcp::model::LoggingLevel::Notice => 2,
82 rmcp::model::LoggingLevel::Warning => 3,
83 rmcp::model::LoggingLevel::Error => 4,
84 rmcp::model::LoggingLevel::Critical => 5,
85 rmcp::model::LoggingLevel::Alert => 6,
86 rmcp::model::LoggingLevel::Emergency => 7,
87 }
88}
89
90fn server_implementation() -> rmcp::model::Implementation {
91 rmcp::model::Implementation::new("futu-mcp", env!("CARGO_PKG_VERSION"))
92}
93
94#[allow(deprecated)]
95fn legacy_server_info() -> rmcp::model::ServerInfo {
96 rmcp::model::ServerInfo::new(
97 rmcp::model::ServerCapabilities::builder()
98 .enable_tools()
99 .enable_logging()
100 .build(),
101 )
102 .with_server_info(server_implementation())
103}
104
105fn modern_server_info(resources_enabled: bool) -> rmcp::model::ServerInfo {
106 let capabilities = if resources_enabled {
107 rmcp::model::ServerCapabilities::builder()
108 .enable_tools()
109 .enable_resources()
110 .enable_resources_subscribe()
111 .build()
112 } else {
113 rmcp::model::ServerCapabilities::builder()
114 .enable_tools()
115 .build()
116 };
117 rmcp::model::ServerInfo::new(capabilities).with_server_info(server_implementation())
118}
119
120#[allow(deprecated)]
121fn server_info() -> rmcp::model::ServerInfo {
122 rmcp::model::ServerInfo::new(
123 rmcp::model::ServerCapabilities::builder()
124 .enable_tools()
125 .enable_logging()
126 .enable_resources()
127 .enable_resources_subscribe()
128 .build(),
129 )
130 .with_server_info(server_implementation())
131}
132
133include!(concat!(env!("OUT_DIR"), "/generated_mcp_catalog.rs"));
134
135#[rmcp::tool_handler(router = Self::tool_router())]
136impl rmcp::ServerHandler for FutuServer {
137 fn get_info(&self) -> rmcp::model::ServerInfo {
138 server_info()
139 }
140
141 fn supported_protocol_versions(&self) -> Cow<'static, [rmcp::model::ProtocolVersion]> {
142 Cow::Borrowed(rmcp::model::ProtocolVersion::KNOWN_VERSIONS)
143 }
144
145 async fn initialize(
146 &self,
147 request: rmcp::model::InitializeRequestParams,
148 context: rmcp::service::RequestContext<rmcp::RoleServer>,
149 ) -> Result<rmcp::model::InitializeResult, rmcp::ErrorData> {
150 if request.protocol_version == rmcp::model::ProtocolVersion::V_2026_07_28 {
151 return Err(rmcp::ErrorData::method_not_found::<
152 rmcp::model::InitializeResultMethod,
153 >());
154 }
155 context.peer.set_peer_info(request.clone());
156 let supported = self.supported_protocol_versions();
157 let negotiated = if supported.contains(&request.protocol_version) {
158 request.protocol_version
159 } else {
160 rmcp::model::ProtocolVersion::LATEST
161 };
162 Ok(legacy_server_info().with_protocol_version(negotiated))
163 }
164
165 async fn discover(
166 &self,
167 context: rmcp::service::RequestContext<rmcp::RoleServer>,
168 ) -> Result<rmcp::model::DiscoverResult, rmcp::ErrorData> {
169 Ok(rmcp::model::DiscoverResult::from_server_info(
170 self.supported_protocol_versions().into_owned(),
171 modern_server_info(self.has_reusable_push_identity(&context)),
172 ))
173 }
174
175 #[allow(deprecated)]
176 async fn set_level(
177 &self,
178 request: rmcp::model::SetLevelRequestParams,
179 context: rmcp::service::RequestContext<rmcp::RoleServer>,
180 ) -> Result<(), rmcp::ErrorData> {
181 if context.protocol_version() == Some(rmcp::model::ProtocolVersion::V_2026_07_28) {
182 return Err(rmcp::ErrorData::method_not_found::<
183 rmcp::model::SetLevelRequestMethod,
184 >());
185 }
186 self.legacy_log_level
187 .store(legacy_logging_level_rank(request.level), Ordering::Release);
188 Ok(())
189 }
190
191 async fn list_prompts(
192 &self,
193 _request: Option<rmcp::model::PaginatedRequestParams>,
194 _context: rmcp::service::RequestContext<rmcp::RoleServer>,
195 ) -> Result<rmcp::model::ListPromptsResult, rmcp::ErrorData> {
196 Err(rmcp::ErrorData::method_not_found::<
200 rmcp::model::ListPromptsRequestMethod,
201 >())
202 }
203
204 async fn list_resources(
205 &self,
206 _request: Option<rmcp::model::PaginatedRequestParams>,
207 context: rmcp::service::RequestContext<rmcp::RoleServer>,
208 ) -> Result<rmcp::model::ListResourcesResult, rmcp::ErrorData> {
209 if context.protocol_version() != Some(rmcp::model::ProtocolVersion::V_2026_07_28) {
210 return Err(rmcp::ErrorData::method_not_found::<
211 rmcp::model::ListResourcesRequestMethod,
212 >());
213 }
214 let caller = self
215 .require_push_continuation("resources/list", &context)
216 .map_err(|_| {
217 rmcp::ErrorData::method_not_found::<rmcp::model::ListResourcesRequestMethod>()
218 })?;
219 let key_id = caller.key_id.as_deref().ok_or_else(|| {
220 rmcp::ErrorData::invalid_params(
221 "resources/list: reusable caller identity required",
222 None,
223 )
224 })?;
225 let allowed_markets = caller
226 .rec
227 .as_ref()
228 .and_then(|rec| rec.allowed_markets.as_ref());
229 let resources = self
230 .state
231 .modern_resources_for_caller(key_id, caller.allowed_acc_ids.as_ref(), allowed_markets)
232 .await;
233 Ok(rmcp::model::ListResourcesResult::with_all_items(resources)
234 .with_ttl_ms(0)
235 .with_cache_scope(rmcp::model::CacheScope::Private))
236 }
237
238 async fn list_resource_templates(
239 &self,
240 _request: Option<rmcp::model::PaginatedRequestParams>,
241 context: rmcp::service::RequestContext<rmcp::RoleServer>,
242 ) -> Result<rmcp::model::ListResourceTemplatesResult, rmcp::ErrorData> {
243 if context.protocol_version() != Some(rmcp::model::ProtocolVersion::V_2026_07_28) {
244 return Err(rmcp::ErrorData::method_not_found::<
245 rmcp::model::ListResourceTemplatesRequestMethod,
246 >());
247 }
248 self.require_push_continuation("resources/templates/list", &context)
249 .map_err(|_| {
250 rmcp::ErrorData::method_not_found::<
251 rmcp::model::ListResourceTemplatesRequestMethod,
252 >()
253 })?;
254 Ok(rmcp::model::ListResourceTemplatesResult::default()
255 .with_ttl_ms(0)
256 .with_cache_scope(rmcp::model::CacheScope::Public))
257 }
258
259 async fn read_resource(
260 &self,
261 request: rmcp::model::ReadResourceRequestParams,
262 context: rmcp::service::RequestContext<rmcp::RoleServer>,
263 ) -> Result<rmcp::model::ReadResourceResponse, rmcp::ErrorData> {
264 if context.protocol_version() != Some(rmcp::model::ProtocolVersion::V_2026_07_28) {
265 return Err(rmcp::ErrorData::method_not_found::<
266 rmcp::model::ReadResourceRequestMethod,
267 >());
268 }
269 let caller = self
270 .require_push_continuation("resources/read", &context)
271 .map_err(|_| {
272 rmcp::ErrorData::method_not_found::<rmcp::model::ReadResourceRequestMethod>()
273 })?;
274 let key_id = caller.key_id.as_deref().ok_or_else(|| {
275 rmcp::ErrorData::invalid_params(
276 "resources/read: reusable caller identity required",
277 None,
278 )
279 })?;
280 let allowed_markets = caller
281 .rec
282 .as_ref()
283 .and_then(|rec| rec.allowed_markets.as_ref());
284 let (session_id, events, dropped_events) = self
285 .state
286 .drain_modern_push_resource(
287 &request.uri,
288 key_id,
289 caller.allowed_acc_ids.as_ref(),
290 allowed_markets,
291 )
292 .await
293 .map_err(|error| rmcp::ErrorData::invalid_params(error, None))?;
294 let text = serde_json::json!({
295 "session_id": session_id,
296 "events": events,
297 "dropped_events": dropped_events,
298 })
299 .to_string();
300 Ok(
301 rmcp::model::ReadResourceResult::new(vec![rmcp::model::ResourceContents::text(
302 text,
303 request.uri,
304 )])
305 .with_ttl_ms(0)
306 .with_cache_scope(rmcp::model::CacheScope::Private)
307 .into(),
308 )
309 }
310
311 fn accepted_subscription_filter(
312 &self,
313 requested: &rmcp::model::SubscriptionFilter,
314 ) -> Option<rmcp::model::SubscriptionFilter> {
315 let Some([uri]) = requested.resource_subscriptions.as_deref() else {
316 return Some(rmcp::model::SubscriptionFilter::new());
317 };
318 Some(match crate::state::parse_push_resource_uri(uri) {
319 Some(_) => rmcp::model::SubscriptionFilter::builder()
320 .resource_subscription(uri.clone())
321 .build(),
322 None => rmcp::model::SubscriptionFilter::new(),
323 })
324 }
325
326 async fn listen(
327 &self,
328 context: rmcp::service::SubscriptionContext,
329 ) -> Result<(), rmcp::ErrorData> {
330 if context.request_context().protocol_version()
331 != Some(rmcp::model::ProtocolVersion::V_2026_07_28)
332 {
333 return Err(rmcp::ErrorData::method_not_found::<
334 rmcp::model::SubscriptionsListenRequestMethod,
335 >());
336 }
337 let caller = self
338 .require_push_continuation("subscriptions/listen", context.request_context())
339 .map_err(|_| {
340 rmcp::ErrorData::method_not_found::<rmcp::model::SubscriptionsListenRequestMethod>()
341 })?;
342 if context.accepted() == &rmcp::model::SubscriptionFilter::new() {
343 context.cancelled().await;
344 return Ok(());
345 }
346 let [uri] = context
347 .accepted()
348 .resource_subscriptions
349 .as_deref()
350 .unwrap_or_default()
351 else {
352 return Err(rmcp::ErrorData::invalid_params(
353 "subscriptions/listen requires exactly one canonical futu push resource URI",
354 None,
355 ));
356 };
357 let key_id = caller.key_id.as_deref().ok_or_else(|| {
358 rmcp::ErrorData::invalid_params(
359 "subscriptions/listen: reusable caller identity required",
360 None,
361 )
362 })?;
363 let allowed_markets = caller
364 .rec
365 .as_ref()
366 .and_then(|rec| rec.allowed_markets.as_ref());
367 let (listener_token, replaced_or_removed) = self
368 .state
369 .attach_modern_push_listener(
370 uri,
371 key_id,
372 caller.allowed_acc_ids.as_ref(),
373 allowed_markets,
374 context.sink().clone(),
375 )
376 .await
377 .map_err(|error| rmcp::ErrorData::invalid_params(error, None))?;
378
379 tokio::select! {
380 () = context.cancelled() => {}
381 () = replaced_or_removed.cancelled() => {}
382 }
383 self.state
384 .detach_modern_push_listener(uri, listener_token)
385 .await;
386 Ok(())
387 }
388}
389
390#[cfg(test)]
391mod tests;