Add env-var support and native HTTP transport for External MCP Servers
External MCP Servers previously only spoke stdio (spawn a local command + args). Adds: - env vars for stdio servers (merged into the subprocess environment, not embedded in the args string), with a masked key-value editor - a native Streamable HTTP transport (URL + Bearer token + custom headers), so HTTP-based MCP servers like Obsidian's Local REST API plugin connect directly without needing npx/Node.js as a bridge Introduces an MCPTransport abstraction (stdio/HTTP) so ExternalMCPClient stays transport-agnostic — mirrors how Provider.swift already abstracts AI backends in this codebase. Also fixes a real crash found via live testing against Obsidian: convertInputSchema force-unwrapped a tool parameter's `type`, which isn't required by JSON Schema — Obsidian's plugin was the first real server to send a parameter without one. Live-verified end to end (vault search/read/write/edit) before this commit, per the project's standing rule to hold external-service-dependent changes until they're actually confirmed working, not just compiling and passing tests.
This commit is contained in:
@@ -5,22 +5,18 @@ import Foundation
|
||||
|
||||
// MARK: - ExternalMCPClient
|
||||
|
||||
/// Manages one MCP stdio server process. All state is MainActor-isolated
|
||||
/// (consistent with SWIFT_DEFAULT_ACTOR_ISOLATION = MainActor project setting).
|
||||
/// Background I/O runs in Task.detached; state mutations hop back to MainActor.
|
||||
/// Owns one MCP server connection's lifecycle and JSON-RPC message framing. Delivery (stdio
|
||||
/// subprocess vs Streamable HTTP) is delegated to a `MCPTransport` — this class only builds
|
||||
/// envelopes, decodes typed results, and tracks connection state.
|
||||
/// All state is MainActor-isolated (consistent with SWIFT_DEFAULT_ACTOR_ISOLATION = MainActor
|
||||
/// project setting).
|
||||
@MainActor
|
||||
final class ExternalMCPClient {
|
||||
let server: ExternalMCPServer
|
||||
weak var stateDelegate: (any ExternalMCPStateDelegate)?
|
||||
|
||||
private var process: Process?
|
||||
private var stdinHandle: FileHandle?
|
||||
private var readTask: Task<Void, Never>?
|
||||
private var stderrTask: Task<Void, Never>?
|
||||
|
||||
private let transport: any MCPTransport
|
||||
private var nextRequestId: Int = 1
|
||||
private var pendingCalls: [Int: CheckedContinuation<Data, Error>] = [:]
|
||||
private var lineBuffer = Data()
|
||||
|
||||
private(set) var state: MCPClientState = .idle
|
||||
private(set) var discoveredTools: [MCPToolDefinition] = []
|
||||
@@ -28,6 +24,31 @@ final class ExternalMCPClient {
|
||||
init(server: ExternalMCPServer, stateDelegate: (any ExternalMCPStateDelegate)?) {
|
||||
self.server = server
|
||||
self.stateDelegate = stateDelegate
|
||||
let stdioTransport: StdioMCPTransport?
|
||||
switch server.transportKind {
|
||||
case .stdio:
|
||||
let t = StdioMCPTransport(server: server)
|
||||
stdioTransport = t
|
||||
self.transport = t
|
||||
case .http:
|
||||
stdioTransport = nil
|
||||
self.transport = HTTPMCPTransport(server: server)
|
||||
}
|
||||
// `self` is only safe to capture once every stored property above has a value —
|
||||
// wire the crash callback here, after `init` would otherwise be considered complete.
|
||||
stdioTransport?.onTerminated = { [weak self] in
|
||||
self?.handleTransportTerminatedUnexpectedly()
|
||||
}
|
||||
}
|
||||
|
||||
/// Called by a stdio transport whose subprocess died on its own — as opposed to a
|
||||
/// deliberate `stop()` call, or a failure already handled inline within `start()`.
|
||||
/// No HTTP equivalent: a Streamable HTTP connection has no persistent process to crash;
|
||||
/// its failures surface per-request instead (handled in `start()`/`callTool()` directly).
|
||||
private func handleTransportTerminatedUnexpectedly() {
|
||||
guard state != .stopped else { return }
|
||||
state = .crashed
|
||||
stateDelegate?.clientDidChangeState(id: server.id, state: .crashed)
|
||||
}
|
||||
|
||||
// MARK: - Lifecycle
|
||||
@@ -37,54 +58,28 @@ final class ExternalMCPClient {
|
||||
state = .connecting
|
||||
stateDelegate?.clientDidChangeState(id: server.id, state: .connecting)
|
||||
|
||||
let proc = Process()
|
||||
if server.command.hasPrefix("/") {
|
||||
proc.executableURL = URL(fileURLWithPath: server.command)
|
||||
proc.arguments = server.args
|
||||
} else {
|
||||
proc.executableURL = URL(fileURLWithPath: "/usr/bin/env")
|
||||
proc.arguments = [server.command] + server.args
|
||||
}
|
||||
proc.environment = ProcessInfo.processInfo.environment
|
||||
|
||||
let stdinPipe = Pipe()
|
||||
let stdoutPipe = Pipe()
|
||||
let stderrPipe = Pipe()
|
||||
proc.standardInput = stdinPipe
|
||||
proc.standardOutput = stdoutPipe
|
||||
proc.standardError = stderrPipe
|
||||
|
||||
proc.terminationHandler = { [weak self] _ in
|
||||
Task { @MainActor [weak self] in self?.handleProcessTerminated() }
|
||||
}
|
||||
|
||||
do {
|
||||
try proc.run()
|
||||
} catch {
|
||||
state = .error(error.localizedDescription)
|
||||
stateDelegate?.clientDidChangeState(id: server.id, state: .error(error.localizedDescription))
|
||||
throw MCPClientError.processLaunchFailed(error.localizedDescription)
|
||||
}
|
||||
try await transport.prepare()
|
||||
|
||||
process = proc
|
||||
stdinHandle = stdinPipe.fileHandleForWriting
|
||||
startReadLoop(pipe: stdoutPipe)
|
||||
startStderrLoop(pipe: stderrPipe)
|
||||
|
||||
do {
|
||||
let _: MCPInitializeResult = try await timedRequest(seconds: 15, method: "initialize", params: [
|
||||
"protocolVersion": "2024-11-05",
|
||||
"capabilities": [:] as [String: Any],
|
||||
"clientInfo": ["name": "Confab", "version": "1.0"] as [String: Any]
|
||||
])
|
||||
try sendNotification(method: "notifications/initialized")
|
||||
try await transport.sendNotification(["jsonrpc": "2.0", "method": "notifications/initialized"])
|
||||
|
||||
let toolsResult: MCPToolsListResult = try await timedRequest(seconds: 15, method: "tools/list", params: nil)
|
||||
discoveredTools = toolsResult.tools
|
||||
} catch {
|
||||
state = .error(error.localizedDescription)
|
||||
stateDelegate?.clientDidChangeState(id: server.id, state: .error(error.localizedDescription))
|
||||
proc.terminate()
|
||||
// Uniformly route every start() failure (bad config, launch failure, handshake
|
||||
// failure — for either transport) through .crashed, not .error, so
|
||||
// ExternalMCPManager's restart-with-backoff drives from exactly one place.
|
||||
// (Previously, stdio relied on the subprocess's termination handler firing
|
||||
// asynchronously to reach .crashed; that path doesn't exist for HTTP, so failures
|
||||
// there would otherwise get stuck at .error with no retry.)
|
||||
state = .crashed
|
||||
transport.stop()
|
||||
stateDelegate?.clientDidChangeState(id: server.id, state: .crashed)
|
||||
throw error
|
||||
}
|
||||
|
||||
@@ -94,14 +89,7 @@ final class ExternalMCPClient {
|
||||
|
||||
func stop() {
|
||||
state = .stopped
|
||||
readTask?.cancel()
|
||||
stderrTask?.cancel()
|
||||
process?.terminate()
|
||||
process = nil
|
||||
stdinHandle = nil
|
||||
lineBuffer = Data()
|
||||
for (_, cont) in pendingCalls { cont.resume(throwing: MCPClientError.notConnected) }
|
||||
pendingCalls.removeAll()
|
||||
transport.stop()
|
||||
}
|
||||
|
||||
// MARK: - Tool Execution
|
||||
@@ -128,125 +116,17 @@ final class ExternalMCPClient {
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - I/O Loops (detached from MainActor)
|
||||
|
||||
private func startReadLoop(pipe: Pipe) {
|
||||
readTask = Task.detached { [weak self] in
|
||||
let handle = pipe.fileHandleForReading
|
||||
while true {
|
||||
let data = handle.availableData
|
||||
if data.isEmpty { break }
|
||||
await self?.receiveData(data)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private func startStderrLoop(pipe: Pipe) {
|
||||
let name = server.name
|
||||
stderrTask = Task.detached {
|
||||
let handle = pipe.fileHandleForReading
|
||||
var buf = Data()
|
||||
while true {
|
||||
let data = handle.availableData
|
||||
if data.isEmpty { break }
|
||||
buf.append(data)
|
||||
while let idx = buf.firstIndex(of: UInt8(ascii: "\n")) {
|
||||
let line = String(data: buf[buf.startIndex..<idx], encoding: .utf8) ?? ""
|
||||
buf = Data(buf[buf.index(after: idx)...])
|
||||
if !line.trimmingCharacters(in: .whitespaces).isEmpty {
|
||||
Log.extMcp.warning("[\(name)] \(line)")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Data Processing (MainActor)
|
||||
|
||||
private func receiveData(_ data: Data) {
|
||||
lineBuffer.append(data)
|
||||
while let idx = lineBuffer.firstIndex(of: UInt8(ascii: "\n")) {
|
||||
let lineData = Data(lineBuffer[lineBuffer.startIndex..<idx])
|
||||
lineBuffer = Data(lineBuffer[lineBuffer.index(after: idx)...])
|
||||
processLine(lineData)
|
||||
}
|
||||
}
|
||||
|
||||
private func processLine(_ data: Data) {
|
||||
guard !data.isEmpty,
|
||||
let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any],
|
||||
let id = json["id"] as? Int,
|
||||
let cont = pendingCalls.removeValue(forKey: id) else { return }
|
||||
|
||||
if let err = json["error"] as? [String: Any] {
|
||||
cont.resume(throwing: MCPClientError.invalidResponse(err["message"] as? String ?? "Unknown error"))
|
||||
} else if let result = json["result"],
|
||||
let resultData = try? JSONSerialization.data(withJSONObject: result) {
|
||||
cont.resume(returning: resultData)
|
||||
} else {
|
||||
cont.resume(throwing: MCPClientError.invalidResponse("Missing result field"))
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - JSON-RPC
|
||||
|
||||
/// Send a JSON-RPC request with a per-call timeout. The timeout fires a cancellation
|
||||
/// directly into the pending-calls table rather than using a task group (which would
|
||||
/// pass the generic T through a @Sendable closure and trigger an isolated-conformance warning).
|
||||
private func timedRequest<T: Decodable>(seconds: Double, method: String, params: [String: Any]?) async throws -> T {
|
||||
let id = nextRequestId
|
||||
nextRequestId += 1
|
||||
var message: [String: Any] = ["jsonrpc": "2.0", "method": method, "id": id]
|
||||
if let params { message["params"] = params }
|
||||
try writeJSON(message)
|
||||
|
||||
// Schedule timeout: cancels the specific pending call by ID
|
||||
let timeoutId = id
|
||||
Task { [weak self, timeoutId] in
|
||||
try? await Task.sleep(nanoseconds: UInt64(seconds * 1_000_000_000))
|
||||
self?.cancelPendingCall(id: timeoutId, with: MCPClientError.timeout)
|
||||
}
|
||||
|
||||
// Await response data, then decode on MainActor
|
||||
let resultData: Data = try await withCheckedThrowingContinuation { cont in
|
||||
pendingCalls[id] = cont
|
||||
}
|
||||
let resultData = try await transport.sendRequest(message, id: id, timeoutSeconds: seconds)
|
||||
return try JSONDecoder().decode(T.self, from: resultData)
|
||||
}
|
||||
|
||||
private func cancelPendingCall(id: Int, with error: Error) {
|
||||
pendingCalls.removeValue(forKey: id)?.resume(throwing: error)
|
||||
}
|
||||
|
||||
private func sendNotification(method: String) throws {
|
||||
try writeJSON(["jsonrpc": "2.0", "method": method])
|
||||
}
|
||||
|
||||
private func writeJSON(_ message: [String: Any]) throws {
|
||||
guard let handle = stdinHandle, process?.isRunning == true else {
|
||||
throw MCPClientError.writeFailed
|
||||
}
|
||||
guard let data = try? JSONSerialization.data(withJSONObject: message),
|
||||
let line = String(data: data, encoding: .utf8) else {
|
||||
throw MCPClientError.writeFailed
|
||||
}
|
||||
do {
|
||||
try handle.write(contentsOf: Data((line + "\n").utf8))
|
||||
} catch {
|
||||
throw MCPClientError.writeFailed
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Process termination
|
||||
|
||||
private func handleProcessTerminated() {
|
||||
guard state != .stopped else { return }
|
||||
state = .crashed
|
||||
for (_, cont) in pendingCalls { cont.resume(throwing: MCPClientError.notConnected) }
|
||||
pendingCalls.removeAll()
|
||||
stateDelegate?.clientDidChangeState(id: server.id, state: .crashed)
|
||||
}
|
||||
|
||||
// MARK: - Result conversion
|
||||
|
||||
private func convertMCPResult(_ result: MCPToolCallResult) -> [String: Any] {
|
||||
|
||||
Reference in New Issue
Block a user