From 946363e9194eb540221a72e03e53950551f7a403 Mon Sep 17 00:00:00 2001 From: Philippe Branchu Date: Fri, 27 Mar 2026 13:57:25 +0000 Subject: [PATCH 1/5] Fix silent reinforcement: recognize [SILENT] token in agent replies Extract is_silent_token() helper for case-insensitive [SILENT] detection. Revert unrelated Cargo.lock and formatting changes. Add unit tests. Co-Authored-By: Claude Opus 4.6 (1M context) --- Cargo.lock | 423 +++---------- crates/openfang-runtime/src/agent_loop.rs | 49 +- crates/openfang-runtime/src/mcp.rs | 698 +++++++++++++++------ crates/openfang-runtime/src/tool_runner.rs | 11 +- 4 files changed, 620 insertions(+), 561 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index df329192..4621940c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -139,7 +139,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -150,7 +150,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -384,28 +384,6 @@ version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" -[[package]] -name = "aws-lc-rs" -version = "1.16.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a054912289d18629dc78375ba2c3726a3afe3ff71b4edba9dedfca0e3446d1fc" -dependencies = [ - "aws-lc-sys", - "zeroize", -] - -[[package]] -name = "aws-lc-sys" -version = "0.39.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1fa7e52a4c5c547c741610a2c6f123f3881e409b714cd27e6798ef020c514f0a" -dependencies = [ - "cc", - "cmake", - "dunce", - "fs_extra", -] - [[package]] name = "axum" version = "0.8.8" @@ -882,15 +860,6 @@ version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" -[[package]] -name = "cmake" -version = "0.1.57" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "75443c44cd6b379beb8c5b45d85d0773baf31cce901fe7bb252f4eff3008ef7d" -dependencies = [ - "cc", -] - [[package]] name = "cobs" version = "0.3.0" @@ -912,7 +881,7 @@ version = "3.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "faf9468729b8cbcea668e36183cb69d317348c2e08e994829fb56ebfdfbaac34" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -988,16 +957,6 @@ dependencies = [ "version_check", ] -[[package]] -name = "core-foundation" -version = "0.9.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "91e195e091a93c46f7102ec7818a2aa394e1e1771c3ab4825963fa03e45afb8f" -dependencies = [ - "core-foundation-sys", - "libc", -] - [[package]] name = "core-foundation" version = "0.10.1" @@ -1021,7 +980,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "064badf302c3194842cf2c5d61f56cc88e54a759313879cdf03abdd27d0c3b97" dependencies = [ "bitflags 2.11.0", - "core-foundation 0.10.1", + "core-foundation", "core-graphics-types", "foreign-types 0.5.0", "libc", @@ -1034,7 +993,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3d44a101f213f6c4cdc1853d4b78aef6db6bdfa3468798cc1d9912f4735013eb" dependencies = [ "bitflags 2.11.0", - "core-foundation 0.10.1", + "core-foundation", "libc", ] @@ -1579,7 +1538,7 @@ dependencies = [ "libc", "option-ext", "redox_users 0.5.2", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -1829,7 +1788,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -1930,17 +1889,6 @@ dependencies = [ "zlib-rs", ] -[[package]] -name = "flume" -version = "0.11.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" -dependencies = [ - "futures-core", - "futures-sink", - "spin", -] - [[package]] name = "fnv" version = "1.0.7" @@ -2010,12 +1958,6 @@ dependencies = [ "percent-encoding", ] -[[package]] -name = "fs_extra" -version = "1.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" - [[package]] name = "futf" version = "0.1.5" @@ -2463,9 +2405,9 @@ dependencies = [ [[package]] name = "governor" -version = "0.10.4" +version = "0.8.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9efcab3c1958580ff1f25a2a41be1668f7603d849bb63af523b208a3cc1223b8" +checksum = "be93b4ec2e4710b04d9264c0c7350cdd62a8c20e5e4ac732552ebb8f0debe8eb" dependencies = [ "cfg-if", "dashmap", @@ -2473,7 +2415,7 @@ dependencies = [ "futures-timer", "futures-util", "getrandom 0.3.4", - "hashbrown 0.16.1", + "no-std-compat", "nonzero_ext", "parking_lot", "portable-atomic", @@ -2569,11 +2511,6 @@ name = "hashbrown" version = "0.16.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "841d1cc9bed7f9236f321df977030373f4a4163ae1a7dbfe1a51a2c1a51d9100" -dependencies = [ - "allocator-api2", - "equivalent", - "foldhash 0.2.0", -] [[package]] name = "hashlink" @@ -2735,10 +2672,10 @@ dependencies = [ "http", "hyper", "hyper-util", - "rustls 0.23.37", + "rustls", "rustls-pki-types", "tokio", - "tokio-rustls 0.26.4", + "tokio-rustls", "tower-service", "webpki-roots", ] @@ -3269,10 +3206,10 @@ dependencies = [ "nom 8.0.0", "percent-encoding", "quoted_printable", - "rustls 0.23.37", + "rustls", "socket2 0.6.3", "tokio", - "tokio-rustls 0.26.4", + "tokio-rustls", "url", "webpki-roots", ] @@ -3627,10 +3564,10 @@ dependencies = [ "libc", "log", "openssl", - "openssl-probe 0.2.1", + "openssl-probe", "openssl-sys", "schannel", - "security-framework 3.7.0", + "security-framework", "security-framework-sys", "tempfile", ] @@ -3672,16 +3609,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086" [[package]] -name = "nix" -version = "0.31.2" +name = "no-std-compat" +version = "0.4.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5d6d0705320c1e6ba1d912b5e37cf18071b6c2e9b7fa8215a1e8a7651966f5d3" -dependencies = [ - "bitflags 2.11.0", - "cfg-if", - "cfg_aliases", - "libc", -] +checksum = "b93853da6d84c2e3c7d730d6473e8817692dd89be387eb01b94d7f108ecb5b8c" [[package]] name = "nodrop" @@ -3735,7 +3666,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -3986,7 +3917,7 @@ dependencies = [ "tokio", "tokio-stream", "tokio-test", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tower", "tower-http", "tracing", @@ -4017,7 +3948,6 @@ dependencies = [ "regex-lite", "reqwest 0.12.28", "roxmltree", - "rumqttc", "serde", "serde_json", "sha1", @@ -4053,7 +3983,7 @@ dependencies = [ "serde_json", "tempfile", "tokio", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "tracing-subscriber", "uuid", @@ -4081,7 +4011,7 @@ dependencies = [ "tauri-plugin-single-instance", "tauri-plugin-updater", "tokio", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "tracing-subscriber", ] @@ -4107,7 +4037,7 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tokio-test", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "url", "uuid", @@ -4126,7 +4056,7 @@ dependencies = [ "tempfile", "thiserror 2.0.18", "tokio-test", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "uuid", ] @@ -4162,7 +4092,7 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tokio-test", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "tracing-subscriber", "uuid", @@ -4202,7 +4132,7 @@ dependencies = [ "serde_yaml", "tempfile", "thiserror 2.0.18", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "uuid", "walkdir", @@ -4220,13 +4150,11 @@ dependencies = [ "dashmap", "futures", "hex", - "http", "openfang-memory", "openfang-skills", "openfang-types", "regex-lite", "reqwest 0.12.28", - "rmcp", "rusqlite", "serde", "serde_json", @@ -4260,7 +4188,7 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tokio-test", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "tracing", "uuid", "walkdir", @@ -4282,7 +4210,7 @@ dependencies = [ "serde_json", "sha2", "thiserror 2.0.18", - "toml 0.9.12+spec-1.1.0", + "toml 0.8.2", "uuid", ] @@ -4334,12 +4262,6 @@ dependencies = [ "syn 2.0.117", ] -[[package]] -name = "openssl-probe" -version = "0.1.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d05e27ee213611ffe7d6348b942e8f942b37114c00cc03cec254295a4a17852e" - [[package]] name = "openssl-probe" version = "0.2.1" @@ -4391,7 +4313,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7d8fae84b431384b68627d0f9b3b1245fcf9f46f6c0e3dc902e9dce64edd1967" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.45.0", ] [[package]] @@ -4985,20 +4907,6 @@ dependencies = [ "unicode-ident", ] -[[package]] -name = "process-wrap" -version = "9.1.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2e842efad9119158434d193c6682e2ebee4b44d6ad801d7b349623b3f57cdf55" -dependencies = [ - "futures", - "indexmap 2.13.0", - "nix", - "tokio", - "tracing", - "windows 0.62.2", -] - [[package]] name = "prost" version = "0.13.5" @@ -5106,7 +5014,7 @@ dependencies = [ "quinn-proto", "quinn-udp", "rustc-hash", - "rustls 0.23.37", + "rustls", "socket2 0.6.3", "thiserror 2.0.18", "tokio", @@ -5120,14 +5028,13 @@ version = "0.11.14" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "434b42fec591c96ef50e21e886936e66d3cc3f737104fdb9b737c40ffb94c098" dependencies = [ - "aws-lc-rs", "bytes", "getrandom 0.3.4", "lru-slab", "rand 0.9.2", "ring", "rustc-hash", - "rustls 0.23.37", + "rustls", "rustls-pki-types", "slab", "thiserror 2.0.18", @@ -5484,14 +5391,14 @@ dependencies = [ "percent-encoding", "pin-project-lite", "quinn", - "rustls 0.23.37", + "rustls", "rustls-pki-types", "serde", "serde_json", "serde_urlencoded", "sync_wrapper", "tokio", - "tokio-rustls 0.26.4", + "tokio-rustls", "tokio-util", "tower", "tower-http", @@ -5524,15 +5431,14 @@ dependencies = [ "log", "percent-encoding", "pin-project-lite", - "quinn", - "rustls 0.23.37", + "rustls", "rustls-pki-types", "rustls-platform-verifier", "serde", "serde_json", "sync_wrapper", "tokio", - "tokio-rustls 0.26.4", + "tokio-rustls", "tokio-util", "tower", "tower-http", @@ -5582,29 +5488,6 @@ dependencies = [ "windows-sys 0.52.0", ] -[[package]] -name = "rmcp" -version = "1.2.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ba6b9d2f0efe2258b23767f1f9e0054cfbcac9c2d6f81a031214143096d7864f" -dependencies = [ - "async-trait", - "chrono", - "futures", - "http", - "pin-project-lite", - "process-wrap", - "reqwest 0.13.2", - "serde", - "serde_json", - "sse-stream", - "thiserror 2.0.18", - "tokio", - "tokio-stream", - "tokio-util", - "tracing", -] - [[package]] name = "rmp" version = "0.8.15" @@ -5633,24 +5516,6 @@ dependencies = [ "memchr", ] -[[package]] -name = "rumqttc" -version = "0.24.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e1568e15fab2d546f940ed3a21f48bbbd1c494c90c99c4481339364a497f94a9" -dependencies = [ - "bytes", - "flume", - "futures-util", - "log", - "rustls-native-certs 0.7.3", - "rustls-pemfile", - "rustls-webpki 0.102.8", - "thiserror 1.0.69", - "tokio", - "tokio-rustls 0.25.0", -] - [[package]] name = "rusqlite" version = "0.31.0" @@ -5710,21 +5575,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.61.2", -] - -[[package]] -name = "rustls" -version = "0.22.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bf4ef73721ac7bcd79b2b315da7779d8fc09718c6b3d2d1b2d94850eb8c18432" -dependencies = [ - "log", - "ring", - "rustls-pki-types", - "rustls-webpki 0.102.8", - "subtle", - "zeroize", + "windows-sys 0.52.0", ] [[package]] @@ -5733,48 +5584,25 @@ version = "0.23.37" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "758025cb5fccfd3bc2fd74708fd4682be41d99e5dff73c377c0646c6012c73a4" dependencies = [ - "aws-lc-rs", "log", "once_cell", "ring", "rustls-pki-types", - "rustls-webpki 0.103.10", + "rustls-webpki", "subtle", "zeroize", ] -[[package]] -name = "rustls-native-certs" -version = "0.7.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e5bfb394eeed242e909609f56089eecfe5fda225042e8b171791b9c95f5931e5" -dependencies = [ - "openssl-probe 0.1.6", - "rustls-pemfile", - "rustls-pki-types", - "schannel", - "security-framework 2.11.1", -] - [[package]] name = "rustls-native-certs" version = "0.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "612460d5f7bea540c490b2b6395d8e34a953e52b491accd6c86c8164c5932a63" dependencies = [ - "openssl-probe 0.2.1", + "openssl-probe", "rustls-pki-types", "schannel", - "security-framework 3.7.0", -] - -[[package]] -name = "rustls-pemfile" -version = "2.2.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dce314e5fee3f39953d46bb63bb8a46d40c2f8fb7cc5a3b6cab2bde9721d6e50" -dependencies = [ - "rustls-pki-types", + "security-framework", ] [[package]] @@ -5793,19 +5621,19 @@ version = "0.6.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d99feebc72bae7ab76ba994bb5e121b8d83d910ca40b36e0921f53becc41784" dependencies = [ - "core-foundation 0.10.1", + "core-foundation", "core-foundation-sys", "jni", "log", "once_cell", - "rustls 0.23.37", - "rustls-native-certs 0.8.3", + "rustls", + "rustls-native-certs", "rustls-platform-verifier-android", - "rustls-webpki 0.103.10", - "security-framework 3.7.0", + "rustls-webpki", + "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -5816,27 +5644,15 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.102.8" +version = "0.103.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "64ca1bc8749bd4cf37b5ce386cc146580777b4e8572c7b97baf22c83f444bee9" +checksum = "d7df23109aa6c1567d1c575b9952556388da57401e4ace1d15f79eedad0d8f53" dependencies = [ "ring", "rustls-pki-types", "untrusted", ] -[[package]] -name = "rustls-webpki" -version = "0.103.10" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df33b2b81ac578cabaf06b89b0631153a3f416b0a886e8a7a1707fb51abbd1ef" -dependencies = [ - "aws-lc-rs", - "ring", - "rustls-pki-types", - "untrusted", -] - [[package]] name = "rustversion" version = "1.0.22" @@ -5924,19 +5740,6 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" -[[package]] -name = "security-framework" -version = "2.11.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "897b2245f0b511c87893af39b033e5ca9cce68824c4d7e7630b5a1d339658d02" -dependencies = [ - "bitflags 2.11.0", - "core-foundation 0.9.4", - "core-foundation-sys", - "libc", - "security-framework-sys", -] - [[package]] name = "security-framework" version = "3.7.0" @@ -5944,7 +5747,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d" dependencies = [ "bitflags 2.11.0", - "core-foundation 0.10.1", + "core-foundation", "core-foundation-sys", "libc", "security-framework-sys", @@ -6375,7 +6178,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3a766e1110788c36f4fa1c2b71b387a7815aa65f88ce0229841826633d93723e" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -6431,9 +6234,6 @@ name = "spin" version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" -dependencies = [ - "lock_api", -] [[package]] name = "spinning_top" @@ -6454,19 +6254,6 @@ dependencies = [ "der", ] -[[package]] -name = "sse-stream" -version = "0.2.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "eb4dc4d33c68ec1f27d386b5610a351922656e1fdf5c05bbaad930cd1519479a" -dependencies = [ - "bytes", - "futures-util", - "http-body", - "http-body-util", - "pin-project-lite", -] - [[package]] name = "stable_deref_trait" version = "1.2.1" @@ -6649,7 +6436,7 @@ checksum = "6e06d52c379e63da659a483a958110bbde891695a0ecb53e48cc7786d5eda7bb" dependencies = [ "bitflags 2.11.0", "block2", - "core-foundation 0.10.1", + "core-foundation", "core-graphics", "crossbeam-channel", "dispatch2", @@ -6673,7 +6460,7 @@ dependencies = [ "tao-macros", "unicode-segmentation", "url", - "windows 0.61.3", + "windows", "windows-core 0.61.2", "windows-version", "x11-dl", @@ -6692,9 +6479,9 @@ dependencies = [ [[package]] name = "tar" -version = "0.4.45" +version = "0.4.44" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "22692a6476a21fa75fdfc11d452fda482af402c008cdbaf3476414e122040973" +checksum = "1d863878d212c87a19c1a610eb53bb01fe12951c0501cf5a0d65f724914a667a" dependencies = [ "filetime", "libc", @@ -6762,7 +6549,7 @@ dependencies = [ "webkit2gtk", "webview2-com", "window-vibrancy", - "windows 0.61.3", + "windows", ] [[package]] @@ -6986,7 +6773,7 @@ dependencies = [ "osakit", "percent-encoding", "reqwest 0.13.2", - "rustls 0.23.37", + "rustls", "semver", "serde", "serde_json", @@ -7024,7 +6811,7 @@ dependencies = [ "url", "webkit2gtk", "webview2-com", - "windows 0.61.3", + "windows", ] [[package]] @@ -7049,7 +6836,7 @@ dependencies = [ "url", "webkit2gtk", "webview2-com", - "windows 0.61.3", + "windows", "wry", ] @@ -7110,7 +6897,7 @@ checksum = "0b1e66e07de489fe43a46678dd0b8df65e0c973909df1b60ba33874e297ba9b9" dependencies = [ "quick-xml 0.37.5", "thiserror 2.0.18", - "windows 0.61.3", + "windows", "windows-version", ] @@ -7124,7 +6911,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -7280,24 +7067,13 @@ dependencies = [ "syn 2.0.117", ] -[[package]] -name = "tokio-rustls" -version = "0.25.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "775e0c0f0adb3a2f22a00c4745d728b479985fc15ee7ca6a2608388c5569860f" -dependencies = [ - "rustls 0.22.4", - "rustls-pki-types", - "tokio", -] - [[package]] name = "tokio-rustls" version = "0.26.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" dependencies = [ - "rustls 0.23.37", + "rustls", "tokio", ] @@ -7331,11 +7107,11 @@ checksum = "edc5f74e248dc973e0dbb7b74c7e0d6fcc301c694ff50049504004ef4d0cdcd9" dependencies = [ "futures-util", "log", - "rustls 0.23.37", - "rustls-native-certs 0.8.3", + "rustls", + "rustls-native-certs", "rustls-pki-types", "tokio", - "tokio-rustls 0.26.4", + "tokio-rustls", "tungstenite 0.24.0", ] @@ -7637,7 +7413,7 @@ dependencies = [ "httparse", "log", "rand 0.8.5", - "rustls 0.23.37", + "rustls", "rustls-pki-types", "sha1", "thiserror 1.0.69", @@ -7687,7 +7463,7 @@ checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ "memoffset", "tempfile", - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -8534,7 +8310,7 @@ checksum = "7130243a7a5b33c54a444e54842e6a9e133de08b5ad7b5861cd8ed9a6a5bc96a" dependencies = [ "webview2-com-macros", "webview2-com-sys", - "windows 0.61.3", + "windows", "windows-core 0.61.2", "windows-implement", "windows-interface", @@ -8558,7 +8334,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "381336cfffd772377d291702245447a5251a2ffa5bad679c99e61bc48bacbf9c" dependencies = [ "thiserror 2.0.18", - "windows 0.61.3", + "windows", "windows-core 0.61.2", ] @@ -8584,7 +8360,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -8634,23 +8410,11 @@ version = "0.61.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9babd3a767a4c1aef6900409f85f5d53ce2544ccdfaa86dad48c91782c6d6893" dependencies = [ - "windows-collections 0.2.0", + "windows-collections", "windows-core 0.61.2", - "windows-future 0.2.1", + "windows-future", "windows-link 0.1.3", - "windows-numerics 0.2.0", -] - -[[package]] -name = "windows" -version = "0.62.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "527fadee13e0c05939a6a05d5bd6eec6cd2e3dbd648b9f8e447c6518133d8580" -dependencies = [ - "windows-collections 0.3.2", - "windows-core 0.62.2", - "windows-future 0.3.2", - "windows-numerics 0.3.1", + "windows-numerics", ] [[package]] @@ -8662,15 +8426,6 @@ dependencies = [ "windows-core 0.61.2", ] -[[package]] -name = "windows-collections" -version = "0.3.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "23b2d95af1a8a14a3c7367e1ed4fc9c20e0a26e79551b1454d72583c97cc6610" -dependencies = [ - "windows-core 0.62.2", -] - [[package]] name = "windows-core" version = "0.61.2" @@ -8705,18 +8460,7 @@ checksum = "fc6a41e98427b19fe4b73c550f060b59fa592d7d686537eebf9385621bfbad8e" dependencies = [ "windows-core 0.61.2", "windows-link 0.1.3", - "windows-threading 0.1.0", -] - -[[package]] -name = "windows-future" -version = "0.3.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e1d6f90251fe18a279739e78025bd6ddc52a7e22f921070ccdc67dde84c605cb" -dependencies = [ - "windows-core 0.62.2", - "windows-link 0.2.1", - "windows-threading 0.2.1", + "windows-threading", ] [[package]] @@ -8763,16 +8507,6 @@ dependencies = [ "windows-link 0.1.3", ] -[[package]] -name = "windows-numerics" -version = "0.3.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6e2e40844ac143cdb44aead537bbf727de9b044e107a0f1220392177d15b0f26" -dependencies = [ - "windows-core 0.62.2", - "windows-link 0.2.1", -] - [[package]] name = "windows-result" version = "0.3.4" @@ -8911,15 +8645,6 @@ dependencies = [ "windows-link 0.1.3", ] -[[package]] -name = "windows-threading" -version = "0.2.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3949bd5b99cafdf1c7ca86b43ca564028dfe27d66958f2470940f73d86d75b37" -dependencies = [ - "windows-link 0.2.1", -] - [[package]] name = "windows-version" version = "0.1.7" @@ -9263,7 +8988,7 @@ dependencies = [ "webkit2gtk", "webkit2gtk-sys", "webview2-com", - "windows 0.61.3", + "windows", "windows-core 0.61.2", "windows-version", "x11-dl", diff --git a/crates/openfang-runtime/src/agent_loop.rs b/crates/openfang-runtime/src/agent_loop.rs index f773def4..53c702cf 100644 --- a/crates/openfang-runtime/src/agent_loop.rs +++ b/crates/openfang-runtime/src/agent_loop.rs @@ -86,6 +86,13 @@ fn phantom_action_detected(text: &str) -> bool { has_action && has_channel } +/// Returns true when the agent response text indicates an intentional silent completion. +/// Matches `NO_REPLY` (exact) and `[SILENT]` (case-insensitive). +fn is_silent_token(text: &str) -> bool { + let trimmed = text.trim(); + trimmed == "NO_REPLY" || trimmed.eq_ignore_ascii_case("[silent]") +} + /// Extra guidance injected after failed tool calls to prevent fabricated follow-up actions. const TOOL_ERROR_GUIDANCE: &str = "[System: One or more tool calls failed. Failed tools did not produce usable data. Do NOT invent missing results, cite nonexistent search results, or pretend failed tools succeeded. If your next steps depend on a failed tool, either retry with a materially different approach or explain the failure to the user and stop. Do not write files, store memory, or take downstream actions based on failed tool outputs.]"; @@ -463,8 +470,9 @@ pub async fn run_agent_loop( crate::reply_directives::parse_directives(&text); let text = cleaned_text; - // NO_REPLY: agent intentionally chose not to reply - if text.trim() == "NO_REPLY" || parsed_directives.silent { + // NO_REPLY / [SILENT]: agent intentionally chose not to reply. + // [SILENT] must not be stored literally — it reinforces silence in future turns. + if is_silent_token(&text) || parsed_directives.silent { debug!(agent = %manifest.name, "Agent chose NO_REPLY/silent — silent completion"); session .messages @@ -1641,8 +1649,9 @@ pub async fn run_agent_loop_streaming( crate::reply_directives::parse_directives(&text); let text = cleaned_text_s; - // NO_REPLY: agent intentionally chose not to reply - if text.trim() == "NO_REPLY" || parsed_directives_s.silent { + // NO_REPLY / [SILENT]: agent intentionally chose not to reply. + // [SILENT] must not be stored literally — it reinforces silence in future turns. + if is_silent_token(&text) || parsed_directives_s.silent { debug!(agent = %manifest.name, "Agent chose NO_REPLY/silent (streaming) — silent completion"); session .messages @@ -4648,4 +4657,36 @@ mod tests { } assert!(!events.is_empty(), "Should have received stream events"); } + + #[test] + fn test_silent_detection_uppercase() { + assert!(is_silent_token("[SILENT]")); + } + + #[test] + fn test_silent_detection_lowercase() { + assert!(is_silent_token("[silent]")); + } + + #[test] + fn test_silent_detection_mixed_case() { + assert!(is_silent_token("[Silent]")); + } + + #[test] + fn test_silent_detection_with_whitespace() { + assert!(is_silent_token(" [SILENT] ")); + } + + #[test] + fn test_silent_detection_no_reply() { + assert!(is_silent_token("NO_REPLY")); + } + + #[test] + fn test_silent_detection_rejects_normal_text() { + assert!(!is_silent_token("Hello, how can I help?")); + assert!(!is_silent_token("SILENT")); + assert!(!is_silent_token("")); + } } diff --git a/crates/openfang-runtime/src/mcp.rs b/crates/openfang-runtime/src/mcp.rs index b9f5f381..6792593b 100644 --- a/crates/openfang-runtime/src/mcp.rs +++ b/crates/openfang-runtime/src/mcp.rs @@ -1,20 +1,16 @@ //! MCP (Model Context Protocol) client — connect to external MCP servers. //! -//! Uses the official `rmcp` SDK for protocol handling. Supports: -//! - **stdio**: subprocess with JSON-RPC over stdin/stdout -//! - **sse**: deprecated HTTP+SSE transport (protocol version 2024-11-05) -//! - **http**: Streamable HTTP transport (protocol version 2025-03-26+) +//! MCP uses JSON-RPC 2.0 over stdio or HTTP+SSE. This module lets OpenFang +//! agents use tools from any MCP server (100+ available: GitHub, filesystem, +//! databases, APIs, etc.). //! //! All MCP tools are namespaced with `mcp_{server}_{tool}` to prevent collisions. -use http::{HeaderName, HeaderValue}; use openfang_types::tool::ToolDefinition; -use rmcp::model::{CallToolRequestParams, ClientCapabilities, ClientInfo, Implementation}; -use rmcp::service::RunningService; -use rmcp::{RoleClient, ServiceExt}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; -use std::sync::Arc; +use std::process::Stdio; +use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tracing::{debug, info}; // --------------------------------------------------------------------------- @@ -34,12 +30,6 @@ pub struct McpServerConfig { /// Environment variables to pass through to the subprocess (sandboxed). #[serde(default)] pub env: Vec, - /// Extra HTTP headers to send with every SSE / Streamable-HTTP request. - /// Each entry is `"Header-Name: value"`. Useful for authentication - /// (`Authorization: Bearer `), API keys (`X-Api-Key: ...`), - /// or any custom headers required by a remote MCP server. - #[serde(default)] - pub headers: Vec, } fn default_timeout() -> u64 { @@ -56,14 +46,8 @@ pub enum McpTransport { #[serde(default)] args: Vec, }, - /// Deprecated HTTP+SSE transport (protocol version 2024-11-05). - /// Uses POST for sending and SSE for receiving. + /// HTTP Server-Sent Events. Sse { url: String }, - /// Streamable HTTP transport (MCP 2025-03-26+). - /// Single endpoint, client MUST send Accept: application/json, text/event-stream. - /// Server responds with either JSON or SSE stream. - /// Supports Mcp-Session-Id for session management. - Http { url: String }, } // --------------------------------------------------------------------------- @@ -80,9 +64,59 @@ pub struct McpConnection { /// Needed because `normalize_name` replaces hyphens with underscores, /// but the server expects the original name (e.g. "list-connections"). original_names: HashMap, - /// The rmcp client handle — type-erased because the concrete type - /// depends on which transport was used (stdio vs HTTP). - client: RunningService, + /// Transport handle for sending requests. + transport: McpTransportHandle, + /// Next JSON-RPC request ID. + next_id: u64, +} + +/// Transport handle — abstraction over stdio subprocess or HTTP. +enum McpTransportHandle { + Stdio { + child: Box, + stdin: tokio::process::ChildStdin, + stdout: BufReader, + }, + Sse { + client: reqwest::Client, + url: String, + }, +} + +/// JSON-RPC 2.0 request. +#[derive(Serialize)] +struct JsonRpcRequest { + jsonrpc: &'static str, + id: u64, + method: String, + #[serde(skip_serializing_if = "Option::is_none")] + params: Option, +} + +/// JSON-RPC 2.0 response. +#[derive(Deserialize)] +struct JsonRpcResponse { + #[allow(dead_code)] + jsonrpc: String, + #[allow(dead_code)] + id: Option, + result: Option, + error: Option, +} + +/// JSON-RPC 2.0 error object. +#[derive(Debug, Deserialize)] +pub struct JsonRpcError { + pub code: i64, + pub message: String, + #[allow(dead_code)] + pub data: Option, +} + +impl std::fmt::Display for JsonRpcError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "JSON-RPC error {}: {}", self.code, self.message) + } } // --------------------------------------------------------------------------- @@ -92,17 +126,13 @@ pub struct McpConnection { impl McpConnection { /// Connect to an MCP server, perform handshake, and discover tools. pub async fn connect(config: McpServerConfig) -> Result { - let client_info = ClientInfo::new( - ClientCapabilities::default(), - Implementation::new("openfang", env!("CARGO_PKG_VERSION")), - ); - - let client = match &config.transport { + let transport = match &config.transport { McpTransport::Stdio { command, args } => { - Self::connect_stdio(command, args, &config.env, client_info).await? + Self::connect_stdio(command, args, &config.env).await? } - McpTransport::Sse { url } | McpTransport::Http { url } => { - Self::connect_http(url, &config.headers, client_info).await? + McpTransport::Sse { url } => { + // SSRF check: reject private/localhost URLs unless explicitly configured + Self::connect_sse(url).await? } }; @@ -110,9 +140,13 @@ impl McpConnection { config, tools: Vec::new(), original_names: HashMap::new(), - client, + transport, + next_id: 1, }; + // Initialize handshake + conn.initialize().await?; + // Discover tools conn.discover_tools().await?; @@ -125,34 +159,76 @@ impl McpConnection { Ok(conn) } + /// Send the MCP `initialize` handshake. + async fn initialize(&mut self) -> Result<(), String> { + let params = serde_json::json!({ + "protocolVersion": "2024-11-05", + "capabilities": {}, + "clientInfo": { + "name": "openfang", + "version": env!("CARGO_PKG_VERSION") + } + }); + + let response = self.send_request("initialize", Some(params)).await?; + + if let Some(result) = response { + debug!( + server = %self.config.name, + server_info = %result, + "MCP initialize response" + ); + } + + // Send initialized notification (no response expected) + self.send_notification("notifications/initialized", None) + .await?; + + Ok(()) + } + /// Discover available tools via `tools/list`. async fn discover_tools(&mut self) -> Result<(), String> { - let tools = self - .client - .list_all_tools() - .await - .map_err(|e| format!("Failed to list MCP tools: {e}"))?; + let response = self.send_request("tools/list", None).await?; - let server_name = &self.config.name; - for tool in &tools { - let raw_name = &tool.name; - let description = tool.description.as_deref().unwrap_or(""); + if let Some(result) = response { + if let Some(tools_array) = result.get("tools").and_then(|t| t.as_array()) { + let server_name = &self.config.name; + for tool in tools_array { + let raw_name = tool["name"].as_str().unwrap_or("unnamed"); + let description = tool["description"].as_str().unwrap_or(""); + let input_schema = tool + .get("inputSchema") + .cloned() + .and_then(|v| { + // Ensure input_schema is a JSON object. MCP servers may + // return it as a string, null, or omit it entirely. + match &v { + serde_json::Value::Object(_) => Some(v), + serde_json::Value::String(s) => { + serde_json::from_str::(s) + .ok() + .filter(|p| p.is_object()) + } + _ => None, + } + }) + .unwrap_or(serde_json::json!({"type": "object"})); - let input_schema = serde_json::to_value(&tool.input_schema) - .unwrap_or(serde_json::json!({"type": "object"})); + // Namespace: mcp_{server}_{tool} + let namespaced = format_mcp_tool_name(server_name, raw_name); - // Namespace: mcp_{server}_{tool} - let namespaced = format_mcp_tool_name(server_name, raw_name); + // Store original name so we can send it back to the server + self.original_names + .insert(namespaced.clone(), raw_name.to_string()); - // Store original name so we can send it back to the server - self.original_names - .insert(namespaced.clone(), raw_name.to_string()); - - self.tools.push(ToolDefinition { - name: namespaced, - description: format!("[MCP:{server_name}] {description}"), - input_schema, - }); + self.tools.push(ToolDefinition { + name: namespaced, + description: format!("[MCP:{server_name}] {description}"), + input_schema, + }); + } + } } Ok(()) @@ -167,38 +243,40 @@ impl McpConnection { arguments: &serde_json::Value, ) -> Result { // Look up the original tool name from the server (preserves hyphens etc.) - let raw_name: String = self + let raw_name = self .original_names .get(name) - .cloned() - .or_else(|| strip_mcp_prefix(&self.config.name, name).map(|s| s.to_string())) - .unwrap_or_else(|| name.to_string()); + .map(|s| s.as_str()) + .or_else(|| strip_mcp_prefix(&self.config.name, name)) + .unwrap_or(name); - let args = arguments.as_object().cloned().unwrap_or_default(); + let params = serde_json::json!({ + "name": raw_name, + "arguments": arguments, + }); - debug!(tool = %raw_name, server = %self.config.name, "MCP tool call"); + let response = self.send_request("tools/call", Some(params)).await?; - let params = CallToolRequestParams::new(raw_name).with_arguments(args); - - let result = self - .client - .call_tool(params) - .await - .map_err(|e| format!("MCP tool call failed: {e}"))?; - - // Extract text content from the response. - // `Content` is `Annotated` which Derefs to `RawContent`. - let texts: Vec<&str> = result - .content - .iter() - .filter_map(|item| item.as_text().map(|tc| tc.text.as_str())) - .collect(); - - if texts.is_empty() { - // Fallback: serialize the entire result - Ok(serde_json::to_string(&result).unwrap_or_default()) - } else { - Ok(texts.join("\n")) + match response { + Some(result) => { + // Extract text content from the response + if let Some(content) = result.get("content").and_then(|c| c.as_array()) { + let texts: Vec<&str> = content + .iter() + .filter_map(|item| { + if item["type"].as_str() == Some("text") { + item["text"].as_str() + } else { + None + } + }) + .collect(); + Ok(texts.join("\n")) + } else { + Ok(result.to_string()) + } + } + None => Err("No result from MCP tools/call".to_string()), } } @@ -212,124 +290,299 @@ impl McpConnection { &self.config.name } - // -- Transport constructors ----------------------------------------------- + // --- Transport helpers --- + + async fn send_request( + &mut self, + method: &str, + params: Option, + ) -> Result, String> { + let id = self.next_id; + self.next_id += 1; + + let request = JsonRpcRequest { + jsonrpc: "2.0", + id, + method: method.to_string(), + params, + }; + + let request_json = serde_json::to_string(&request) + .map_err(|e| format!("Failed to serialize request: {e}"))?; + + debug!(method, id, "MCP request"); + + match &mut self.transport { + McpTransportHandle::Stdio { stdin, stdout, .. } => { + // Write request + newline + stdin + .write_all(request_json.as_bytes()) + .await + .map_err(|e| format!("Failed to write to MCP stdin: {e}"))?; + stdin + .write_all(b"\n") + .await + .map_err(|e| format!("Failed to write newline: {e}"))?; + stdin + .flush() + .await + .map_err(|e| format!("Failed to flush stdin: {e}"))?; + + // Read response lines until we find one matching our request ID. + // MCP servers may send notifications or log lines before the + // actual response. + let timeout_dur = tokio::time::Duration::from_secs(self.config.timeout_secs); + let deadline = tokio::time::Instant::now() + timeout_dur; + + loop { + let mut line = String::new(); + let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); + if remaining.is_zero() { + return Err("MCP request timed out".to_string()); + } + match tokio::time::timeout(remaining, stdout.read_line(&mut line)).await { + Ok(Ok(0)) => return Err("MCP server closed connection".to_string()), + Ok(Ok(_)) => {} + Ok(Err(e)) => return Err(format!("Failed to read MCP response: {e}")), + Err(_) => return Err("MCP request timed out".to_string()), + } + + let trimmed = line.trim(); + if trimmed.is_empty() { + continue; + } + + // Try to parse as JSON-RPC response + let parsed: Result = serde_json::from_str(trimmed); + match parsed { + Ok(response) if response.id == Some(id) => { + if let Some(err) = response.error { + return Err(format!("{err}")); + } + return Ok(response.result); + } + Ok(_other) => { + // Response for a different ID or notification — skip + debug!("MCP: skipping non-matching response line"); + continue; + } + Err(_) => { + // Not valid JSON-RPC — skip (could be a log line) + debug!("MCP: skipping non-JSON line from server"); + continue; + } + } + } + } + McpTransportHandle::Sse { client, url } => { + let response = client + .post(url.as_str()) + .json(&request) + .timeout(std::time::Duration::from_secs(self.config.timeout_secs)) + .send() + .await + .map_err(|e| format!("MCP SSE request failed: {e}"))?; + + if !response.status().is_success() { + return Err(format!("MCP SSE returned {}", response.status())); + } + + let body = response + .text() + .await + .map_err(|e| format!("Failed to read SSE response: {e}"))?; + + // Handle Streamable HTTP MCP responses that use SSE framing + // (e.g. "event: message\ndata: {...}\n\n"). Extract the JSON + // from the last `data:` line if the body looks like SSE. + let json_body = if body.trim_start().starts_with("event:") + || body.trim_start().starts_with("data:") + { + body.lines() + .filter_map(|line| { + line.strip_prefix("data: ") + .or_else(|| line.strip_prefix("data:")) + .filter(|s| !s.is_empty()) + }) + .next_back() + .unwrap_or(&body) + .to_string() + } else { + body + }; + + let rpc_response: JsonRpcResponse = serde_json::from_str(&json_body) + .map_err(|e| format!("Invalid MCP SSE JSON-RPC response: {e}"))?; + + if let Some(err) = rpc_response.error { + return Err(format!("{err}")); + } + + Ok(rpc_response.result) + } + } + } + + async fn send_notification( + &mut self, + method: &str, + params: Option, + ) -> Result<(), String> { + let notification = serde_json::json!({ + "jsonrpc": "2.0", + "method": method, + "params": params.unwrap_or(serde_json::json!({})), + }); + + let json = serde_json::to_string(¬ification) + .map_err(|e| format!("Failed to serialize notification: {e}"))?; + + match &mut self.transport { + McpTransportHandle::Stdio { stdin, .. } => { + stdin + .write_all(json.as_bytes()) + .await + .map_err(|e| format!("Write notification: {e}"))?; + stdin + .write_all(b"\n") + .await + .map_err(|e| format!("Write newline: {e}"))?; + stdin.flush().await.map_err(|e| format!("Flush: {e}"))?; + } + McpTransportHandle::Sse { client, url } => { + let _ = client.post(url.as_str()).json(¬ification).send().await; + } + } + + Ok(()) + } - /// Connect using stdio transport (subprocess). async fn connect_stdio( command: &str, args: &[String], env_whitelist: &[String], - client_info: ClientInfo, - ) -> Result, String> { - use rmcp::transport::{ConfigureCommandExt, TokioChildProcess}; - use tokio::process::Command; - + ) -> Result { // Validate command path (no path traversal) if command.contains("..") { return Err("MCP command path contains '..': rejected".to_string()); } - let cmd_str = command.to_string(); - let args_vec: Vec = args.to_vec(); - let env_list: Vec = env_whitelist.to_vec(); - - let transport = TokioChildProcess::new(Command::new(&cmd_str).configure(move |cmd| { - for arg in &args_vec { - cmd.arg(arg); - } - // Sandbox: clear environment, only pass whitelisted vars - cmd.env_clear(); - for var_name in &env_list { - if let Ok(val) = std::env::var(var_name) { - cmd.env(var_name, val); + // On Windows, npm/npx install as .cmd batch wrappers. Detect and adapt. + let resolved_command: String = if cfg!(windows) { + // If the user already specified .cmd/.bat, use as-is + if command.ends_with(".cmd") || command.ends_with(".bat") { + command.to_string() + } else { + // Check if the .cmd variant exists on PATH + let cmd_variant = format!("{command}.cmd"); + let has_cmd = std::env::var("PATH") + .unwrap_or_default() + .split(';') + .any(|dir| std::path::Path::new(dir).join(&cmd_variant).exists()); + if has_cmd { + cmd_variant + } else { + command.to_string() } } - // Always pass PATH for binary resolution - if let Ok(path) = std::env::var("PATH") { - cmd.env("PATH", path); - } - // On Windows, npm/node need extra vars - if cfg!(windows) { - for var in &[ - "APPDATA", - "LOCALAPPDATA", - "USERPROFILE", - "SystemRoot", - "TEMP", - "TMP", - "HOME", - "HOMEDRIVE", - "HOMEPATH", - ] { - if let Ok(val) = std::env::var(var) { - cmd.env(var, val); - } - } - } - })) - .map_err(|e| format!("Failed to spawn MCP server '{cmd_str}': {e}"))?; - - let client = client_info - .serve(transport) - .await - .map_err(|e| format!("MCP stdio handshake failed: {e}"))?; - - Ok(client) - } - - /// Connect using Streamable HTTP transport (or SSE fallback via the same endpoint). - /// - /// The `rmcp` SDK's `StreamableHttpClientTransport` handles the full - /// Streamable HTTP protocol: Accept headers, Mcp-Session-Id tracking, - /// SSE stream parsing, and content-type negotiation. - async fn connect_http( - url: &str, - headers: &[String], - client_info: ClientInfo, - ) -> Result, String> { - use rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig; - use rmcp::transport::StreamableHttpClientTransport; - - Self::check_ssrf(url)?; - - // Parse custom headers (e.g., "Authorization: Bearer "). - let mut custom_headers: HashMap = HashMap::new(); - for header_str in headers { - if let Some((name, value)) = header_str.split_once(':') { - let name = name.trim(); - let value = value.trim(); - if let (Ok(hn), Ok(hv)) = ( - HeaderName::from_bytes(name.as_bytes()), - HeaderValue::from_str(value), - ) { - custom_headers.insert(hn, hv); - } - } - } - - let config = StreamableHttpClientTransportConfig { - uri: Arc::from(url), - custom_headers, - ..Default::default() + } else { + command.to_string() }; - let transport = StreamableHttpClientTransport::from_config(config); + let mut cmd = tokio::process::Command::new(&resolved_command); + cmd.args(args); + cmd.stdin(Stdio::piped()); + cmd.stdout(Stdio::piped()); + cmd.stderr(Stdio::piped()); - let client = client_info - .serve(transport) - .await - .map_err(|e| format!("MCP HTTP connection failed: {e}"))?; + // Sandbox: clear environment, only pass whitelisted vars + cmd.env_clear(); + for var_name in env_whitelist { + if let Ok(val) = std::env::var(var_name) { + cmd.env(var_name, val); + } + } + // Always pass PATH for binary resolution + if let Ok(path) = std::env::var("PATH") { + cmd.env("PATH", path); + } + // On Windows, npm/node need APPDATA, USERPROFILE, LOCALAPPDATA, and SystemRoot + if cfg!(windows) { + for var in &[ + "APPDATA", + "LOCALAPPDATA", + "USERPROFILE", + "SystemRoot", + "TEMP", + "TMP", + "HOME", + "HOMEDRIVE", + "HOMEPATH", + ] { + if let Ok(val) = std::env::var(var) { + cmd.env(var, val); + } + } + } - Ok(client) + let mut child = cmd + .spawn() + .map_err(|e| format!("Failed to spawn MCP server '{resolved_command}': {e}"))?; + + // Log stderr in background for debugging MCP server issues + if let Some(stderr) = child.stderr.take() { + let cmd_name = resolved_command.clone(); + tokio::spawn(async move { + use tokio::io::AsyncBufReadExt; + let reader = tokio::io::BufReader::new(stderr); + let mut lines = reader.lines(); + while let Ok(Some(line)) = lines.next_line().await { + tracing::debug!(mcp_server = %cmd_name, "stderr: {line}"); + } + }); + } + + let stdin = child + .stdin + .take() + .ok_or("Failed to capture MCP server stdin")?; + let stdout = child + .stdout + .take() + .ok_or("Failed to capture MCP server stdout")?; + + Ok(McpTransportHandle::Stdio { + child: Box::new(child), + stdin, + stdout: BufReader::new(stdout), + }) } - /// Basic SSRF check: reject obviously private/metadata URLs. - fn check_ssrf(url: &str) -> Result<(), String> { + async fn connect_sse(url: &str) -> Result { + // Basic SSRF check: reject obviously private URLs let lower = url.to_lowercase(); if lower.contains("169.254.169.254") || lower.contains("metadata.google") { - return Err("SSRF: MCP URL targets metadata endpoint".to_string()); + return Err("SSRF: MCP SSE URL targets metadata endpoint".to_string()); + } + + let client = reqwest::Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .build() + .map_err(|e| format!("Failed to create HTTP client: {e}"))?; + + Ok(McpTransportHandle::Sse { + client, + url: url.to_string(), + }) + } +} + +impl Drop for McpConnection { + fn drop(&mut self) { + if let McpTransportHandle::Stdio { ref mut child, .. } = self.transport { + // Best-effort kill of the subprocess + let _ = child.start_kill(); } - Ok(()) } } @@ -420,12 +673,16 @@ mod tests { #[test] fn test_hyphenated_tool_name_preserved() { + // Tool names with hyphens get normalized to underscores for namespacing, + // but original_names map preserves the original for call_tool dispatch. let namespaced = format_mcp_tool_name("sqlcl", "list-connections"); assert_eq!(namespaced, "mcp_sqlcl_list_connections"); + // Simulate what discover_tools does let mut original_names = HashMap::new(); original_names.insert(namespaced.clone(), "list-connections".to_string()); + // call_tool should resolve to original hyphenated name let raw = original_names .get(&namespaced) .map(|s| s.as_str()) @@ -444,21 +701,25 @@ mod tests { #[test] fn test_extract_mcp_server_from_known_with_hyphens() { + // Server "bocha-search" normalized to "bocha_search" in tool prefix let servers = vec!["bocha-search", "github"]; let tool = "mcp_bocha_search_bocha_web_search"; assert_eq!( extract_mcp_server_from_known(tool, &servers), Some("bocha-search") ); + // Simple server name still works assert_eq!( extract_mcp_server_from_known("mcp_github_create_issue", &servers), Some("github") ); + // Non-MCP tool returns None assert_eq!(extract_mcp_server_from_known("file_read", &servers), None); } #[test] fn test_extract_mcp_server_from_known_longest_match() { + // "my-api" and "my-api-v2" — should match the longer one let servers = vec!["my-api", "my-api-v2"]; assert_eq!( extract_mcp_server_from_known("mcp_my_api_v2_get_users", &servers), @@ -470,6 +731,60 @@ mod tests { ); } + #[test] + fn test_mcp_jsonrpc_initialize() { + // Verify the initialize request structure + let request = JsonRpcRequest { + jsonrpc: "2.0", + id: 1, + method: "initialize".to_string(), + params: Some(serde_json::json!({ + "protocolVersion": "2024-11-05", + "capabilities": {}, + "clientInfo": { + "name": "openfang", + "version": "0.1.0" + } + })), + }; + let json = serde_json::to_string(&request).unwrap(); + assert!(json.contains("initialize")); + assert!(json.contains("protocolVersion")); + assert!(json.contains("openfang")); + } + + #[test] + fn test_mcp_jsonrpc_tools_list() { + // Simulate a tools/list response + let response_json = r#"{ + "jsonrpc": "2.0", + "id": 2, + "result": { + "tools": [ + { + "name": "create_issue", + "description": "Create a GitHub issue", + "inputSchema": { + "type": "object", + "properties": { + "title": {"type": "string"}, + "body": {"type": "string"} + }, + "required": ["title"] + } + } + ] + } + }"#; + + let response: JsonRpcResponse = serde_json::from_str(response_json).unwrap(); + assert!(response.error.is_none()); + let result = response.result.unwrap(); + let tools = result["tools"].as_array().unwrap(); + assert_eq!(tools.len(), 1); + assert_eq!(tools[0]["name"].as_str().unwrap(), "create_issue"); + } + #[test] fn test_mcp_transport_config_serde() { let config = McpServerConfig { @@ -483,7 +798,6 @@ mod tests { }, timeout_secs: 30, env: vec!["GITHUB_PERSONAL_ACCESS_TOKEN".to_string()], - headers: vec![], }; let json = serde_json::to_string(&config).unwrap(); @@ -508,7 +822,6 @@ mod tests { }, timeout_secs: 60, env: vec![], - headers: vec![], }; let json = serde_json::to_string(&sse_config).unwrap(); let back: McpServerConfig = serde_json::from_str(&json).unwrap(); @@ -516,24 +829,5 @@ mod tests { McpTransport::Sse { url } => assert_eq!(url, "https://example.com/mcp"), _ => panic!("Expected SSE transport"), } - - // HTTP (Streamable HTTP) variant - let http_config = McpServerConfig { - name: "atlassian".to_string(), - transport: McpTransport::Http { - url: "https://mcp.atlassian.com/v1/mcp".to_string(), - }, - timeout_secs: 120, - env: vec![], - headers: vec!["Authorization: Bearer test-token-456".to_string()], - }; - let json = serde_json::to_string(&http_config).unwrap(); - let back: McpServerConfig = serde_json::from_str(&json).unwrap(); - match back.transport { - McpTransport::Http { url } => { - assert_eq!(url, "https://mcp.atlassian.com/v1/mcp") - } - _ => panic!("Expected Http transport"), - } } } diff --git a/crates/openfang-runtime/src/tool_runner.rs b/crates/openfang-runtime/src/tool_runner.rs index 664735c8..f6c024dc 100644 --- a/crates/openfang-runtime/src/tool_runner.rs +++ b/crates/openfang-runtime/src/tool_runner.rs @@ -3677,13 +3677,12 @@ mod tests { None, // process_manager ) .await; - // Should NOT be the capability-check denial — it should normalize to file_write - // and pass the capability check. It may fail for other reasons (path validation, - // OS-level errors), but not the agent capability gate. + // Should NOT be the capability-enforcement "Permission denied" — it should + // normalize to file_write and pass the capability check. It may still fail + // for filesystem reasons (e.g. OS "Permission denied (os error 13)"), so we + // check specifically for the capability-gate message. assert!( - !result - .content - .contains("does not have capability to use tool"), + !result.content.contains("Permission denied: agent"), "fs-write should normalize to file_write and pass capability check, got: {}", result.content ); From 4b5aba28cfceb98ff1529abb71e2dacfb7125f2c Mon Sep 17 00:00:00 2001 From: Philippe Branchu Date: Sat, 28 Mar 2026 00:42:52 +0000 Subject: [PATCH 2/5] Revert unrelated mcp.rs rewrite and Cargo.lock changes The previous commit accidentally included a complete MCP module rewrite that removed Http transport and headers support, breaking compilation against upstream kernel.rs tests. Co-Authored-By: Claude Opus 4.6 (1M context) --- Cargo.lock | 423 ++++++++++++++--- crates/openfang-runtime/src/mcp.rs | 698 +++++++++-------------------- 2 files changed, 551 insertions(+), 570 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 4621940c..df329192 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -139,7 +139,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -150,7 +150,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -384,6 +384,28 @@ version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" +[[package]] +name = "aws-lc-rs" +version = "1.16.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a054912289d18629dc78375ba2c3726a3afe3ff71b4edba9dedfca0e3446d1fc" +dependencies = [ + "aws-lc-sys", + "zeroize", +] + +[[package]] +name = "aws-lc-sys" +version = "0.39.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fa7e52a4c5c547c741610a2c6f123f3881e409b714cd27e6798ef020c514f0a" +dependencies = [ + "cc", + "cmake", + "dunce", + "fs_extra", +] + [[package]] name = "axum" version = "0.8.8" @@ -860,6 +882,15 @@ version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" +[[package]] +name = "cmake" +version = "0.1.57" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75443c44cd6b379beb8c5b45d85d0773baf31cce901fe7bb252f4eff3008ef7d" +dependencies = [ + "cc", +] + [[package]] name = "cobs" version = "0.3.0" @@ -881,7 +912,7 @@ version = "3.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "faf9468729b8cbcea668e36183cb69d317348c2e08e994829fb56ebfdfbaac34" dependencies = [ - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -957,6 +988,16 @@ dependencies = [ "version_check", ] +[[package]] +name = "core-foundation" +version = "0.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91e195e091a93c46f7102ec7818a2aa394e1e1771c3ab4825963fa03e45afb8f" +dependencies = [ + "core-foundation-sys", + "libc", +] + [[package]] name = "core-foundation" version = "0.10.1" @@ -980,7 +1021,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "064badf302c3194842cf2c5d61f56cc88e54a759313879cdf03abdd27d0c3b97" dependencies = [ "bitflags 2.11.0", - "core-foundation", + "core-foundation 0.10.1", "core-graphics-types", "foreign-types 0.5.0", "libc", @@ -993,7 +1034,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3d44a101f213f6c4cdc1853d4b78aef6db6bdfa3468798cc1d9912f4735013eb" dependencies = [ "bitflags 2.11.0", - "core-foundation", + "core-foundation 0.10.1", "libc", ] @@ -1538,7 +1579,7 @@ dependencies = [ "libc", "option-ext", "redox_users 0.5.2", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -1788,7 +1829,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1889,6 +1930,17 @@ dependencies = [ "zlib-rs", ] +[[package]] +name = "flume" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" +dependencies = [ + "futures-core", + "futures-sink", + "spin", +] + [[package]] name = "fnv" version = "1.0.7" @@ -1958,6 +2010,12 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "fs_extra" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" + [[package]] name = "futf" version = "0.1.5" @@ -2405,9 +2463,9 @@ dependencies = [ [[package]] name = "governor" -version = "0.8.1" +version = "0.10.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "be93b4ec2e4710b04d9264c0c7350cdd62a8c20e5e4ac732552ebb8f0debe8eb" +checksum = "9efcab3c1958580ff1f25a2a41be1668f7603d849bb63af523b208a3cc1223b8" dependencies = [ "cfg-if", "dashmap", @@ -2415,7 +2473,7 @@ dependencies = [ "futures-timer", "futures-util", "getrandom 0.3.4", - "no-std-compat", + "hashbrown 0.16.1", "nonzero_ext", "parking_lot", "portable-atomic", @@ -2511,6 +2569,11 @@ name = "hashbrown" version = "0.16.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "841d1cc9bed7f9236f321df977030373f4a4163ae1a7dbfe1a51a2c1a51d9100" +dependencies = [ + "allocator-api2", + "equivalent", + "foldhash 0.2.0", +] [[package]] name = "hashlink" @@ -2672,10 +2735,10 @@ dependencies = [ "http", "hyper", "hyper-util", - "rustls", + "rustls 0.23.37", "rustls-pki-types", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tower-service", "webpki-roots", ] @@ -3206,10 +3269,10 @@ dependencies = [ "nom 8.0.0", "percent-encoding", "quoted_printable", - "rustls", + "rustls 0.23.37", "socket2 0.6.3", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "url", "webpki-roots", ] @@ -3564,10 +3627,10 @@ dependencies = [ "libc", "log", "openssl", - "openssl-probe", + "openssl-probe 0.2.1", "openssl-sys", "schannel", - "security-framework", + "security-framework 3.7.0", "security-framework-sys", "tempfile", ] @@ -3609,10 +3672,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086" [[package]] -name = "no-std-compat" -version = "0.4.1" +name = "nix" +version = "0.31.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b93853da6d84c2e3c7d730d6473e8817692dd89be387eb01b94d7f108ecb5b8c" +checksum = "5d6d0705320c1e6ba1d912b5e37cf18071b6c2e9b7fa8215a1e8a7651966f5d3" +dependencies = [ + "bitflags 2.11.0", + "cfg-if", + "cfg_aliases", + "libc", +] [[package]] name = "nodrop" @@ -3666,7 +3735,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -3917,7 +3986,7 @@ dependencies = [ "tokio", "tokio-stream", "tokio-test", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tower", "tower-http", "tracing", @@ -3948,6 +4017,7 @@ dependencies = [ "regex-lite", "reqwest 0.12.28", "roxmltree", + "rumqttc", "serde", "serde_json", "sha1", @@ -3983,7 +4053,7 @@ dependencies = [ "serde_json", "tempfile", "tokio", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "tracing-subscriber", "uuid", @@ -4011,7 +4081,7 @@ dependencies = [ "tauri-plugin-single-instance", "tauri-plugin-updater", "tokio", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "tracing-subscriber", ] @@ -4037,7 +4107,7 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tokio-test", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "url", "uuid", @@ -4056,7 +4126,7 @@ dependencies = [ "tempfile", "thiserror 2.0.18", "tokio-test", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "uuid", ] @@ -4092,7 +4162,7 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tokio-test", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "tracing-subscriber", "uuid", @@ -4132,7 +4202,7 @@ dependencies = [ "serde_yaml", "tempfile", "thiserror 2.0.18", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "uuid", "walkdir", @@ -4150,11 +4220,13 @@ dependencies = [ "dashmap", "futures", "hex", + "http", "openfang-memory", "openfang-skills", "openfang-types", "regex-lite", "reqwest 0.12.28", + "rmcp", "rusqlite", "serde", "serde_json", @@ -4188,7 +4260,7 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tokio-test", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "tracing", "uuid", "walkdir", @@ -4210,7 +4282,7 @@ dependencies = [ "serde_json", "sha2", "thiserror 2.0.18", - "toml 0.8.2", + "toml 0.9.12+spec-1.1.0", "uuid", ] @@ -4262,6 +4334,12 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "openssl-probe" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d05e27ee213611ffe7d6348b942e8f942b37114c00cc03cec254295a4a17852e" + [[package]] name = "openssl-probe" version = "0.2.1" @@ -4313,7 +4391,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7d8fae84b431384b68627d0f9b3b1245fcf9f46f6c0e3dc902e9dce64edd1967" dependencies = [ "libc", - "windows-sys 0.45.0", + "windows-sys 0.61.2", ] [[package]] @@ -4907,6 +4985,20 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "process-wrap" +version = "9.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e842efad9119158434d193c6682e2ebee4b44d6ad801d7b349623b3f57cdf55" +dependencies = [ + "futures", + "indexmap 2.13.0", + "nix", + "tokio", + "tracing", + "windows 0.62.2", +] + [[package]] name = "prost" version = "0.13.5" @@ -5014,7 +5106,7 @@ dependencies = [ "quinn-proto", "quinn-udp", "rustc-hash", - "rustls", + "rustls 0.23.37", "socket2 0.6.3", "thiserror 2.0.18", "tokio", @@ -5028,13 +5120,14 @@ version = "0.11.14" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "434b42fec591c96ef50e21e886936e66d3cc3f737104fdb9b737c40ffb94c098" dependencies = [ + "aws-lc-rs", "bytes", "getrandom 0.3.4", "lru-slab", "rand 0.9.2", "ring", "rustc-hash", - "rustls", + "rustls 0.23.37", "rustls-pki-types", "slab", "thiserror 2.0.18", @@ -5391,14 +5484,14 @@ dependencies = [ "percent-encoding", "pin-project-lite", "quinn", - "rustls", + "rustls 0.23.37", "rustls-pki-types", "serde", "serde_json", "serde_urlencoded", "sync_wrapper", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-util", "tower", "tower-http", @@ -5431,14 +5524,15 @@ dependencies = [ "log", "percent-encoding", "pin-project-lite", - "rustls", + "quinn", + "rustls 0.23.37", "rustls-pki-types", "rustls-platform-verifier", "serde", "serde_json", "sync_wrapper", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-util", "tower", "tower-http", @@ -5488,6 +5582,29 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rmcp" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba6b9d2f0efe2258b23767f1f9e0054cfbcac9c2d6f81a031214143096d7864f" +dependencies = [ + "async-trait", + "chrono", + "futures", + "http", + "pin-project-lite", + "process-wrap", + "reqwest 0.13.2", + "serde", + "serde_json", + "sse-stream", + "thiserror 2.0.18", + "tokio", + "tokio-stream", + "tokio-util", + "tracing", +] + [[package]] name = "rmp" version = "0.8.15" @@ -5516,6 +5633,24 @@ dependencies = [ "memchr", ] +[[package]] +name = "rumqttc" +version = "0.24.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e1568e15fab2d546f940ed3a21f48bbbd1c494c90c99c4481339364a497f94a9" +dependencies = [ + "bytes", + "flume", + "futures-util", + "log", + "rustls-native-certs 0.7.3", + "rustls-pemfile", + "rustls-webpki 0.102.8", + "thiserror 1.0.69", + "tokio", + "tokio-rustls 0.25.0", +] + [[package]] name = "rusqlite" version = "0.31.0" @@ -5575,7 +5710,21 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.52.0", + "windows-sys 0.61.2", +] + +[[package]] +name = "rustls" +version = "0.22.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bf4ef73721ac7bcd79b2b315da7779d8fc09718c6b3d2d1b2d94850eb8c18432" +dependencies = [ + "log", + "ring", + "rustls-pki-types", + "rustls-webpki 0.102.8", + "subtle", + "zeroize", ] [[package]] @@ -5584,25 +5733,48 @@ version = "0.23.37" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "758025cb5fccfd3bc2fd74708fd4682be41d99e5dff73c377c0646c6012c73a4" dependencies = [ + "aws-lc-rs", "log", "once_cell", "ring", "rustls-pki-types", - "rustls-webpki", + "rustls-webpki 0.103.10", "subtle", "zeroize", ] +[[package]] +name = "rustls-native-certs" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5bfb394eeed242e909609f56089eecfe5fda225042e8b171791b9c95f5931e5" +dependencies = [ + "openssl-probe 0.1.6", + "rustls-pemfile", + "rustls-pki-types", + "schannel", + "security-framework 2.11.1", +] + [[package]] name = "rustls-native-certs" version = "0.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "612460d5f7bea540c490b2b6395d8e34a953e52b491accd6c86c8164c5932a63" dependencies = [ - "openssl-probe", + "openssl-probe 0.2.1", "rustls-pki-types", "schannel", - "security-framework", + "security-framework 3.7.0", +] + +[[package]] +name = "rustls-pemfile" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dce314e5fee3f39953d46bb63bb8a46d40c2f8fb7cc5a3b6cab2bde9721d6e50" +dependencies = [ + "rustls-pki-types", ] [[package]] @@ -5621,19 +5793,19 @@ version = "0.6.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d99feebc72bae7ab76ba994bb5e121b8d83d910ca40b36e0921f53becc41784" dependencies = [ - "core-foundation", + "core-foundation 0.10.1", "core-foundation-sys", "jni", "log", "once_cell", - "rustls", - "rustls-native-certs", + "rustls 0.23.37", + "rustls-native-certs 0.8.3", "rustls-platform-verifier-android", - "rustls-webpki", - "security-framework", + "rustls-webpki 0.103.10", + "security-framework 3.7.0", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -5644,15 +5816,27 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.9" +version = "0.102.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d7df23109aa6c1567d1c575b9952556388da57401e4ace1d15f79eedad0d8f53" +checksum = "64ca1bc8749bd4cf37b5ce386cc146580777b4e8572c7b97baf22c83f444bee9" dependencies = [ "ring", "rustls-pki-types", "untrusted", ] +[[package]] +name = "rustls-webpki" +version = "0.103.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df33b2b81ac578cabaf06b89b0631153a3f416b0a886e8a7a1707fb51abbd1ef" +dependencies = [ + "aws-lc-rs", + "ring", + "rustls-pki-types", + "untrusted", +] + [[package]] name = "rustversion" version = "1.0.22" @@ -5740,6 +5924,19 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" +[[package]] +name = "security-framework" +version = "2.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "897b2245f0b511c87893af39b033e5ca9cce68824c4d7e7630b5a1d339658d02" +dependencies = [ + "bitflags 2.11.0", + "core-foundation 0.9.4", + "core-foundation-sys", + "libc", + "security-framework-sys", +] + [[package]] name = "security-framework" version = "3.7.0" @@ -5747,7 +5944,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d" dependencies = [ "bitflags 2.11.0", - "core-foundation", + "core-foundation 0.10.1", "core-foundation-sys", "libc", "security-framework-sys", @@ -6178,7 +6375,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3a766e1110788c36f4fa1c2b71b387a7815aa65f88ce0229841826633d93723e" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -6234,6 +6431,9 @@ name = "spin" version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +dependencies = [ + "lock_api", +] [[package]] name = "spinning_top" @@ -6254,6 +6454,19 @@ dependencies = [ "der", ] +[[package]] +name = "sse-stream" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eb4dc4d33c68ec1f27d386b5610a351922656e1fdf5c05bbaad930cd1519479a" +dependencies = [ + "bytes", + "futures-util", + "http-body", + "http-body-util", + "pin-project-lite", +] + [[package]] name = "stable_deref_trait" version = "1.2.1" @@ -6436,7 +6649,7 @@ checksum = "6e06d52c379e63da659a483a958110bbde891695a0ecb53e48cc7786d5eda7bb" dependencies = [ "bitflags 2.11.0", "block2", - "core-foundation", + "core-foundation 0.10.1", "core-graphics", "crossbeam-channel", "dispatch2", @@ -6460,7 +6673,7 @@ dependencies = [ "tao-macros", "unicode-segmentation", "url", - "windows", + "windows 0.61.3", "windows-core 0.61.2", "windows-version", "x11-dl", @@ -6479,9 +6692,9 @@ dependencies = [ [[package]] name = "tar" -version = "0.4.44" +version = "0.4.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1d863878d212c87a19c1a610eb53bb01fe12951c0501cf5a0d65f724914a667a" +checksum = "22692a6476a21fa75fdfc11d452fda482af402c008cdbaf3476414e122040973" dependencies = [ "filetime", "libc", @@ -6549,7 +6762,7 @@ dependencies = [ "webkit2gtk", "webview2-com", "window-vibrancy", - "windows", + "windows 0.61.3", ] [[package]] @@ -6773,7 +6986,7 @@ dependencies = [ "osakit", "percent-encoding", "reqwest 0.13.2", - "rustls", + "rustls 0.23.37", "semver", "serde", "serde_json", @@ -6811,7 +7024,7 @@ dependencies = [ "url", "webkit2gtk", "webview2-com", - "windows", + "windows 0.61.3", ] [[package]] @@ -6836,7 +7049,7 @@ dependencies = [ "url", "webkit2gtk", "webview2-com", - "windows", + "windows 0.61.3", "wry", ] @@ -6897,7 +7110,7 @@ checksum = "0b1e66e07de489fe43a46678dd0b8df65e0c973909df1b60ba33874e297ba9b9" dependencies = [ "quick-xml 0.37.5", "thiserror 2.0.18", - "windows", + "windows 0.61.3", "windows-version", ] @@ -6911,7 +7124,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -7067,13 +7280,24 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "tokio-rustls" +version = "0.25.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "775e0c0f0adb3a2f22a00c4745d728b479985fc15ee7ca6a2608388c5569860f" +dependencies = [ + "rustls 0.22.4", + "rustls-pki-types", + "tokio", +] + [[package]] name = "tokio-rustls" version = "0.26.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" dependencies = [ - "rustls", + "rustls 0.23.37", "tokio", ] @@ -7107,11 +7331,11 @@ checksum = "edc5f74e248dc973e0dbb7b74c7e0d6fcc301c694ff50049504004ef4d0cdcd9" dependencies = [ "futures-util", "log", - "rustls", - "rustls-native-certs", + "rustls 0.23.37", + "rustls-native-certs 0.8.3", "rustls-pki-types", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tungstenite 0.24.0", ] @@ -7413,7 +7637,7 @@ dependencies = [ "httparse", "log", "rand 0.8.5", - "rustls", + "rustls 0.23.37", "rustls-pki-types", "sha1", "thiserror 1.0.69", @@ -7463,7 +7687,7 @@ checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ "memoffset", "tempfile", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -8310,7 +8534,7 @@ checksum = "7130243a7a5b33c54a444e54842e6a9e133de08b5ad7b5861cd8ed9a6a5bc96a" dependencies = [ "webview2-com-macros", "webview2-com-sys", - "windows", + "windows 0.61.3", "windows-core 0.61.2", "windows-implement", "windows-interface", @@ -8334,7 +8558,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "381336cfffd772377d291702245447a5251a2ffa5bad679c99e61bc48bacbf9c" dependencies = [ "thiserror 2.0.18", - "windows", + "windows 0.61.3", "windows-core 0.61.2", ] @@ -8360,7 +8584,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -8410,11 +8634,23 @@ version = "0.61.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9babd3a767a4c1aef6900409f85f5d53ce2544ccdfaa86dad48c91782c6d6893" dependencies = [ - "windows-collections", + "windows-collections 0.2.0", "windows-core 0.61.2", - "windows-future", + "windows-future 0.2.1", "windows-link 0.1.3", - "windows-numerics", + "windows-numerics 0.2.0", +] + +[[package]] +name = "windows" +version = "0.62.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "527fadee13e0c05939a6a05d5bd6eec6cd2e3dbd648b9f8e447c6518133d8580" +dependencies = [ + "windows-collections 0.3.2", + "windows-core 0.62.2", + "windows-future 0.3.2", + "windows-numerics 0.3.1", ] [[package]] @@ -8426,6 +8662,15 @@ dependencies = [ "windows-core 0.61.2", ] +[[package]] +name = "windows-collections" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "23b2d95af1a8a14a3c7367e1ed4fc9c20e0a26e79551b1454d72583c97cc6610" +dependencies = [ + "windows-core 0.62.2", +] + [[package]] name = "windows-core" version = "0.61.2" @@ -8460,7 +8705,18 @@ checksum = "fc6a41e98427b19fe4b73c550f060b59fa592d7d686537eebf9385621bfbad8e" dependencies = [ "windows-core 0.61.2", "windows-link 0.1.3", - "windows-threading", + "windows-threading 0.1.0", +] + +[[package]] +name = "windows-future" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e1d6f90251fe18a279739e78025bd6ddc52a7e22f921070ccdc67dde84c605cb" +dependencies = [ + "windows-core 0.62.2", + "windows-link 0.2.1", + "windows-threading 0.2.1", ] [[package]] @@ -8507,6 +8763,16 @@ dependencies = [ "windows-link 0.1.3", ] +[[package]] +name = "windows-numerics" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e2e40844ac143cdb44aead537bbf727de9b044e107a0f1220392177d15b0f26" +dependencies = [ + "windows-core 0.62.2", + "windows-link 0.2.1", +] + [[package]] name = "windows-result" version = "0.3.4" @@ -8645,6 +8911,15 @@ dependencies = [ "windows-link 0.1.3", ] +[[package]] +name = "windows-threading" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3949bd5b99cafdf1c7ca86b43ca564028dfe27d66958f2470940f73d86d75b37" +dependencies = [ + "windows-link 0.2.1", +] + [[package]] name = "windows-version" version = "0.1.7" @@ -8988,7 +9263,7 @@ dependencies = [ "webkit2gtk", "webkit2gtk-sys", "webview2-com", - "windows", + "windows 0.61.3", "windows-core 0.61.2", "windows-version", "x11-dl", diff --git a/crates/openfang-runtime/src/mcp.rs b/crates/openfang-runtime/src/mcp.rs index 6792593b..b9f5f381 100644 --- a/crates/openfang-runtime/src/mcp.rs +++ b/crates/openfang-runtime/src/mcp.rs @@ -1,16 +1,20 @@ //! MCP (Model Context Protocol) client — connect to external MCP servers. //! -//! MCP uses JSON-RPC 2.0 over stdio or HTTP+SSE. This module lets OpenFang -//! agents use tools from any MCP server (100+ available: GitHub, filesystem, -//! databases, APIs, etc.). +//! Uses the official `rmcp` SDK for protocol handling. Supports: +//! - **stdio**: subprocess with JSON-RPC over stdin/stdout +//! - **sse**: deprecated HTTP+SSE transport (protocol version 2024-11-05) +//! - **http**: Streamable HTTP transport (protocol version 2025-03-26+) //! //! All MCP tools are namespaced with `mcp_{server}_{tool}` to prevent collisions. +use http::{HeaderName, HeaderValue}; use openfang_types::tool::ToolDefinition; +use rmcp::model::{CallToolRequestParams, ClientCapabilities, ClientInfo, Implementation}; +use rmcp::service::RunningService; +use rmcp::{RoleClient, ServiceExt}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; -use std::process::Stdio; -use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; +use std::sync::Arc; use tracing::{debug, info}; // --------------------------------------------------------------------------- @@ -30,6 +34,12 @@ pub struct McpServerConfig { /// Environment variables to pass through to the subprocess (sandboxed). #[serde(default)] pub env: Vec, + /// Extra HTTP headers to send with every SSE / Streamable-HTTP request. + /// Each entry is `"Header-Name: value"`. Useful for authentication + /// (`Authorization: Bearer `), API keys (`X-Api-Key: ...`), + /// or any custom headers required by a remote MCP server. + #[serde(default)] + pub headers: Vec, } fn default_timeout() -> u64 { @@ -46,8 +56,14 @@ pub enum McpTransport { #[serde(default)] args: Vec, }, - /// HTTP Server-Sent Events. + /// Deprecated HTTP+SSE transport (protocol version 2024-11-05). + /// Uses POST for sending and SSE for receiving. Sse { url: String }, + /// Streamable HTTP transport (MCP 2025-03-26+). + /// Single endpoint, client MUST send Accept: application/json, text/event-stream. + /// Server responds with either JSON or SSE stream. + /// Supports Mcp-Session-Id for session management. + Http { url: String }, } // --------------------------------------------------------------------------- @@ -64,59 +80,9 @@ pub struct McpConnection { /// Needed because `normalize_name` replaces hyphens with underscores, /// but the server expects the original name (e.g. "list-connections"). original_names: HashMap, - /// Transport handle for sending requests. - transport: McpTransportHandle, - /// Next JSON-RPC request ID. - next_id: u64, -} - -/// Transport handle — abstraction over stdio subprocess or HTTP. -enum McpTransportHandle { - Stdio { - child: Box, - stdin: tokio::process::ChildStdin, - stdout: BufReader, - }, - Sse { - client: reqwest::Client, - url: String, - }, -} - -/// JSON-RPC 2.0 request. -#[derive(Serialize)] -struct JsonRpcRequest { - jsonrpc: &'static str, - id: u64, - method: String, - #[serde(skip_serializing_if = "Option::is_none")] - params: Option, -} - -/// JSON-RPC 2.0 response. -#[derive(Deserialize)] -struct JsonRpcResponse { - #[allow(dead_code)] - jsonrpc: String, - #[allow(dead_code)] - id: Option, - result: Option, - error: Option, -} - -/// JSON-RPC 2.0 error object. -#[derive(Debug, Deserialize)] -pub struct JsonRpcError { - pub code: i64, - pub message: String, - #[allow(dead_code)] - pub data: Option, -} - -impl std::fmt::Display for JsonRpcError { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "JSON-RPC error {}: {}", self.code, self.message) - } + /// The rmcp client handle — type-erased because the concrete type + /// depends on which transport was used (stdio vs HTTP). + client: RunningService, } // --------------------------------------------------------------------------- @@ -126,13 +92,17 @@ impl std::fmt::Display for JsonRpcError { impl McpConnection { /// Connect to an MCP server, perform handshake, and discover tools. pub async fn connect(config: McpServerConfig) -> Result { - let transport = match &config.transport { + let client_info = ClientInfo::new( + ClientCapabilities::default(), + Implementation::new("openfang", env!("CARGO_PKG_VERSION")), + ); + + let client = match &config.transport { McpTransport::Stdio { command, args } => { - Self::connect_stdio(command, args, &config.env).await? + Self::connect_stdio(command, args, &config.env, client_info).await? } - McpTransport::Sse { url } => { - // SSRF check: reject private/localhost URLs unless explicitly configured - Self::connect_sse(url).await? + McpTransport::Sse { url } | McpTransport::Http { url } => { + Self::connect_http(url, &config.headers, client_info).await? } }; @@ -140,13 +110,9 @@ impl McpConnection { config, tools: Vec::new(), original_names: HashMap::new(), - transport, - next_id: 1, + client, }; - // Initialize handshake - conn.initialize().await?; - // Discover tools conn.discover_tools().await?; @@ -159,76 +125,34 @@ impl McpConnection { Ok(conn) } - /// Send the MCP `initialize` handshake. - async fn initialize(&mut self) -> Result<(), String> { - let params = serde_json::json!({ - "protocolVersion": "2024-11-05", - "capabilities": {}, - "clientInfo": { - "name": "openfang", - "version": env!("CARGO_PKG_VERSION") - } - }); - - let response = self.send_request("initialize", Some(params)).await?; - - if let Some(result) = response { - debug!( - server = %self.config.name, - server_info = %result, - "MCP initialize response" - ); - } - - // Send initialized notification (no response expected) - self.send_notification("notifications/initialized", None) - .await?; - - Ok(()) - } - /// Discover available tools via `tools/list`. async fn discover_tools(&mut self) -> Result<(), String> { - let response = self.send_request("tools/list", None).await?; + let tools = self + .client + .list_all_tools() + .await + .map_err(|e| format!("Failed to list MCP tools: {e}"))?; - if let Some(result) = response { - if let Some(tools_array) = result.get("tools").and_then(|t| t.as_array()) { - let server_name = &self.config.name; - for tool in tools_array { - let raw_name = tool["name"].as_str().unwrap_or("unnamed"); - let description = tool["description"].as_str().unwrap_or(""); - let input_schema = tool - .get("inputSchema") - .cloned() - .and_then(|v| { - // Ensure input_schema is a JSON object. MCP servers may - // return it as a string, null, or omit it entirely. - match &v { - serde_json::Value::Object(_) => Some(v), - serde_json::Value::String(s) => { - serde_json::from_str::(s) - .ok() - .filter(|p| p.is_object()) - } - _ => None, - } - }) - .unwrap_or(serde_json::json!({"type": "object"})); + let server_name = &self.config.name; + for tool in &tools { + let raw_name = &tool.name; + let description = tool.description.as_deref().unwrap_or(""); - // Namespace: mcp_{server}_{tool} - let namespaced = format_mcp_tool_name(server_name, raw_name); + let input_schema = serde_json::to_value(&tool.input_schema) + .unwrap_or(serde_json::json!({"type": "object"})); - // Store original name so we can send it back to the server - self.original_names - .insert(namespaced.clone(), raw_name.to_string()); + // Namespace: mcp_{server}_{tool} + let namespaced = format_mcp_tool_name(server_name, raw_name); - self.tools.push(ToolDefinition { - name: namespaced, - description: format!("[MCP:{server_name}] {description}"), - input_schema, - }); - } - } + // Store original name so we can send it back to the server + self.original_names + .insert(namespaced.clone(), raw_name.to_string()); + + self.tools.push(ToolDefinition { + name: namespaced, + description: format!("[MCP:{server_name}] {description}"), + input_schema, + }); } Ok(()) @@ -243,40 +167,38 @@ impl McpConnection { arguments: &serde_json::Value, ) -> Result { // Look up the original tool name from the server (preserves hyphens etc.) - let raw_name = self + let raw_name: String = self .original_names .get(name) - .map(|s| s.as_str()) - .or_else(|| strip_mcp_prefix(&self.config.name, name)) - .unwrap_or(name); + .cloned() + .or_else(|| strip_mcp_prefix(&self.config.name, name).map(|s| s.to_string())) + .unwrap_or_else(|| name.to_string()); - let params = serde_json::json!({ - "name": raw_name, - "arguments": arguments, - }); + let args = arguments.as_object().cloned().unwrap_or_default(); - let response = self.send_request("tools/call", Some(params)).await?; + debug!(tool = %raw_name, server = %self.config.name, "MCP tool call"); - match response { - Some(result) => { - // Extract text content from the response - if let Some(content) = result.get("content").and_then(|c| c.as_array()) { - let texts: Vec<&str> = content - .iter() - .filter_map(|item| { - if item["type"].as_str() == Some("text") { - item["text"].as_str() - } else { - None - } - }) - .collect(); - Ok(texts.join("\n")) - } else { - Ok(result.to_string()) - } - } - None => Err("No result from MCP tools/call".to_string()), + let params = CallToolRequestParams::new(raw_name).with_arguments(args); + + let result = self + .client + .call_tool(params) + .await + .map_err(|e| format!("MCP tool call failed: {e}"))?; + + // Extract text content from the response. + // `Content` is `Annotated` which Derefs to `RawContent`. + let texts: Vec<&str> = result + .content + .iter() + .filter_map(|item| item.as_text().map(|tc| tc.text.as_str())) + .collect(); + + if texts.is_empty() { + // Fallback: serialize the entire result + Ok(serde_json::to_string(&result).unwrap_or_default()) + } else { + Ok(texts.join("\n")) } } @@ -290,299 +212,124 @@ impl McpConnection { &self.config.name } - // --- Transport helpers --- - - async fn send_request( - &mut self, - method: &str, - params: Option, - ) -> Result, String> { - let id = self.next_id; - self.next_id += 1; - - let request = JsonRpcRequest { - jsonrpc: "2.0", - id, - method: method.to_string(), - params, - }; - - let request_json = serde_json::to_string(&request) - .map_err(|e| format!("Failed to serialize request: {e}"))?; - - debug!(method, id, "MCP request"); - - match &mut self.transport { - McpTransportHandle::Stdio { stdin, stdout, .. } => { - // Write request + newline - stdin - .write_all(request_json.as_bytes()) - .await - .map_err(|e| format!("Failed to write to MCP stdin: {e}"))?; - stdin - .write_all(b"\n") - .await - .map_err(|e| format!("Failed to write newline: {e}"))?; - stdin - .flush() - .await - .map_err(|e| format!("Failed to flush stdin: {e}"))?; - - // Read response lines until we find one matching our request ID. - // MCP servers may send notifications or log lines before the - // actual response. - let timeout_dur = tokio::time::Duration::from_secs(self.config.timeout_secs); - let deadline = tokio::time::Instant::now() + timeout_dur; - - loop { - let mut line = String::new(); - let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); - if remaining.is_zero() { - return Err("MCP request timed out".to_string()); - } - match tokio::time::timeout(remaining, stdout.read_line(&mut line)).await { - Ok(Ok(0)) => return Err("MCP server closed connection".to_string()), - Ok(Ok(_)) => {} - Ok(Err(e)) => return Err(format!("Failed to read MCP response: {e}")), - Err(_) => return Err("MCP request timed out".to_string()), - } - - let trimmed = line.trim(); - if trimmed.is_empty() { - continue; - } - - // Try to parse as JSON-RPC response - let parsed: Result = serde_json::from_str(trimmed); - match parsed { - Ok(response) if response.id == Some(id) => { - if let Some(err) = response.error { - return Err(format!("{err}")); - } - return Ok(response.result); - } - Ok(_other) => { - // Response for a different ID or notification — skip - debug!("MCP: skipping non-matching response line"); - continue; - } - Err(_) => { - // Not valid JSON-RPC — skip (could be a log line) - debug!("MCP: skipping non-JSON line from server"); - continue; - } - } - } - } - McpTransportHandle::Sse { client, url } => { - let response = client - .post(url.as_str()) - .json(&request) - .timeout(std::time::Duration::from_secs(self.config.timeout_secs)) - .send() - .await - .map_err(|e| format!("MCP SSE request failed: {e}"))?; - - if !response.status().is_success() { - return Err(format!("MCP SSE returned {}", response.status())); - } - - let body = response - .text() - .await - .map_err(|e| format!("Failed to read SSE response: {e}"))?; - - // Handle Streamable HTTP MCP responses that use SSE framing - // (e.g. "event: message\ndata: {...}\n\n"). Extract the JSON - // from the last `data:` line if the body looks like SSE. - let json_body = if body.trim_start().starts_with("event:") - || body.trim_start().starts_with("data:") - { - body.lines() - .filter_map(|line| { - line.strip_prefix("data: ") - .or_else(|| line.strip_prefix("data:")) - .filter(|s| !s.is_empty()) - }) - .next_back() - .unwrap_or(&body) - .to_string() - } else { - body - }; - - let rpc_response: JsonRpcResponse = serde_json::from_str(&json_body) - .map_err(|e| format!("Invalid MCP SSE JSON-RPC response: {e}"))?; - - if let Some(err) = rpc_response.error { - return Err(format!("{err}")); - } - - Ok(rpc_response.result) - } - } - } - - async fn send_notification( - &mut self, - method: &str, - params: Option, - ) -> Result<(), String> { - let notification = serde_json::json!({ - "jsonrpc": "2.0", - "method": method, - "params": params.unwrap_or(serde_json::json!({})), - }); - - let json = serde_json::to_string(¬ification) - .map_err(|e| format!("Failed to serialize notification: {e}"))?; - - match &mut self.transport { - McpTransportHandle::Stdio { stdin, .. } => { - stdin - .write_all(json.as_bytes()) - .await - .map_err(|e| format!("Write notification: {e}"))?; - stdin - .write_all(b"\n") - .await - .map_err(|e| format!("Write newline: {e}"))?; - stdin.flush().await.map_err(|e| format!("Flush: {e}"))?; - } - McpTransportHandle::Sse { client, url } => { - let _ = client.post(url.as_str()).json(¬ification).send().await; - } - } - - Ok(()) - } + // -- Transport constructors ----------------------------------------------- + /// Connect using stdio transport (subprocess). async fn connect_stdio( command: &str, args: &[String], env_whitelist: &[String], - ) -> Result { + client_info: ClientInfo, + ) -> Result, String> { + use rmcp::transport::{ConfigureCommandExt, TokioChildProcess}; + use tokio::process::Command; + // Validate command path (no path traversal) if command.contains("..") { return Err("MCP command path contains '..': rejected".to_string()); } - // On Windows, npm/npx install as .cmd batch wrappers. Detect and adapt. - let resolved_command: String = if cfg!(windows) { - // If the user already specified .cmd/.bat, use as-is - if command.ends_with(".cmd") || command.ends_with(".bat") { - command.to_string() - } else { - // Check if the .cmd variant exists on PATH - let cmd_variant = format!("{command}.cmd"); - let has_cmd = std::env::var("PATH") - .unwrap_or_default() - .split(';') - .any(|dir| std::path::Path::new(dir).join(&cmd_variant).exists()); - if has_cmd { - cmd_variant - } else { - command.to_string() + let cmd_str = command.to_string(); + let args_vec: Vec = args.to_vec(); + let env_list: Vec = env_whitelist.to_vec(); + + let transport = TokioChildProcess::new(Command::new(&cmd_str).configure(move |cmd| { + for arg in &args_vec { + cmd.arg(arg); + } + // Sandbox: clear environment, only pass whitelisted vars + cmd.env_clear(); + for var_name in &env_list { + if let Ok(val) = std::env::var(var_name) { + cmd.env(var_name, val); } } - } else { - command.to_string() + // Always pass PATH for binary resolution + if let Ok(path) = std::env::var("PATH") { + cmd.env("PATH", path); + } + // On Windows, npm/node need extra vars + if cfg!(windows) { + for var in &[ + "APPDATA", + "LOCALAPPDATA", + "USERPROFILE", + "SystemRoot", + "TEMP", + "TMP", + "HOME", + "HOMEDRIVE", + "HOMEPATH", + ] { + if let Ok(val) = std::env::var(var) { + cmd.env(var, val); + } + } + } + })) + .map_err(|e| format!("Failed to spawn MCP server '{cmd_str}': {e}"))?; + + let client = client_info + .serve(transport) + .await + .map_err(|e| format!("MCP stdio handshake failed: {e}"))?; + + Ok(client) + } + + /// Connect using Streamable HTTP transport (or SSE fallback via the same endpoint). + /// + /// The `rmcp` SDK's `StreamableHttpClientTransport` handles the full + /// Streamable HTTP protocol: Accept headers, Mcp-Session-Id tracking, + /// SSE stream parsing, and content-type negotiation. + async fn connect_http( + url: &str, + headers: &[String], + client_info: ClientInfo, + ) -> Result, String> { + use rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig; + use rmcp::transport::StreamableHttpClientTransport; + + Self::check_ssrf(url)?; + + // Parse custom headers (e.g., "Authorization: Bearer "). + let mut custom_headers: HashMap = HashMap::new(); + for header_str in headers { + if let Some((name, value)) = header_str.split_once(':') { + let name = name.trim(); + let value = value.trim(); + if let (Ok(hn), Ok(hv)) = ( + HeaderName::from_bytes(name.as_bytes()), + HeaderValue::from_str(value), + ) { + custom_headers.insert(hn, hv); + } + } + } + + let config = StreamableHttpClientTransportConfig { + uri: Arc::from(url), + custom_headers, + ..Default::default() }; - let mut cmd = tokio::process::Command::new(&resolved_command); - cmd.args(args); - cmd.stdin(Stdio::piped()); - cmd.stdout(Stdio::piped()); - cmd.stderr(Stdio::piped()); + let transport = StreamableHttpClientTransport::from_config(config); - // Sandbox: clear environment, only pass whitelisted vars - cmd.env_clear(); - for var_name in env_whitelist { - if let Ok(val) = std::env::var(var_name) { - cmd.env(var_name, val); - } - } - // Always pass PATH for binary resolution - if let Ok(path) = std::env::var("PATH") { - cmd.env("PATH", path); - } - // On Windows, npm/node need APPDATA, USERPROFILE, LOCALAPPDATA, and SystemRoot - if cfg!(windows) { - for var in &[ - "APPDATA", - "LOCALAPPDATA", - "USERPROFILE", - "SystemRoot", - "TEMP", - "TMP", - "HOME", - "HOMEDRIVE", - "HOMEPATH", - ] { - if let Ok(val) = std::env::var(var) { - cmd.env(var, val); - } - } - } + let client = client_info + .serve(transport) + .await + .map_err(|e| format!("MCP HTTP connection failed: {e}"))?; - let mut child = cmd - .spawn() - .map_err(|e| format!("Failed to spawn MCP server '{resolved_command}': {e}"))?; - - // Log stderr in background for debugging MCP server issues - if let Some(stderr) = child.stderr.take() { - let cmd_name = resolved_command.clone(); - tokio::spawn(async move { - use tokio::io::AsyncBufReadExt; - let reader = tokio::io::BufReader::new(stderr); - let mut lines = reader.lines(); - while let Ok(Some(line)) = lines.next_line().await { - tracing::debug!(mcp_server = %cmd_name, "stderr: {line}"); - } - }); - } - - let stdin = child - .stdin - .take() - .ok_or("Failed to capture MCP server stdin")?; - let stdout = child - .stdout - .take() - .ok_or("Failed to capture MCP server stdout")?; - - Ok(McpTransportHandle::Stdio { - child: Box::new(child), - stdin, - stdout: BufReader::new(stdout), - }) + Ok(client) } - async fn connect_sse(url: &str) -> Result { - // Basic SSRF check: reject obviously private URLs + /// Basic SSRF check: reject obviously private/metadata URLs. + fn check_ssrf(url: &str) -> Result<(), String> { let lower = url.to_lowercase(); if lower.contains("169.254.169.254") || lower.contains("metadata.google") { - return Err("SSRF: MCP SSE URL targets metadata endpoint".to_string()); - } - - let client = reqwest::Client::builder() - .timeout(std::time::Duration::from_secs(30)) - .build() - .map_err(|e| format!("Failed to create HTTP client: {e}"))?; - - Ok(McpTransportHandle::Sse { - client, - url: url.to_string(), - }) - } -} - -impl Drop for McpConnection { - fn drop(&mut self) { - if let McpTransportHandle::Stdio { ref mut child, .. } = self.transport { - // Best-effort kill of the subprocess - let _ = child.start_kill(); + return Err("SSRF: MCP URL targets metadata endpoint".to_string()); } + Ok(()) } } @@ -673,16 +420,12 @@ mod tests { #[test] fn test_hyphenated_tool_name_preserved() { - // Tool names with hyphens get normalized to underscores for namespacing, - // but original_names map preserves the original for call_tool dispatch. let namespaced = format_mcp_tool_name("sqlcl", "list-connections"); assert_eq!(namespaced, "mcp_sqlcl_list_connections"); - // Simulate what discover_tools does let mut original_names = HashMap::new(); original_names.insert(namespaced.clone(), "list-connections".to_string()); - // call_tool should resolve to original hyphenated name let raw = original_names .get(&namespaced) .map(|s| s.as_str()) @@ -701,25 +444,21 @@ mod tests { #[test] fn test_extract_mcp_server_from_known_with_hyphens() { - // Server "bocha-search" normalized to "bocha_search" in tool prefix let servers = vec!["bocha-search", "github"]; let tool = "mcp_bocha_search_bocha_web_search"; assert_eq!( extract_mcp_server_from_known(tool, &servers), Some("bocha-search") ); - // Simple server name still works assert_eq!( extract_mcp_server_from_known("mcp_github_create_issue", &servers), Some("github") ); - // Non-MCP tool returns None assert_eq!(extract_mcp_server_from_known("file_read", &servers), None); } #[test] fn test_extract_mcp_server_from_known_longest_match() { - // "my-api" and "my-api-v2" — should match the longer one let servers = vec!["my-api", "my-api-v2"]; assert_eq!( extract_mcp_server_from_known("mcp_my_api_v2_get_users", &servers), @@ -731,60 +470,6 @@ mod tests { ); } - #[test] - fn test_mcp_jsonrpc_initialize() { - // Verify the initialize request structure - let request = JsonRpcRequest { - jsonrpc: "2.0", - id: 1, - method: "initialize".to_string(), - params: Some(serde_json::json!({ - "protocolVersion": "2024-11-05", - "capabilities": {}, - "clientInfo": { - "name": "openfang", - "version": "0.1.0" - } - })), - }; - let json = serde_json::to_string(&request).unwrap(); - assert!(json.contains("initialize")); - assert!(json.contains("protocolVersion")); - assert!(json.contains("openfang")); - } - - #[test] - fn test_mcp_jsonrpc_tools_list() { - // Simulate a tools/list response - let response_json = r#"{ - "jsonrpc": "2.0", - "id": 2, - "result": { - "tools": [ - { - "name": "create_issue", - "description": "Create a GitHub issue", - "inputSchema": { - "type": "object", - "properties": { - "title": {"type": "string"}, - "body": {"type": "string"} - }, - "required": ["title"] - } - } - ] - } - }"#; - - let response: JsonRpcResponse = serde_json::from_str(response_json).unwrap(); - assert!(response.error.is_none()); - let result = response.result.unwrap(); - let tools = result["tools"].as_array().unwrap(); - assert_eq!(tools.len(), 1); - assert_eq!(tools[0]["name"].as_str().unwrap(), "create_issue"); - } - #[test] fn test_mcp_transport_config_serde() { let config = McpServerConfig { @@ -798,6 +483,7 @@ mod tests { }, timeout_secs: 30, env: vec!["GITHUB_PERSONAL_ACCESS_TOKEN".to_string()], + headers: vec![], }; let json = serde_json::to_string(&config).unwrap(); @@ -822,6 +508,7 @@ mod tests { }, timeout_secs: 60, env: vec![], + headers: vec![], }; let json = serde_json::to_string(&sse_config).unwrap(); let back: McpServerConfig = serde_json::from_str(&json).unwrap(); @@ -829,5 +516,24 @@ mod tests { McpTransport::Sse { url } => assert_eq!(url, "https://example.com/mcp"), _ => panic!("Expected SSE transport"), } + + // HTTP (Streamable HTTP) variant + let http_config = McpServerConfig { + name: "atlassian".to_string(), + transport: McpTransport::Http { + url: "https://mcp.atlassian.com/v1/mcp".to_string(), + }, + timeout_secs: 120, + env: vec![], + headers: vec!["Authorization: Bearer test-token-456".to_string()], + }; + let json = serde_json::to_string(&http_config).unwrap(); + let back: McpServerConfig = serde_json::from_str(&json).unwrap(); + match back.transport { + McpTransport::Http { url } => { + assert_eq!(url, "https://mcp.atlassian.com/v1/mcp") + } + _ => panic!("Expected Http transport"), + } } } From 6ed6d3ac3b2989c1a7164e939542d2cd157e0c73 Mon Sep 17 00:00:00 2001 From: Philippe Branchu Date: Sat, 28 Mar 2026 04:15:46 +0000 Subject: [PATCH 3/5] Fix cargo fmt formatting Co-Authored-By: Claude Opus 4.6 (1M context) --- Cargo.lock | 4 ++++ crates/openfang-api/src/channel_bridge.rs | 2 +- crates/openfang-api/src/routes.rs | 3 ++- crates/openfang-channels/src/lib.rs | 2 +- crates/openfang-channels/src/line.rs | 6 ++--- crates/openfang-channels/src/mqtt.rs | 8 +++++-- crates/openfang-kernel/src/kernel.rs | 24 ++++++++----------- crates/openfang-runtime/src/compactor.rs | 10 ++++++-- .../openfang-runtime/src/context_overflow.rs | 24 ++++++++++++++----- 9 files changed, 52 insertions(+), 31 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index df329192..a5b30bfc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7687,7 +7687,11 @@ checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ "memoffset", "tempfile", +<<<<<<< Updated upstream "windows-sys 0.61.2", +======= + "windows-sys 0.60.2", +>>>>>>> Stashed changes ] [[package]] diff --git a/crates/openfang-api/src/channel_bridge.rs b/crates/openfang-api/src/channel_bridge.rs index 39fa42e2..f8336ba5 100644 --- a/crates/openfang-api/src/channel_bridge.rs +++ b/crates/openfang-api/src/channel_bridge.rs @@ -49,8 +49,8 @@ use openfang_channels::discourse::DiscourseAdapter; use openfang_channels::gitter::GitterAdapter; use openfang_channels::gotify::GotifyAdapter; use openfang_channels::linkedin::LinkedInAdapter; -use openfang_channels::mumble::MumbleAdapter; use openfang_channels::mqtt::MqttAdapter; +use openfang_channels::mumble::MumbleAdapter; use openfang_channels::ntfy::NtfyAdapter; use openfang_channels::webhook::WebhookAdapter; use openfang_channels::wecom::WeComAdapter; diff --git a/crates/openfang-api/src/routes.rs b/crates/openfang-api/src/routes.rs index e547a896..b3ce41dc 100644 --- a/crates/openfang-api/src/routes.rs +++ b/crates/openfang-api/src/routes.rs @@ -579,7 +579,8 @@ pub async fn get_agent_session( msg.get_mut("tools").and_then(|v| v.as_array_mut()) { if let Some(tool_obj) = tools_arr.get_mut(tool_idx) { - tool_obj["result"] = serde_json::Value::String(result.clone()); + tool_obj["result"] = + serde_json::Value::String(result.clone()); tool_obj["is_error"] = serde_json::Value::Bool(*is_error); } diff --git a/crates/openfang-channels/src/lib.rs b/crates/openfang-channels/src/lib.rs index 949d3935..7b122d2a 100644 --- a/crates/openfang-channels/src/lib.rs +++ b/crates/openfang-channels/src/lib.rs @@ -48,8 +48,8 @@ pub mod discourse; pub mod gitter; pub mod gotify; pub mod linkedin; -pub mod mumble; pub mod mqtt; +pub mod mumble; pub mod ntfy; pub mod webhook; pub mod wecom; diff --git a/crates/openfang-channels/src/line.rs b/crates/openfang-channels/src/line.rs index b20294af..479cd427 100644 --- a/crates/openfang-channels/src/line.rs +++ b/crates/openfang-channels/src/line.rs @@ -381,8 +381,7 @@ impl ChannelAdapter for LineAdapter { axum::routing::post({ let secret = Arc::clone(&channel_secret); let tx = Arc::clone(&tx); - move |headers: axum::http::HeaderMap, - body: axum::body::Bytes| { + move |headers: axum::http::HeaderMap, body: axum::body::Bytes| { let secret = Arc::clone(&secret); let tx = Arc::clone(&tx); async move { @@ -404,8 +403,7 @@ impl ChannelAdapter for LineAdapter { shutdown_rx: watch::channel(false).1, }; - if !signature.is_empty() - && !adapter.verify_signature(&body, signature) + if !signature.is_empty() && !adapter.verify_signature(&body, signature) { warn!("LINE: invalid webhook signature"); return axum::http::StatusCode::UNAUTHORIZED; diff --git a/crates/openfang-channels/src/mqtt.rs b/crates/openfang-channels/src/mqtt.rs index 69bb6349..a3e5b154 100644 --- a/crates/openfang-channels/src/mqtt.rs +++ b/crates/openfang-channels/src/mqtt.rs @@ -152,7 +152,10 @@ impl MqttAdapter { } /// Parse host:port string. - fn parse_host_port(s: &str, default_port: u16) -> Result<(String, u16), Box> { + fn parse_host_port( + s: &str, + default_port: u16, + ) -> Result<(String, u16), Box> { let s = s.trim(); if let Some(colon_pos) = s.rfind(':') { let host = s[..colon_pos].to_string(); @@ -239,7 +242,8 @@ impl ChannelAdapter for MqttAdapter { async fn start( &self, - ) -> Result + Send>>, Box> { + ) -> Result + Send>>, Box> + { let options = self.build_mqtt_options()?; let (client, mut eventloop) = AsyncClient::new(options, 10); diff --git a/crates/openfang-kernel/src/kernel.rs b/crates/openfang-kernel/src/kernel.rs index f449adda..52bf6c6c 100644 --- a/crates/openfang-kernel/src/kernel.rs +++ b/crates/openfang-kernel/src/kernel.rs @@ -2889,20 +2889,16 @@ impl OpenFangKernel { model: &str, explicit_provider: Option<&str>, ) -> KernelResult<()> { - let catalog_entry = self - .model_catalog - .read() - .ok() - .and_then(|catalog| { - // When the caller specifies a provider, use provider-aware lookup - // so we resolve the model on the correct provider — not a builtin - // from a different provider that happens to share the same name (#833). - if let Some(ep) = explicit_provider { - catalog.find_model_for_provider(model, ep).cloned() - } else { - catalog.find_model(model).cloned() - } - }); + let catalog_entry = self.model_catalog.read().ok().and_then(|catalog| { + // When the caller specifies a provider, use provider-aware lookup + // so we resolve the model on the correct provider — not a builtin + // from a different provider that happens to share the same name (#833). + if let Some(ep) = explicit_provider { + catalog.find_model_for_provider(model, ep).cloned() + } else { + catalog.find_model(model).cloned() + } + }); let provider = if let Some(ep) = explicit_provider { // User explicitly set the provider — use it as-is Some(ep.to_string()) diff --git a/crates/openfang-runtime/src/compactor.rs b/crates/openfang-runtime/src/compactor.rs index 05f75f95..3186e4f4 100644 --- a/crates/openfang-runtime/src/compactor.rs +++ b/crates/openfang-runtime/src/compactor.rs @@ -1478,7 +1478,10 @@ mod tests { Message::assistant("Done reading."), ]; let adjusted = adjust_split_for_tool_pairs(&messages, 2); - assert_eq!(adjusted, 1, "Should pull back split to keep ToolUse + ToolResult together"); + assert_eq!( + adjusted, 1, + "Should pull back split to keep ToolUse + ToolResult together" + ); } #[test] @@ -1489,7 +1492,10 @@ mod tests { Message::user("c"), ]; let adjusted = adjust_split_for_tool_pairs(&messages, 1); - assert_eq!(adjusted, 1, "Should not change split for plain text messages"); + assert_eq!( + adjusted, 1, + "Should not change split for plain text messages" + ); } #[test] diff --git a/crates/openfang-runtime/src/context_overflow.rs b/crates/openfang-runtime/src/context_overflow.rs index 397bf9d4..22d14642 100644 --- a/crates/openfang-runtime/src/context_overflow.rs +++ b/crates/openfang-runtime/src/context_overflow.rs @@ -32,10 +32,14 @@ fn safe_drain_boundary(messages: &[Message], mut boundary: usize) -> usize { // is in the last drained message (boundary - 1). Pull boundary back by 1. if messages[boundary].role == Role::User { if let MessageContent::Blocks(blocks) = &messages[boundary].content { - let has_tool_result = blocks.iter().any(|b| matches!(b, ContentBlock::ToolResult { .. })); + let has_tool_result = blocks + .iter() + .any(|b| matches!(b, ContentBlock::ToolResult { .. })); if has_tool_result && boundary > 0 && messages[boundary - 1].role == Role::Assistant { if let MessageContent::Blocks(asst_blocks) = &messages[boundary - 1].content { - let has_tool_use = asst_blocks.iter().any(|b| matches!(b, ContentBlock::ToolUse { .. })); + let has_tool_use = asst_blocks + .iter() + .any(|b| matches!(b, ContentBlock::ToolUse { .. })); if has_tool_use { boundary -= 1; debug!( @@ -135,7 +139,8 @@ pub fn recover_from_overflow( debug!( estimated_tokens = estimated, removing = remove, - "Stage 1: moderate trim to last {} messages", messages.len() - remove + "Stage 1: moderate trim to last {} messages", + messages.len() - remove ); messages.drain(..remove); // Re-check after trim @@ -156,7 +161,8 @@ pub fn recover_from_overflow( warn!( estimated_tokens = estimate_tokens(messages, system_prompt, tools), removing = remove, - "Stage 2: aggressive overflow compaction to last {} messages", messages.len() - remove + "Stage 2: aggressive overflow compaction to last {} messages", + messages.len() - remove ); let summary = Message::user(format!( "[System: {} earlier messages were removed due to context overflow. \ @@ -373,7 +379,10 @@ mod tests { ]; // Boundary 2 would cut between the assistant(ToolUse) at [1] and user(ToolResult) at [2]. let adjusted = safe_drain_boundary(&msgs, 2); - assert_eq!(adjusted, 1, "Should pull boundary back to keep the ToolUse/ToolResult pair together"); + assert_eq!( + adjusted, 1, + "Should pull boundary back to keep the ToolUse/ToolResult pair together" + ); } #[test] @@ -385,7 +394,10 @@ mod tests { Message::assistant("d"), ]; let adjusted = safe_drain_boundary(&msgs, 2); - assert_eq!(adjusted, 2, "Should not change boundary for plain text messages"); + assert_eq!( + adjusted, 2, + "Should not change boundary for plain text messages" + ); } #[test] From 613a7d4a3b33966366feafa68ae0d240bb5aed47 Mon Sep 17 00:00:00 2001 From: Philippe Branchu Date: Sat, 28 Mar 2026 04:19:43 +0000 Subject: [PATCH 4/5] Fix corrupt Cargo.lock (resolve merge conflict markers) Co-Authored-By: Claude Opus 4.6 (1M context) --- Cargo.lock | 4 ---- 1 file changed, 4 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index a5b30bfc..df329192 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7687,11 +7687,7 @@ checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ "memoffset", "tempfile", -<<<<<<< Updated upstream "windows-sys 0.61.2", -======= - "windows-sys 0.60.2", ->>>>>>> Stashed changes ] [[package]] From a3cefa424c65faa40ca8ba7cf9233a9002fa5c1e Mon Sep 17 00:00:00 2001 From: Philippe Branchu Date: Sat, 28 Mar 2026 04:22:10 +0000 Subject: [PATCH 5/5] Fix clippy: remove needless borrow in line.rs Co-Authored-By: Claude Opus 4.6 (1M context) --- crates/openfang-channels/src/line.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/openfang-channels/src/line.rs b/crates/openfang-channels/src/line.rs index 479cd427..fa3ecd80 100644 --- a/crates/openfang-channels/src/line.rs +++ b/crates/openfang-channels/src/line.rs @@ -108,7 +108,7 @@ impl LineAdapter { diff |= a ^ b; } if diff != 0 { - let computed = base64::engine::general_purpose::STANDARD.encode(&result); + let computed = base64::engine::general_purpose::STANDARD.encode(result); // Log first/last 4 chars of each signature for debugging without leaking full HMAC let comp_redacted = format!( "{}...{}",