Skip to main content

futu_mcp/
main.rs

1//! FutuOpenD-rs MCP 服务器
2//!
3//! 通过 Model Context Protocol 把 Futu 行情/账户能力暴露给 Claude / LLM 客户端。
4//!
5//! 授权有两种模式:
6//!
7//! - **Scope 模式**:`--keys-file <path>` 启用。stdio 客户端通过
8//!   `FUTU_MCP_API_KEY` 提供进程级 key;HTTP 客户端的每个 `/mcp` 请求必须
9//!   携带有效 `Authorization: Bearer`。服务器按 scope + 限额放行。
10//! - **Legacy 模式**:未提供 keys-file 时回退到旧的
11//!   `--enable-trading` / `--allow-real-trading` 两级开关。
12
13mod 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;
26// v1.4.90 P0-A: resilient stdio transport — recovers from malformed JSON
27// (e.g. `{"price": Infinity}`) instead of `exit(0)`-ing the whole server.
28// See crates/futu-mcp/src/transport.rs for full rationale.
29mod 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;
38// v1.4.90 P0-A: stdio() (rmcp default) treats parse errors as fatal — see transport.rs.
39use 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
135/// 初始化 stderr 日志 + 可选 audit JSONL 层
136///
137/// - 常规事件走 stderr(no-ansi,因为 MCP client 的 stderr 往往不是 tty)
138/// - 如果 `audit_path` 传了,加一个 target=futu_audit 的 JSON 层写到文件/目录
139fn 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/// FutuOpenD-rs MCP server
179#[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    /// 网关地址(可用 FUTU_GATEWAY 环境变量覆盖)
188    #[arg(short, long, env = "FUTU_GATEWAY", default_value = "127.0.0.1:11111")]
189    gateway: String,
190
191    /// 启用 debug 日志
192    #[arg(short, long)]
193    verbose: bool,
194
195    /// Scope 模式:加载 keys.json 文件(API Key 授权)。
196    ///
197    /// stdio 模式需通过 FUTU_MCP_API_KEY / --api-key 提供默认身份;HTTP 模式
198    /// 则要求每个 /mcp 请求携带有效 Authorization Bearer。scope / 限额由
199    /// keys.json 配置决定;此时 --enable-trading / --allow-real-trading 被忽略。
200    #[arg(long)]
201    keys_file: Option<PathBuf>,
202
203    /// stdio 进程级 API Key 明文(等价于 FUTU_MCP_API_KEY 环境变量)。
204    /// HTTP scope 模式不会用它替代请求自己的 Authorization Bearer。
205    ///
206    /// 生产环境强烈建议用环境变量而非命令行参数(后者会进 `ps` 输出)。
207    #[arg(long, env = "FUTU_MCP_API_KEY", hide_env_values = true)]
208    api_key: Option<String>,
209
210    /// 交易密码所属登录账号,用于兼容读取账号级 credential-store 条目。
211    ///
212    /// 读取优先级为 FUTU_TRADE_PWD(trim 后非空)> 账号级条目 > 旧全局条目;
213    /// 未设置账号时会尝试 FUTU_ACCOUNT 作为账号 hint。macOS 发布包不支持
214    /// `futucli` 跨 binary 写入,MCP 请使用 FUTU_TRADE_PWD。
215    #[arg(long, env = "FUTU_TRADE_PWD_ACCOUNT")]
216    trade_pwd_account: Option<String>,
217
218    /// [Legacy] 启用交易写工具(place / modify / cancel)。默认关闭。
219    ///
220    /// 开启后默认仅允许 simulate 环境;要操作真实账户需额外 --allow-real-trading。
221    /// 注意:下单前网关必须已 unlock_trade(密码不经过 MCP / LLM)。
222    /// 若提供了 --keys-file,此开关被忽略,改由 key 的 scope 决定。
223    #[arg(long)]
224    enable_trading: bool,
225
226    /// [Legacy] 允许交易写工具对 real 环境执行。必须与 --enable-trading 搭配。
227    #[arg(long, requires = "enable_trading")]
228    allow_real_trading: bool,
229
230    /// 审计日志输出:JSONL 文件路径或目录
231    ///
232    /// - 带扩展名(`/var/log/futu-mcp-audit.jsonl`)→ 单文件 append
233    /// - 不带扩展名 / 以 `/` 结尾 → 每日滚动 `futu-audit.log` + 日期
234    ///
235    /// 只记录 auth / 交易 事件(target = `futu_audit`)。
236    #[arg(long)]
237    audit_log: Option<PathBuf>,
238
239    /// 以 HTTP transport 启动(streamable HTTP),监听该端口(格式 `host:port` 或 `:port`)
240    ///
241    /// 默认 stdio:LLM 客户端启子进程走 stdin/stdout。开 HTTP 后可以让多个
242    /// 客户端连同一个 MCP 进程,并同时暴露 `/metrics`。HTTP 调用统一通过
243    /// `Authorization: Bearer <api-key>` 鉴权;只有 schema 显式声明 `api_key`
244    /// 的工具才支持参数级 key 覆盖。
245    ///
246    /// 例:`--http-listen 127.0.0.1:3000` / `--http-listen :3000`
247    #[arg(long)]
248    http_listen: Option<String>,
249
250    /// TLS 证书文件路径(PEM 格式;需与 --tls-key 配合)
251    ///
252    /// 启用后 HTTP transport 走 HTTPS。若不设置,走纯 HTTP(建议前置 Caddy / Nginx
253    /// 做 TLS 终止)。
254    #[arg(long, requires = "tls_key")]
255    tls_cert: Option<PathBuf>,
256
257    /// TLS 私钥文件路径(PEM 格式;需与 --tls-cert 配合)
258    #[arg(long, requires = "tls_cert")]
259    tls_key: Option<PathBuf>,
260
261    /// TOML 配置文件路径(字段名与 CLI 参数一致,CLI 参数覆盖配置文件)
262    ///
263    /// 示例:
264    /// ```toml
265    /// gateway = "10.0.0.1:11111"
266    /// http_listen = ":3000"
267    /// keys_file = "/etc/futu/keys.json"
268    /// audit_log = "/var/log/futu-mcp-audit.jsonl"
269    /// tls_cert = "/etc/futu/cert.pem"
270    /// tls_key  = "/etc/futu/key.pem"
271    /// ```
272    #[arg(long)]
273    config: Option<PathBuf>,
274}
275
276/// TOML 配置文件映射——字段名与 CLI 参数完全一致
277///
278/// codex 0547 F4 (P2) fix: 加 `#[serde(deny_unknown_fields)]` — 与 `futu-opend`
279/// XmlConfig (BUG-006 v1.4.102 加的) 同语义级别. typo (e.g. `keys_flie` /
280/// `auditlog`) 之前 silent drop, 用户配置 silent 失效:
281/// - `key_file` typo → keystore 不加载 → MCP 进 legacy mode (无 scope)
282/// - `http_litsen` typo → HTTP transport 不启动 (默认 stdio)
283/// - `auditlog` typo → 审计文件无事件
284///
285/// 修后: 任何 unknown field / typo / `[server]` 类未支持 section 立即 parse
286/// fatal, daemon abort + 清晰错误.
287///
288/// **不 break 老 deprecated alias**: 没有 alias 历史, MCP TOML schema 自 v1.0
289/// 起字段名稳定, 升级用户无 typo 不会被影响.
290#[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
306/// codex 0547 F5 (P3) fix: clap `ValueSource` 区分 "CLI 显式传" vs "默认值".
307///
308/// 之前用 `self.gateway == "127.0.0.1:11111"` 判 "CLI 没传" — 当用户**显式**
309/// 传 `--gateway 127.0.0.1:11111` (与默认值相等) 时被错判为 "未传" → TOML
310/// 配置 gateway override 反向. 违背 "CLI 始终覆盖配置文件" 契约.
311///
312/// 同模式: bool 字段 (`verbose` / `enable_trading` / `allow_real_trading`)
313/// 之前用 `if !self.field` 判, 用户显式 `--verbose` 时反复 = false 也不能
314/// 区分 "CLI 没传 + TOML 也没设" 与 "CLI 显式 false (无 --no-flag)". clap
315/// derive 不天然支持 `--no-*`, 所以 bool 字段的 explicit-false override 是
316/// 设计 limitation; 但 explicit-true 一定要尊重.
317///
318/// 本 helper 接受 `&ArgMatches` 与 `arg_id`, 返 true 仅在 user explicitly 传
319/// (而不是 default / env). 见
320/// <https://docs.rs/clap/latest/clap/parser/enum.ValueSource.html>.
321fn 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    /// 如果指定了 `--config`,先从文件读取默认值,再让 CLI 参数覆盖。
330    ///
331    /// codex 0547 F5 (P3): 用 `&ArgMatches` 精准判断 "CLI/env 是否显式传"
332    /// 而非旧的 "值 == 默认 → 当作没传" 启发式. CLI 显式传 = 显式 (即使值
333    /// 等于默认). TOML 文件值仅在 CLI / env 都没显式传时才采用.
334    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        // codex 0547 F5: gateway 用 ValueSource 精准判. 之前 "self.gateway ==
344        // 默认值" 启发式在用户**显式**传默认值时反向 (TOML 覆盖 CLI).
345        if let Some(g) = fc.gateway
346            && !is_cli_explicit(matches, "gateway")
347        {
348            self.gateway = g;
349        }
350        // codex 0547 F5: Option 字段用 None check (CLI 没传 = None, 不会与
351        // 默认值 ambiguity).
352        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        // codex 0547 F5: bool 字段用 ValueSource 区分 "未传" vs "显式 false".
364        // clap derive 没 `--no-verbose`, 所以无法 explicit-set false; 但用户
365        // **显式 true** (e.g. `--verbose` 在 CLI) 要保留 (TOML 即使写 false
366        // 也不能覆盖 explicit-true).
367        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        // 此时 tracing 可能还没初始化,写 stderr 即可
389        eprintln!("[config] loaded {}", config_path.display());
390        Ok(self)
391    }
392}
393
394#[tokio::main]
395async fn main() -> Result<()> {
396    // codex 0547 F5 (P3): 解析两次 — 用 ArgMatches 区分 explicit vs default
397    // 后再用 derive 反向 build Cli 结构. 单次 parse 走不通 (Cli::parse 不暴露
398    // ArgMatches), 但解析+from_arg_matches 是 0-allocation cycle.
399    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    // MCP 用 stdout 传协议帧,所有日志必须写 stderr
405    let default_level = if cli.verbose { "debug" } else { "info" };
406
407    // audit 日志 guard 必须活到 main 返回;否则后台 flush 可能丢事件
408    let _audit_guard = setup_logging(default_level, cli.audit_log.as_deref())?;
409
410    // ---------- 加载 KeyStore ----------
411    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    // ---------- 校验调用方 API key ----------
435    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    // v1.4.105 external reviewer #4 (BUG-v1.4.104-002, P1) fix: standalone MCP 启动时把
495    // `allowed_card_nums` (string format, e.g. ["0757"]) resolve 成
496    // `allowed_acc_ids` (numeric set), 行为与 futu-opend daemon 启动 +
497    // SIGHUP 路径**byte-identical** (resolver 4-suffix / 16-exact 双匹配
498    // card_num + uni_card_num, 1 个 → resolved, 0 → unresolved warn,
499    // ≥2 → ambiguous warn). 不做 expand 时 KeyStore::load_file 注入的
500    // fail-closed sentinel `allowed_acc_ids = {0}` 会让真账户 acc_id ≠ 0
501    // 永远 reject "not in allowed list {0}".
502    //
503    // 设计要点:
504    // - 仅在 KeyStore 至少一条 key 配置了 allowed_card_nums 时才连 daemon
505    //   (避免无意义 GetAccList 请求)
506    // - daemon 可能尚未起来 / connect race → 后台 task 重试 6 次 × 10s
507    //   覆盖 daemon 启动 ~60s 窗口, 与 daemon 内部 trd_cache 加载 retry 同节奏
508    // - expand 失败不阻塞 MCP server 启动 — sentinel 仍生效保护
509    // - SIGHUP 重载 keys.json 后必须重新 expand (sentinel/旧 acc_ids 会失效)
510    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    // SIGHUP 热重载 keys.json(unix only)— v1.4.105 external reviewer #4: reload 后 +
519    // re-expand card_num (与 daemon `card_num_reload_and_expand_fn` 同语义)
520    #[cfg(unix)]
521    spawn_sighup_reload(key_store.clone(), state.clone());
522
523    // v1.0:install 全局 metrics registry,让 audit::* 的 counter hook 起作用
524    // (HTTP 模式下 /metrics 端点消费这套;stdio 模式虽然没 HTTP,但写进内存
525    // 方便 debug 和后续加 transport)
526    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
541/// stdio 模式:MCP 客户端启动子进程,stdin/stdout 传协议帧
542///
543/// v1.4.90 P0-A: uses `resilient_stdio()` instead of `rmcp::transport::stdio()`
544/// — a malformed JSON line (e.g. `{"price": Infinity}` from an LLM client)
545/// now produces a `-32700 Parse error` response instead of `exit(0)`-ing the
546/// entire server. See crates/futu-mcp/src/transport.rs for full background.
547async 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            // rmcp can return an initialize error before the background stdout
554            // writer flushes its JSON-RPC error frame. Keep the runtime alive
555            // until the transport confirms every accepted response ID was
556            // written, or the existing bounded drain grace expires.
557            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;