MCPSprut
一种Go服务,通过多个MCP(模型上下文协议)服务器维护持久的SSE流,监听工具更改通知,并通过批写入将工具存储在本地数据库中。
建筑
flowchart TD
main["cmd/sprut/main.go"] --> hub["Hub (orchestrator)"]
hub --> connector["Connector pool"]
connector --> c1["Conn: server 0"]
connector --> c2["Conn: server 1"]
connector --> cN["Conn: server N"]
c1 --> mcpClient["MCP Client"]
hub --> batcher["Batcher"]
batcher --> storage["Storage (interface)"]
storage --> bolt["BoltDB impl"]
storage --> future["Future: PG, Redis..."]运作原理
- 启动时,从存储中加载MCP服务器列表
- 对于每台服务器,执行MCP握手(
initialize+tools/list)并持续使用工具 - 向每个服务器打开GET SSE流,监听
notifications/tools/list_changed - 收到通知后(根据MCP规范,有效载荷为空),重新获取
tools/list并更新存储 - 批处理写入存储(可配置:
1为了实时,100用于高通量) - 在运行时检测添加到存储中的新服务器——无需重新启动
数据流
sequenceDiagram
participant Hub
participant Connector
participant MCPServer
participant Batcher
participant Storage
Hub->>Storage: LoadServers()
Hub->>Connector: Start(serverConfig)
Connector->>MCPServer: POST initialize
MCPServer-->>Connector: capabilities
Connector->>MCPServer: POST tools/list
MCPServer-->>Connector: tools[]
Connector->>Batcher: Submit(serverID, tools)
Connector->>MCPServer: GET /mcp (SSE stream)
loop On notification
MCPServer-->>Connector: notifications/tools/list_changed
Connector->>MCPServer: POST tools/list
MCPServer-->>Connector: tools[]
Connector->>Batcher: Submit(serverID, tools)
end
Batcher->>Storage: SaveToolsBatch(updates)项目结构
| 路径 | 描述 |
|---|---|
cmd/sprut/main.go | 入口点、环境配置、优雅关机 |
internal/hub/hub.go | Orchestrator:加载服务器、启动连接器、处理新服务器 |
internal/connector/connector.go | 单MCP服务器连接生命周期 |
internal/mcpclient/client.go | HTTP客户端:POST初始化,POST工具/列表 |
internal/mcpclient/stream.go | SSE流:GET/mcp,事件解析 |
internal/batcher/batcher.go | 批处理工具更新,按大小或计时器刷新 |
internal/storage/interface.go | 存储接口(服务器+工具) |
internal/storage/bolt.go | BoltDB实现 |
internal/jsonrpc/schema.go | JSON-RPC 2.0消息结构 |
配置
所有设置都是从环境变量中读取的。
| 变量 | 描述 | 默认值 |
|---|---|---|
SPRUT_DB_PATH | BoltDB文件的路径 | sprut.db |
SPRUT_BATCH_SIZE | 工具写入的批量大小 | 1 |
SPRUT_BUFFER_SIZE | 连接器和批处理器之间的通道缓冲区 | 256 |
SPRUT_FLUSH_INTERVAL | 冲洗之间的最大间隔 | 5s |
SPRUT_CONNECT_TIMEOUT | MCP握手超时 | 30s |
SPRUT_RETRY_INTERVAL | 故障时重新连接间隔 | 10s |
优雅地关闭
- 信号/信号项→ 取消上下文
- 所有连接器goroutines关闭SSE流并退出
- 批处理程序刷新剩余缓冲区
- BoltDB关闭
- 进程退出
优雅降级
每个连接器独立运行。如果一台MCP服务器发生故障,其余服务器将继续工作。在SSE流断开连接或握手失败时,连接器会等待 SPRUT_RETRY_INTERVAL 并重试。
跑步
# Start the MCP simulator (separate repo)
cd ~/golang/mcp-simulator
go run ./cmd/simulator/ --servers 10 --port 9090
# Start MCPSprut
cd ~/golang/mcp-sprut
SPRUT_DB_PATH=sprut.db go run ./cmd/sprut/调用树
main.go
├── config.Load() // env → Config
├── storage.NewBoltStorage(dbPath) // open sprut.db
├── mcpclient.NewClient(timeout) // HTTP client for MCP
├── batcher.NewBatcher(store, size, interval)
│ └── batcher.Start(ctx) // goroutine: buffer → flush to storage
├── hub.NewHub(store, client, batcher, retry)
│ └── hub.Start(ctx)
│ ├── store.LoadServers() // load servers from BoltDB
│ ├── for each server → startConnector()
│ │ └── go connector.Run(ctx) // goroutine per server
│ │ └── connect() // loop with retry
│ │ ├── client.Initialize()
│ │ ├── client.SendInitialized()
│ │ ├── client.ListTools() → batcher.Submit()
│ │ └── client.SubscribeNotifications() → SSE loop
│ │ └── on notification → client.ListTools() → batcher.Submit()
│ └── store.OnNewServer(callback) // new servers at runtime
├── <-sigCh // wait for SIGINT/SIGTERM
├── cancel() // stop all goroutines
├── batcher.Wait() // wait for remaining flush
└── store.Close() // close BoltDB设计模式
| 模式 | 地点 |
|---|---|
| 存储库 | Storage 接口抽象底层存储 |
| 观察员 | OnNewServer 事件驱动服务器发现的回调 |
| 扇出 | Hub推出N个连接器goroutines |
| 批处理写入器 | 批处理器在提交到存储之前累积更新 |
| 良好的降级 | 单个连接器故障不会导致系统停机 |
