1mod card_num_expand;
14mod guard;
15mod handlers;
16mod http;
17mod qot_sdk_adapter;
18mod state;
19mod tool_account;
20mod tool_args;
21mod tool_auth;
22mod tool_enums;
23mod tools;
24mod trade_pwd;
25mod trd_sdk_adapter;
26mod transport;
30
31use std::path::PathBuf;
32use std::sync::Arc;
33
34use anyhow::{Context, Result};
35use clap::{ArgMatches, CommandFactory, FromArgMatches, Parser, parser::ValueSource};
36use futu_auth::KeyStore;
37use rmcp::ServiceExt;
38use crate::transport::resilient_stdio;
40use tracing_subscriber::{
41 EnvFilter, Layer, filter::filter_fn, fmt, layer::SubscriberExt, util::SubscriberInitExt,
42};
43
44#[cfg(test)]
45pub(crate) use crate::card_num_expand::build_card_num_resolver;
46use crate::card_num_expand::spawn_card_num_expand_retry;
47#[cfg(unix)]
48use crate::card_num_expand::spawn_sighup_reload;
49use crate::http::serve_http;
50#[cfg(test)]
51use crate::http::{
52 authenticate_mcp_transport, inject_www_authenticate, oauth_protected_resource_metadata,
53 render_mcp_metrics_body_for,
54};
55use crate::state::ServerState;
56use crate::tools::FutuServer;
57
58#[derive(Default)]
59struct RmcpDrainEventFields {
60 message: Option<String>,
61 error: Option<String>,
62}
63
64impl tracing::field::Visit for RmcpDrainEventFields {
65 fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
66 match field.name() {
67 "message" => self.message = Some(value.to_string()),
68 "error" => self.error = Some(value.to_string()),
69 _ => {}
70 }
71 }
72
73 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
74 let rendered = format!("{value:?}");
75 let rendered = rendered
76 .strip_prefix('"')
77 .and_then(|value| value.strip_suffix('"'))
78 .unwrap_or(&rendered)
79 .to_string();
80 match field.name() {
81 "message" => self.message = Some(rendered),
82 "error" => self.error = Some(rendered),
83 _ => {}
84 }
85 }
86}
87
88fn is_expected_rmcp_drain_close(
89 target: &str,
90 level: &tracing::Level,
91 message: Option<&str>,
92 error: Option<&str>,
93) -> bool {
94 target == "rmcp::service"
95 && *level == tracing::Level::ERROR
96 && message == Some("failed to send pending response during drain")
97 && error == Some("channel closed")
98}
99
100#[derive(Clone, Copy)]
101struct SuppressExpectedRmcpDrainClose;
102
103impl<S> tracing_subscriber::layer::Filter<S> for SuppressExpectedRmcpDrainClose
104where
105 S: tracing::Subscriber,
106{
107 fn enabled(
108 &self,
109 _meta: &tracing::Metadata<'_>,
110 _cx: &tracing_subscriber::layer::Context<'_, S>,
111 ) -> bool {
112 true
113 }
114
115 fn event_enabled(
116 &self,
117 event: &tracing::Event<'_>,
118 _cx: &tracing_subscriber::layer::Context<'_, S>,
119 ) -> bool {
120 let meta = event.metadata();
121 if meta.target() != "rmcp::service" || *meta.level() != tracing::Level::ERROR {
122 return true;
123 }
124 let mut fields = RmcpDrainEventFields::default();
125 event.record(&mut fields);
126 !is_expected_rmcp_drain_close(
127 meta.target(),
128 meta.level(),
129 fields.message.as_deref(),
130 fields.error.as_deref(),
131 )
132 }
133}
134
135fn setup_logging(
140 default_level: &str,
141 audit_path: Option<&std::path::Path>,
142) -> Result<Option<tracing_appender::non_blocking::WorkerGuard>> {
143 let filter =
144 EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(default_level));
145
146 let fmt_layer = fmt::layer()
147 .with_timer(futu_core::log::LocalRfc3339Timer)
148 .with_writer(std::io::stderr)
149 .with_ansi(false)
150 .with_filter(SuppressExpectedRmcpDrainClose);
151
152 let registry = tracing_subscriber::registry().with(filter).with(fmt_layer);
153
154 if let Some(path) = audit_path {
155 let (writer, guard) = futu_auth::audit::open_writer(path)
156 .with_context(|| format!("open audit log {}", path.display()))?;
157 let audit_layer = fmt::layer()
158 .json()
159 .with_timer(futu_core::log::LocalRfc3339Timer)
160 .flatten_event(true)
161 .with_current_span(false)
162 .with_span_list(false)
163 .with_target(true)
164 .with_writer(writer)
165 .with_filter(filter_fn(|meta| meta.target() == futu_auth::audit::TARGET));
166 registry.with(audit_layer).init();
167 tracing::info!(
168 path = %path.display(),
169 "audit JSONL logger enabled (target=futu_audit)"
170 );
171 Ok(Some(guard))
172 } else {
173 registry.init();
174 Ok(None)
175 }
176}
177
178#[derive(Parser)]
180#[command(
181 name = "futu-mcp",
182 version,
183 about = "FutuOpenD-rs MCP server",
184 long_about = "通过 Model Context Protocol 暴露 Futu 行情/账户工具。默认 stdio transport。"
185)]
186struct Cli {
187 #[arg(short, long, env = "FUTU_GATEWAY", default_value = "127.0.0.1:11111")]
189 gateway: String,
190
191 #[arg(short, long)]
193 verbose: bool,
194
195 #[arg(long)]
201 keys_file: Option<PathBuf>,
202
203 #[arg(long, env = "FUTU_MCP_API_KEY", hide_env_values = true)]
208 api_key: Option<String>,
209
210 #[arg(long, env = "FUTU_TRADE_PWD_ACCOUNT")]
216 trade_pwd_account: Option<String>,
217
218 #[arg(long)]
224 enable_trading: bool,
225
226 #[arg(long, requires = "enable_trading")]
228 allow_real_trading: bool,
229
230 #[arg(long)]
237 audit_log: Option<PathBuf>,
238
239 #[arg(long)]
248 http_listen: Option<String>,
249
250 #[arg(long, requires = "tls_key")]
255 tls_cert: Option<PathBuf>,
256
257 #[arg(long, requires = "tls_cert")]
259 tls_key: Option<PathBuf>,
260
261 #[arg(long)]
273 config: Option<PathBuf>,
274}
275
276#[derive(Debug, Default, serde::Deserialize)]
291#[serde(default, deny_unknown_fields)]
292struct FileConfig {
293 gateway: Option<String>,
294 verbose: Option<bool>,
295 keys_file: Option<PathBuf>,
296 api_key: Option<String>,
297 trade_pwd_account: Option<String>,
298 enable_trading: Option<bool>,
299 allow_real_trading: Option<bool>,
300 audit_log: Option<PathBuf>,
301 http_listen: Option<String>,
302 tls_cert: Option<PathBuf>,
303 tls_key: Option<PathBuf>,
304}
305
306fn is_cli_explicit(matches: &ArgMatches, arg_id: &str) -> bool {
322 matches!(
323 matches.value_source(arg_id),
324 Some(ValueSource::CommandLine) | Some(ValueSource::EnvVariable)
325 )
326}
327
328impl Cli {
329 fn merge_config(mut self, matches: &ArgMatches) -> Result<Self> {
335 let Some(config_path) = &self.config else {
336 return Ok(self);
337 };
338 let content = std::fs::read_to_string(config_path)
339 .with_context(|| format!("read config file {}", config_path.display()))?;
340 let fc: FileConfig = toml::from_str(&content)
341 .with_context(|| format!("parse config file {}", config_path.display()))?;
342
343 if let Some(g) = fc.gateway
346 && !is_cli_explicit(matches, "gateway")
347 {
348 self.gateway = g;
349 }
350 if self.keys_file.is_none() {
353 self.keys_file = fc.keys_file;
354 }
355 if self.api_key.is_none()
356 && let Some(k) = fc.api_key
357 {
358 self.api_key = Some(k);
359 }
360 if self.trade_pwd_account.is_none() {
361 self.trade_pwd_account = fc.trade_pwd_account;
362 }
363 if fc.verbose.is_some() && !is_cli_explicit(matches, "verbose") {
368 self.verbose = fc.verbose.unwrap_or(false);
369 }
370 if fc.enable_trading.is_some() && !is_cli_explicit(matches, "enable_trading") {
371 self.enable_trading = fc.enable_trading.unwrap_or(false);
372 }
373 if fc.allow_real_trading.is_some() && !is_cli_explicit(matches, "allow_real_trading") {
374 self.allow_real_trading = fc.allow_real_trading.unwrap_or(false);
375 }
376 if self.audit_log.is_none() {
377 self.audit_log = fc.audit_log;
378 }
379 if self.http_listen.is_none() {
380 self.http_listen = fc.http_listen;
381 }
382 if self.tls_cert.is_none() {
383 self.tls_cert = fc.tls_cert;
384 }
385 if self.tls_key.is_none() {
386 self.tls_key = fc.tls_key;
387 }
388 eprintln!("[config] loaded {}", config_path.display());
390 Ok(self)
391 }
392}
393
394#[tokio::main]
395async fn main() -> Result<()> {
396 let matches = Cli::command().get_matches();
400 let cli = Cli::from_arg_matches(&matches)
401 .map_err(|e| anyhow::anyhow!("clap derive build failed: {e}"))?
402 .merge_config(&matches)?;
403
404 let default_level = if cli.verbose { "debug" } else { "info" };
406
407 let _audit_guard = setup_logging(default_level, cli.audit_log.as_deref())?;
409
410 let key_store = match &cli.keys_file {
412 Some(path) => {
413 let store = KeyStore::load(path)
414 .with_context(|| format!("load keys file {}", path.display()))?;
415 tracing::info!(
416 path = %path.display(),
417 keys_loaded = store.len(),
418 "scope mode: keys file loaded"
419 );
420 if cli.enable_trading || cli.allow_real_trading {
421 tracing::warn!(
422 "--enable-trading / --allow-real-trading are IGNORED in scope mode; \
423 trading permissions are controlled by API key scopes"
424 );
425 }
426 Arc::new(store)
427 }
428 None => {
429 tracing::info!("legacy mode: no keys file; using --enable-trading switches");
430 Arc::new(KeyStore::empty())
431 }
432 };
433
434 let http_mode = cli.http_listen.is_some();
436 let authed_key = if key_store.is_configured() && http_mode {
437 if cli.api_key.as_deref().is_some_and(|key| !key.is_empty()) {
438 tracing::warn!(
439 "--api-key / FUTU_MCP_API_KEY is stdio-only and IGNORED in HTTP scope mode; \
440 every /mcp request must provide its own Authorization Bearer"
441 );
442 }
443 None
444 } else if key_store.is_configured() {
445 match cli.api_key.as_deref() {
446 Some(plaintext) if !plaintext.is_empty() => match key_store.verify(plaintext) {
447 Some(rec) => {
448 tracing::info!(
449 key_id = %rec.id,
450 scopes = ?rec.scopes.iter().map(|s| s.as_str()).collect::<Vec<_>>(),
451 "API key verified"
452 );
453 Some(rec)
454 }
455 None => {
456 tracing::error!("FUTU_MCP_API_KEY does not match any key in keys.json");
457 None
458 }
459 },
460 _ => {
461 tracing::warn!(
462 "stdio scope mode active but FUTU_MCP_API_KEY not set; \
463 all tool calls will be rejected"
464 );
465 None
466 }
467 }
468 } else {
469 None
470 };
471
472 tracing::info!(
473 gateway = %cli.gateway,
474 scope_mode = key_store.is_configured(),
475 enable_trading = cli.enable_trading,
476 allow_real_trading = cli.allow_real_trading,
477 trade_pwd_account = cli.trade_pwd_account.as_deref().unwrap_or("<legacy/env>"),
478 "futu-mcp starting"
479 );
480 if !key_store.is_configured() && cli.enable_trading {
481 tracing::warn!(
482 allow_real_trading = cli.allow_real_trading,
483 "trading write tools ENABLED (legacy mode)"
484 );
485 }
486
487 let state = ServerState::new(cli.gateway)
488 .with_trading(cli.enable_trading, cli.allow_real_trading)
489 .with_key_store(key_store.clone())
490 .with_authed_key(authed_key)
491 .with_trade_pwd_account(cli.trade_pwd_account);
492 let server = FutuServer::new(state.clone());
493
494 if key_store.is_configured() && key_store.has_any_card_num_restrictions() {
511 spawn_card_num_expand_retry(state.clone(), key_store.clone());
512 } else if key_store.is_configured() {
513 tracing::debug!(
514 "v1.4.105 external report #4: keystore 无 allowed_card_nums 限制, 跳过 daemon expand"
515 );
516 }
517
518 #[cfg(unix)]
521 spawn_sighup_reload(key_store.clone(), state.clone());
522
523 futu_auth::metrics::install(std::sync::Arc::new(futu_auth::MetricsRegistry::default()));
527
528 if let Some(listen) = cli.http_listen {
529 let tls = match (cli.tls_cert, cli.tls_key) {
530 (Some(cert), Some(key)) => Some((cert, key)),
531 _ => None,
532 };
533 serve_http(server, key_store, &listen, tls).await?;
534 } else {
535 serve_stdio(server).await?;
536 }
537
538 Ok(())
539}
540
541async fn serve_stdio(server: tools::FutuServer) -> Result<()> {
548 let transport = resilient_stdio();
549 let drain = transport.pending_response_drain();
550 let service = match server.serve(transport).await {
551 Ok(service) => service,
552 Err(error) => {
553 if !drain.wait_pending_responses().await {
558 tracing::warn!(
559 "MCP stdio: initialize response did not drain before shutdown grace expired"
560 );
561 }
562 return Err(anyhow::anyhow!("MCP service init failed: {error}"));
563 }
564 };
565
566 if let Err(error) = service.waiting().await {
567 if !drain.wait_pending_responses().await {
568 tracing::warn!(
569 "MCP stdio: response did not drain before service shutdown grace expired"
570 );
571 }
572 return Err(anyhow::anyhow!("MCP service error: {error}"));
573 }
574 Ok(())
575}
576
577#[cfg(test)]
578mod tests;