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_MCP_OPEND_REST_URL")]
213 opend_rest_url: Option<String>,
214
215 #[arg(long, env = "FUTU_TRADE_PWD_ACCOUNT")]
221 trade_pwd_account: Option<String>,
222
223 #[arg(long)]
229 enable_trading: bool,
230
231 #[arg(long, requires = "enable_trading")]
233 allow_real_trading: bool,
234
235 #[arg(long)]
242 audit_log: Option<PathBuf>,
243
244 #[arg(long)]
253 http_listen: Option<String>,
254
255 #[arg(long, requires = "tls_key")]
260 tls_cert: Option<PathBuf>,
261
262 #[arg(long, requires = "tls_cert")]
264 tls_key: Option<PathBuf>,
265
266 #[arg(long)]
278 config: Option<PathBuf>,
279}
280
281#[derive(Debug, Default, serde::Deserialize)]
296#[serde(default, deny_unknown_fields)]
297struct FileConfig {
298 gateway: Option<String>,
299 verbose: Option<bool>,
300 keys_file: Option<PathBuf>,
301 api_key: Option<String>,
302 opend_rest_url: Option<String>,
303 trade_pwd_account: Option<String>,
304 enable_trading: Option<bool>,
305 allow_real_trading: Option<bool>,
306 audit_log: Option<PathBuf>,
307 http_listen: Option<String>,
308 tls_cert: Option<PathBuf>,
309 tls_key: Option<PathBuf>,
310}
311
312fn is_cli_explicit(matches: &ArgMatches, arg_id: &str) -> bool {
328 matches!(
329 matches.value_source(arg_id),
330 Some(ValueSource::CommandLine) | Some(ValueSource::EnvVariable)
331 )
332}
333
334impl Cli {
335 fn merge_config(mut self, matches: &ArgMatches) -> Result<Self> {
341 let Some(config_path) = &self.config else {
342 return Ok(self);
343 };
344 let content = std::fs::read_to_string(config_path)
345 .with_context(|| format!("read config file {}", config_path.display()))?;
346 let fc: FileConfig = toml::from_str(&content)
347 .with_context(|| format!("parse config file {}", config_path.display()))?;
348
349 if let Some(g) = fc.gateway
352 && !is_cli_explicit(matches, "gateway")
353 {
354 self.gateway = g;
355 }
356 if self.keys_file.is_none() {
359 self.keys_file = fc.keys_file;
360 }
361 if self.api_key.is_none()
362 && let Some(k) = fc.api_key
363 {
364 self.api_key = Some(k);
365 }
366 if self.opend_rest_url.is_none() {
367 self.opend_rest_url = fc.opend_rest_url;
368 }
369 if self.trade_pwd_account.is_none() {
370 self.trade_pwd_account = fc.trade_pwd_account;
371 }
372 if fc.verbose.is_some() && !is_cli_explicit(matches, "verbose") {
377 self.verbose = fc.verbose.unwrap_or(false);
378 }
379 if fc.enable_trading.is_some() && !is_cli_explicit(matches, "enable_trading") {
380 self.enable_trading = fc.enable_trading.unwrap_or(false);
381 }
382 if fc.allow_real_trading.is_some() && !is_cli_explicit(matches, "allow_real_trading") {
383 self.allow_real_trading = fc.allow_real_trading.unwrap_or(false);
384 }
385 if self.audit_log.is_none() {
386 self.audit_log = fc.audit_log;
387 }
388 if self.http_listen.is_none() {
389 self.http_listen = fc.http_listen;
390 }
391 if self.tls_cert.is_none() {
392 self.tls_cert = fc.tls_cert;
393 }
394 if self.tls_key.is_none() {
395 self.tls_key = fc.tls_key;
396 }
397 eprintln!("[config] loaded {}", config_path.display());
399 Ok(self)
400 }
401}
402
403#[tokio::main]
404async fn main() -> Result<()> {
405 let matches = Cli::command().get_matches();
409 let cli = Cli::from_arg_matches(&matches)
410 .map_err(|e| anyhow::anyhow!("clap derive build failed: {e}"))?
411 .merge_config(&matches)?;
412
413 let default_level = if cli.verbose { "debug" } else { "info" };
415
416 let _audit_guard = setup_logging(default_level, cli.audit_log.as_deref())?;
418
419 let key_store = match &cli.keys_file {
421 Some(path) => {
422 let store = KeyStore::load(path)
423 .with_context(|| format!("load keys file {}", path.display()))?;
424 tracing::info!(
425 path = %path.display(),
426 keys_loaded = store.len(),
427 "scope mode: keys file loaded"
428 );
429 if cli.enable_trading || cli.allow_real_trading {
430 tracing::warn!(
431 "--enable-trading / --allow-real-trading are IGNORED in scope mode; \
432 trading permissions are controlled by API key scopes"
433 );
434 }
435 Arc::new(store)
436 }
437 None => {
438 tracing::info!("legacy mode: no keys file; using --enable-trading switches");
439 Arc::new(KeyStore::empty())
440 }
441 };
442
443 let http_mode = cli.http_listen.is_some();
445 let authed_key = if key_store.is_configured() && http_mode {
446 if cli.api_key.as_deref().is_some_and(|key| !key.is_empty()) {
447 tracing::warn!(
448 "--api-key / FUTU_MCP_API_KEY is stdio-only and IGNORED in HTTP scope mode; \
449 every /mcp request must provide its own Authorization Bearer"
450 );
451 }
452 None
453 } else if key_store.is_configured() {
454 match cli.api_key.as_deref() {
455 Some(plaintext) if !plaintext.is_empty() => match key_store.verify(plaintext) {
456 Some(rec) => {
457 tracing::info!(
458 key_id = %rec.id,
459 scopes = ?rec.scopes.iter().map(|s| s.as_str()).collect::<Vec<_>>(),
460 "API key verified"
461 );
462 Some(rec)
463 }
464 None => {
465 tracing::error!("FUTU_MCP_API_KEY does not match any key in keys.json");
466 None
467 }
468 },
469 _ => {
470 tracing::warn!(
471 "stdio scope mode active but FUTU_MCP_API_KEY not set; \
472 all tool calls will be rejected"
473 );
474 None
475 }
476 }
477 } else {
478 None
479 };
480
481 tracing::info!(
482 gateway = %cli.gateway,
483 scope_mode = key_store.is_configured(),
484 enable_trading = cli.enable_trading,
485 allow_real_trading = cli.allow_real_trading,
486 trade_pwd_account = cli.trade_pwd_account.as_deref().unwrap_or("<legacy/env>"),
487 "futu-mcp starting"
488 );
489 if !key_store.is_configured() && cli.enable_trading {
490 tracing::warn!(
491 allow_real_trading = cli.allow_real_trading,
492 "trading write tools ENABLED (legacy mode)"
493 );
494 }
495
496 let state = ServerState::new(cli.gateway)
497 .with_trading(cli.enable_trading, cli.allow_real_trading)
498 .with_key_store(key_store.clone())
499 .with_authed_key(authed_key)
500 .with_opend_rest(cli.opend_rest_url, cli.api_key.clone())
501 .with_trade_pwd_account(cli.trade_pwd_account);
502 let server = FutuServer::new(state.clone());
503
504 if key_store.is_configured() && key_store.has_any_card_num_restrictions() {
521 spawn_card_num_expand_retry(state.clone(), key_store.clone());
522 } else if key_store.is_configured() {
523 tracing::debug!(
524 "v1.4.105 external report #4: keystore 无 allowed_card_nums 限制, 跳过 daemon expand"
525 );
526 }
527
528 #[cfg(unix)]
531 spawn_sighup_reload(key_store.clone(), state.clone());
532
533 futu_auth::metrics::install(std::sync::Arc::new(futu_auth::MetricsRegistry::default()));
537
538 if let Some(listen) = cli.http_listen {
539 let tls = match (cli.tls_cert, cli.tls_key) {
540 (Some(cert), Some(key)) => Some((cert, key)),
541 _ => None,
542 };
543 serve_http(server, key_store, &listen, tls).await?;
544 } else {
545 serve_stdio(server).await?;
546 }
547
548 Ok(())
549}
550
551async fn serve_stdio(server: tools::FutuServer) -> Result<()> {
558 let transport = resilient_stdio();
559 let drain = transport.pending_response_drain();
560 let service = match server.serve(transport).await {
561 Ok(service) => service,
562 Err(error) => {
563 if !drain.wait_pending_responses().await {
568 tracing::warn!(
569 "MCP stdio: initialize response did not drain before shutdown grace expired"
570 );
571 }
572 return Err(anyhow::anyhow!("MCP service init failed: {error}"));
573 }
574 };
575
576 if let Err(error) = service.waiting().await {
577 if !drain.wait_pending_responses().await {
578 tracing::warn!(
579 "MCP stdio: response did not drain before service shutdown grace expired"
580 );
581 }
582 return Err(anyhow::anyhow!("MCP service error: {error}"));
583 }
584 Ok(())
585}
586
587#[cfg(test)]
588mod tests;