Distributed MCP Gateway Library
项目状态: 早期阶段,固执己见。该库对于下文所述的操作问题具有完整的功能,并通过测试和示例进行了验证,但尚未由指定的生产用户部署。将其视为主动迭代下的设计:我正在寻找愿意在实际工作负载中运行它并从那里塑造它的设计合作伙伴。尽最大努力维护;问题和拉取请求是受欢迎的,并在时间允许的情况下进行审查,但回复可能需要一段时间。如果你今天依赖这个库,那么你可以为任何时间敏感的东西分叉或携带补丁。
goakt mcp 是Go的MCP(模型上下文协议)网关库,旨在超越瘦JSON-RPC代理。它被构造为MCP工作负载的操作控制平面:工具生命周期、会话关联性、凭证代理、策略执行、断路、审计和单一路由背后的集群感知路由 Gateway API
为什么选择goakt mcp
MCP工作负载是有状态的、并发的,并且容易发生故障,这是通用HTTP代理无法很好地处理的:
- 工具是异构的。 有些是stdio上的本地子进程;其他是远程HTTP服务。它们的故障模式、延迟情况和启动成本差异很大。
- 会话携带状态。 每个租户和客户的会话所有权是正确性要求,而不是优化。
- 工具后端出现故障。 连接中断、进程崩溃和呼叫缓慢是常态。如果没有监督层,单个工具故障就会悄无声息地破坏整个请求管道。
- 多租户是运营的必要条件。 速率限制、并发上限和策略执行必须在每次调用中一致应用:而不是在应用程序代码中附加。
- 可观察性不是可选的。 生产操作需要延迟、故障、电路状态、策略决策和审计跟踪。
演员模型非常适合这个问题空间。每个工具都有一个专门的监督角色。每个会话都有自己的参与者,与其他会话隔离开来。控制平面(路由器、凭证代理、策略引擎、日志)作为受监督的参与者运行,具有确定性名称、干净的消息边界和自动重启语义。故障得到控制,恢复是本地的,系统会优雅地降级,而不是灾难性地降级。
特性
- 三重运输出口 :通过stdio(子进程)、HTTP(远程MCP服务器)或gRPC(使用原型描述符集或服务器反射)调用工具
- 四个入口运输 :通过流式HTTP、服务器发送事件、WebSocket或gRPC(支持流式进程)为MCP客户端提供服务
- 流式调用 :
InvokeStream返回aStreamingResult具有进度事件通道和最终结果通道,由gRPC服务器流式RPC支持 - OAuth 2.0/OIDC身份验证 :承载令牌验证,可插拔
TokenVerifier,IdentityMapper,以及RFC 9728受保护资源元数据发现(/.well-known/oauth-protected-resource);包括用于一元RPC和流式RPC的gRPC身份验证拦截器 - 多租户技术 :每个租户配额强制(速率限制+并发上限)和可插拔策略评估
- 凭证经纪 :解析来自任何来源(vault、env、KMS)的机密,并使用可配置的LRU缓存将其注入调用中
- 会话亲和性 :每个租户+客户端+工具的粘性会话所有权;或无状态工具的最小负载平衡
- 会话空闲钝化 :会话在可配置的空闲超时(默认5分钟)后自动钝化,无需客户端干预即可回收资源
- 断路器 :每个工具的断路器具有可配置的故障阈值、断开持续时间和半断开探测
- 透明恢复 :失败的会话会重新创建其执行器并自动重试,而无需等待钝化
- 背压 :
MaxSessionsPerTool硬上限防止过载的工具耗尽共享资源;审计日志上的有界邮箱提供自然的背压,而不会丢失事件 - 健康调查 对每个工具主管进行定期健康检查;反映在刀具状态中的电路状态转换
- 动态工具管理 :注册、更新、启用、禁用、耗尽和删除工具,而无需重新启动网关
- 架构发现 :在注册时自动从后端获取和缓存MCP工具模式
- 资源代理 :通过以下方式从后端发现并缓存MCP资源和资源模板
resources/list和resources/templates/list;代理人resources/read通过具有与工具调用相同的弹性保证的完整参与者管道 - 结构化错误代码 :14已打印
ErrorCode价值观(TOOL_NOT_FOUND,POLICY_DENIED,RATE_LIMITED,TRANSPORT_FAILURE等等)RuntimeError包装和满errors.Is/As支持 - 开放遥测 :通过OTLP导出的跟踪和指标;W3C跟踪上下文在出口上传播
- 结构化日志记录 :可插拔
Logger每个日志行上都有结构化关联字段(租户ID、工具ID、请求ID、跟踪ID)的接口(zap、zerolog、slog、logrus等) - 持久的审计追踪 :每个策略决策、调用、电路状态更改和健康转换都写入可插拔的
AuditSink - 群集模式 :具有八卦成员资格的多节点操作、分布式参与者消息传递、集群单例注册器和可插拔对等体发现
- TLS无处不在 :用于群集远程处理和工具后端HTTP连接的可选双向TLS
- 完全可插拔 :每个操作问题(身份、策略、凭据、审计、发现)都是您实现的接口
建筑
goakt mcp的结构分为三层:入口、网关核心和出口:由监督的参与者运行时连接。
graph TB
subgraph Clients["MCP Clients"]
C1["Claude Desktop"]
C2["LLM Agent"]
C3["AI Application"]
end
subgraph Ingress["Ingress Layer"]
HTTP["Streamable HTTP\n/mcp"]
SSE["Server-Sent Events\n/sse"]
WS["WebSocket\n/ws"]
GRPC["gRPC\nMCPToolService"]
end
subgraph Core["Gateway Core (Actor Runtime)"]
IR["Identity\nResolver"]
Router["RouterActor"]
Policy["PolicyActor"]
CB["CredentialBroker\nActor"]
Reg["RegistrarActor\n(cluster singleton)"]
TS1["ToolSupervisor\n tool-a"]
TS2["ToolSupervisor\n tool-b"]
TS3["ToolSupervisor\n tool-c"]
S1["SessionActor"]
S2["SessionActor"]
S3["SessionActor"]
Journal["JournalActor\n(bounded mailbox)"]
Health["HealthActor"]
end
subgraph Egress["Tool Backends"]
T1["Stdio Tool\nchild process"]
T2["HTTP Tool\nremote MCP server"]
T3["gRPC Tool\nprotobuf service"]
end
subgraph Observability["Observability"]
OT["OpenTelemetry\n(OTLP)"]
Audit["Audit Sink\n(file, custom)"]
end
C1 & C2 & C3 --> HTTP & SSE & WS & GRPC
HTTP & SSE & WS & GRPC --> IR
IR --> Router
Router --> Policy
Router --> CB
Router --> Reg
Reg --> TS1 & TS2 & TS3
TS1 --> S1 --> T1
TS2 --> S2 --> T2
TS3 --> S3 --> T3
Router --> Journal
Health -.->|probe| TS1 & TS2 & TS3
Journal --> Audit
Core --> OT核心概念
网关 是唯一的公共入口点。它拥有actor系统,并提供所有生命周期和管理方法。你只创建一个 Gateway 每个进程(或每个集群节点)。
工具 是注册的MCP后端。工具具有唯一的ID、传输类型(stdio、HTTP或gRPC)、路由模式、断路器配置、凭据要求和生命周期状态。工具在启动时通过以下方式注册 mcp.Config.Tools 或动态通过 RegisterTool.
工具主管 是监督单个工具的所有会话的参与者。它跟踪断路器状态,执行 MaxSessionsPerTool 并管理会话参与者生命周期。每个注册的工具都有一个主管。
会话 是拥有特定租户+客户端对的单个工具会话的参与者。它持有 ToolExecutor (实际的传输连接),处理调用,并透明地从传输故障中恢复。在粘性路由模式下,会话在来自同一客户端的请求之间重用。
路由器 是中央调度角色。它接收每次调用,评估策略,解析凭据,选择适当的工具监督程序,并将请求路由到会话。它还异步写入审计日志。
注册员 是工具注册表参与者。它拥有权威的注册工具及其监督员名单。在集群模式下,它作为集群范围内的单例运行:所有节点共享一个注册器。
调用 是工作单元。它携带工具ID、方法、参数、租户和客户端身份、解析的凭据以及相关元数据(请求ID、跟踪ID、接收到的时间戳)。
执行结果 这就是结果。它携带执行状态(success, failure, timeout, denied, throttled),输出映射、错误详细信息、挂钟持续时间和相关性元数据。
断路器 按工具计算。它从 closed (正常)→ open (连续N次故障后快速失效)→ half-open (有限请求的探测)→ closed 关于恢复。电路状态可通过以下方式查看 GetToolStatus 并作为度量导出。
请求的生命周期
以下序列显示了 tools/call 或 resources/read 从客户端到后端并返回的调用。
sequenceDiagram
participant Client
participant Ingress
participant IdentityResolver
participant RouterActor
participant PolicyActor
participant CredentialBroker
participant ToolSupervisor
participant SessionActor
participant ToolBackend
participant JournalActor
Client->>Ingress: POST /mcp (tools/call, resources/read) or gRPC CallTool
Ingress->>IdentityResolver: ResolveIdentity(request) or ResolveGRPCIdentity(ctx)
IdentityResolver-->>Ingress: TenantID, ClientID
Ingress->>RouterActor: RouteInvocation{inv}
RouterActor->>PolicyActor: EvaluatePolicy(TenantID, ToolID, quotas)
PolicyActor-->>RouterActor: allow / deny
RouterActor->>CredentialBroker: ResolveCredentials(TenantID, ToolID)
CredentialBroker-->>RouterActor: Credentials
RouterActor->>ToolSupervisor: SessionInvoke(inv + credentials)
ToolSupervisor->>SessionActor: Execute(invocation)
SessionActor->>ToolBackend: MCP tools/call or resources/read (stdio or HTTP)
ToolBackend-->>SessionActor: result
SessionActor-->>ToolSupervisor: ExecutionResult
ToolSupervisor-->>RouterActor: ExecutionResult
RouterActor->>JournalActor: AuditEvent (async, non-blocking)
RouterActor-->>Ingress: ExecutionResult
Ingress-->>Client: SSE event / JSON response每个MCP会话进行一次身份解析(在 initialize 时间)用于HTTP/SSE/Webocket传输。会话中的所有后续请求都重用已解析的标识,而无需重新调用解析器。对于gRPC入口,身份在每次请求时都会通过以下方式解析 GRPCIdentityResolver 来自gRPC元数据。每次调用时都会进行策略评估、凭据解析和审计日志记录。
运输支持
入口:服务MCP客户
goakt mcp揭示了三种工厂化方法 Gateway 创建ingress HTTP处理程序,以及gRPC的注册方法。安装返回 http.Handler 在您自己的HTTP服务器或路由器中;自行注册gRPC服务 grpc.Server.
Gateway.Handler:可流式HTTP,MCP 2025-11-25规范,当前MCP客户端,Claude Desktop,代理Gateway.SSEHandler:服务器发送事件、MCP 2024-11-05规范、传统客户端、基于浏览器的代理Gateway.WSHandler:WebSocket、全双工流、延迟敏感、双向流Gateway.RegisterGRPCService:gRPC、Protobuf/gRPC、服务到服务、流式传输进度、多语言
这三个HTTP处理程序接受 mcp.IngressConfig 指定 IdentityResolver、会话空闲超时和状态模式。
在 有状态的 模式(默认),每个客户端都会收到一个唯一的 Mcp-Session-Id;身份解析器在每个连接中运行一次,所有后续请求都重用解析的身份。这是长期代理会话的正确选择。
在 无状态 模式(Stateless: true),为每个HTTP请求创建一个新会话。这在无法保证请求粘性的负载均衡器之后效果很好。
gRPC入口
gRPC入口暴露了 MCPToolService 带有三个RPC的protobuf服务:
ListTools:返回所有已注册的工具及其模式CallTool:同步工具调用:发送请求,获取结果CallToolStream:流式调用:在进度事件到达时交付进度事件,然后交付最终结果
与HTTP处理程序不同(它返回 http.Handler),gRPC服务直接在 grpc.Server。每次请求都会通过以下方式解决身份问题 GRPCIdentityResolver,它从gRPC元数据中读取:相当于HTTP标头的gRPC。
服务器设置:
// 1. Create your gRPC server (add interceptors, TLS, etc. as needed)
grpcServer := grpc.NewServer()
// 2. Register the MCP service with your identity resolver
gw.RegisterGRPCService(grpcServer, mcp.GRPCIngressConfig{
IdentityResolver: &myMetadataResolver{},
})
// 3. Serve
lis, _ := net.Listen("tcp", ":50051")
grpcServer.Serve(lis)客户端使用情况(任何gRPC支持的语言):
conn, _ := grpc.NewClient("localhost:50051", grpc.WithTransportCredentials(insecure.NewCredentials()))
client := pb.NewMCPToolServiceClient(conn)
// Attach identity via gRPC metadata
md := metadata.Pairs("x-tenant-id", "acme", "x-client-id", "agent-1")
ctx := metadata.NewOutgoingContext(context.Background(), md)
// Discover tools
resp, _ := client.ListTools(ctx, &pb.ListToolsRequest{})
// Call a tool
args, _ := json.Marshal(map[string]any{"path": "/tmp"})
result, _ := client.CallTool(ctx, &pb.CallToolRequest{
ToolName: "list_directory",
Arguments: args,
})
// Stream a tool call with progress events
stream, _ := client.CallToolStream(ctx, &pb.CallToolStreamRequest{
ToolName: "list_directory",
Arguments: args,
})
for {
msg, err := stream.Recv()
if err == io.EOF { break }
switch p := msg.GetPayload().(type) {
case *pb.CallToolStreamResponse_Progress:
log.Printf("progress: %s", p.Progress.GetMessage())
case *pb.CallToolStreamResponse_Result:
log.Printf("result: %v", p.Result.GetContent())
}
}工具参数是JSON编码的字节,而不是类型化的原型字段。这允许gRPC入口转发任意工具模式,而不需要编译 .pb.go 每个后端工具的类型。
gRPC的OAuth/OIDC身份验证 使用拦截器。这 GRPCAuthInterceptors 函数返回一元和流拦截器,用于验证来自 authorization gRPC元数据密钥,强制执行所需的作用域,并将验证的令牌信息存储在上下文中:
unary, stream, _ := goaktmcp.GRPCAuthInterceptors(&mcp.EnterpriseAuthConfig{
TokenVerifier: myVerifier,
RequiredScopes: []string{"tools:read"},
})
srv := grpc.NewServer(
grpc.ChainUnaryInterceptor(unary),
grpc.ChainStreamInterceptor(stream),
)
gw.RegisterGRPCService(srv, mcp.GRPCIngressConfig{
EnterpriseAuth: &mcp.EnterpriseAuthConfig{
TokenVerifier: myVerifier,
},
// IdentityResolver auto-installed from token claims
})工具名称缓存: 默认情况下, CallTool 和 CallToolStream 将工具名称缓存到ToolID映射中5秒(DefaultToolCacheTTL)为了避免 ListTools 演员问每一个要求。集 ToolCacheTTL 上 GRPCIngressConfig 调整或禁用(-1)缓存。
出口:连接到工具后端
goakt mcp原生支持三种后端传输类型。
工作室 启动一个子进程,并使用MCP协议通过stdin/stdout与之通信。这是本地安装的MCP服务器(文件系统、shell、代码解释器等)的标准传输。该进程受到监督:如果它崩溃,会话参与者将在下次调用时创建一个新的进程。
超文本传输协议 通过HTTP连接到远程MCP服务器。使用完整的MCP流式HTTP语义,包括会话管理和流式传输。启用跟踪时,W3C跟踪上下文标头会在每次出站调用时传播。
gRPC 使用动态protobuf消息连接到远程gRPC服务。原型描述符可以从本地加载 .binpb 文件或通过gRPC服务器反射获取。JSON模式是从原型消息描述符自动派生出来的。通过以下方式支持服务器流式RPC ToolStreamExecutor 界面。
这三种传输方式都在注册时获取后端的工具模式并缓存它们。网关使用实际的工具名称、描述和JSON模式来构建入口服务器的工具注册表,为MCP客户端提供准确、可发现的模式信息。
对于HTTP和stdio后端,网关还在注册时发现MCP资源和资源模板。看 MCP资源 下面的部分。
MCP资源
MCP规范定义了 资源 作为服务器管理的数据,客户端可以发现和读取。资源补充了工具:在工具执行操作的地方,资源公开数据(文件、数据库行、API响应、配置等)以供LLM代理检索。
goakt mcp代理入口客户端和出口后端之间的完整mcp资源生命周期。
发现
当工具注册时(启动时通过 Config.Tools 或动态通过 RegisterTool),网关连接到后端并调用 resources/list 和 resources/templates/list 与现有 tools/list 模式获取。发现的资源元数据(URI、名称、描述、MIME类型)和资源模板元数据(URI模板、名称、说明、MIME类型 ListTools.
在每个新的MCP客户端会话上,入口层在每个会话SDK服务器上注册发现的资源和资源模板。SDK会自动播发 resources 能力在 initialize 响应,使MCP客户端能够调用 resources/list 浏览可用资源。
gRPC后端不公开MCP资源(它们使用protobuf服务定义),并在发现过程中返回空元数据。
阅读
当MCP客户端呼叫时 resources/read 使用资源URI,请求将通过与以下相同的参与者管道 tools/call :
- 入口 翻译SDK
ReadResourceRequest进入网关Invocation用方法resources/readURI在Params["uri"] - 路由器 评估策略、解析凭据并路由到工具主管
- 会话 发送给遗嘱执行人
ResourceExecutor.ReadResource通过类型断言的方法 - 执行者 电话
ClientSession.ReadResource在后端MCP服务器上 - 响应通过管道返回,具有与工具调用相同的断路、执行器恢复、钝化和审计日志保证
OAuth作用域传播适用:来自已验证承载令牌的作用域附加到调用中,并可供自定义 PolicyEvaluator 实现对资源访问的范围感知授权。
资源模板
资源模板使用 RFC 6570 URI模板(例如。 file:///{path})以公开参数化资源。网关通过以下方式发现模板 resources/templates/list 并通过以下方式在每会话SDK服务器上注册它们 AddResourceTemplate。当客户端读取模板化URI时,SDK会解析模板并将其分派给相同的URI resources/read 处理程序。
执行者支持
这 ResourceExecutor 接口是的可选扩展 ToolExecutor :
type ResourceExecutor interface {
ReadResource(ctx context.Context, inv *Invocation) (*ExecutionResult, error)
}内置的HTTP和stdio执行器实现 ResourceExecutorgRPC执行器没有(gRPC工具是基于protobuf的,而不是基于MCP资源的)。会话参与者检查 ResourceExecutor 在运行时通过类型断言;未实现它的执行器返回以下错误 resources/read 请求:
自定义 ToolExecutor 实现可以通过实现来选择资源支持 ResourceExecutor 一起 ToolExecutor.
多租户和授权
每次调用都归因于 租户ID 和 可点击,在会话创建时由解决 IdentityResolver 你提供。
身份解析
对于 HTTP/SSE/Webocket ingress,解析器接收原始HTTP请求:
type IdentityResolver interface {
ResolveIdentity(r *http.Request) (TenantID, ClientID, error)
}常见的实现读取JWT声明、API密钥头、mTLS证书主题或自定义会话令牌。非nil错误会拒绝使用HTTP 400的传入会话。解析器在每个MCP会话中运行一次:所有后续请求都重用已解析的标识。
对于 gRPC ingress,解析器接收gRPC请求上下文:
type GRPCIdentityResolver interface {
ResolveGRPCIdentity(ctx context.Context) (TenantID, ClientID, error)
}实现通常通过以下方式从gRPC元数据读取 metadata.FromIncomingContext(ctx)非nil错误返回gRPC状态 Unauthenticated 给来电者。与HTTP解析器不同,gRPC解析器在每个请求上运行。
租户配置
每个租户都可以配置独立的配额和策略设置。未在配置中列出的租户将被接受,而不会强制执行配额。
ID:不透明字符串标识符:必须匹配IdentityResolver回报Quotas.RequestsPerMinute:每分钟最大调用次数;零表示无限制Quotas.ConcurrentSessions:任何时候的最大直播时段;零表示无限制Evaluator:可选自定义PolicyEvaluator对于这个租户
政策执行
goakt mcp在每次调用时应用两层策略检查:
- 内置检查 :速率限制、并发上限和工具授权(配置时租户分配列表)
- 自定义评估器 :您的
PolicyEvaluator实现,仅在所有内置检查通过时调用
type PolicyEvaluator interface {
Evaluate(ctx context.Context, input PolicyInput) *RuntimeError
}PolicyInput 携带 TenantID, ToolID、实时会话计数和当前分钟的请求,以便评估人员能够做出基于上下文的决策。返回 nil 允许,或 *RuntimeError 和 ErrCodePolicyDenied 否认。常见的模式包括OPA集成、基于属性的访问控制(ABAC)和特定于环境的允许列表。
每个政策决定:允许或拒绝:都记录在审计日记中。
凭证经纪
一些工具后端需要身份验证机密(API密钥、令牌、证书)。goaktmcp的凭证代理代表会话解析机密,并将其注入到每次调用中,从而将机密从应用程序代码和请求参数中排除。
type CredentialsProvider interface {
ID() string
ResolveCredentials(ctx context.Context, tenantID TenantID, toolID ToolID) (*Credentials, error)
}可以注册多个提供商。代理按顺序尝试每个提供者,并返回第一个非nil结果。解析后的凭据缓存在可配置的LRU缓存中(TTL和max条目都是可调的),以避免在热路径上重复的秘密存储往返。
凭据的范围为租户+工具对。这意味着不同的租户可以对同一个工具拥有不同的凭据,而单个租户可以对不同的工具持有不同的凭据。
这 CredentialPolicy 每个工具上的字段控制凭据是否 optional (如果解决失败,请继续使用空凭据)或 required (如果解析失败,则调用失败)。
韧性
断路器
每个工具都有一个独立的断路器,可以保护系统的其他部分免受后端行为不端的影响。电路状态由跟踪 ToolSupervisor 演员和自动转换。
stateDiagram-v2
[*] --> Closed: Initial state
Closed --> Open: ≥ FailureThreshold consecutive failures
Open --> HalfOpen: OpenDuration elapsed
HalfOpen --> Closed: Probe succeeds
HalfOpen --> Open: Probe failsFailureThreshold(默认值5):打开前连续出现故障OpenDuration(默认值30 s):探测前电路保持打开的时间HalfOpenMaxRequests(默认值1):半开放探测期间的最大并发请求数
当电路断开时,调用会立即失败 ErrCodeToolUnavailable 而无需联系后端。这可以防止共享同一工具的租户之间发生级联故障。
这 ResetCircuit 管理员方法手动关闭开路,在操作员确认后端正常后立即恢复。
会话执行器恢复
当会话参与者检测到传输失败时: ToolExecutor.Execute 或结果 ErrCodeTransportFailure :它透明地尝试恢复:
- 失败的执行器已关闭
- 通过以下方式创建新的执行器
ExecutorFactory - 使用新的超时上下文重试一次调用
- 如果恢复失败,则返回原始错误并递增断路器
这意味着在同一请求中可以恢复stdio进程崩溃或HTTP连接中断,而不需要会话钝化、重新创建或客户端重试。
背压
这 MaxSessionsPerTool 工具定义上的字段为并发会话设置了硬上限。当达到限制时,新的调用将收到 ErrCodeConcurrencyLimitReached 立即。这可以防止单个过载的工具耗尽共享资源。
审计日记参与者使用 有界邮箱 具有可配置的容量(AuditConfig.MailboxSize).当邮箱已满时,发件人会阻止,直到空间可用。这在审计路径上提供了自然的背压,而不会丢失事件。
健康探索
这 HealthActor 定期探测 ToolSupervisor 检查其工具是否响应。探测结果驱动工具的运行 State 字段:
enabled:工具状态良好,可以接受请求degraded:工具响应缓慢或出现间歇性错误unavailable:工具无法访问;电路可能开路disabled:工具已在管理上禁用
健康状态转换记录在审计日志中,并作为指标导出。
可观测性
开放遥测指标
使用 WithMetrics() 选项。指标通过OTLP导出到中配置的端点 mcp.Config.Telemetry.OTLPEndpoint.
- 调用延迟 (柱状图):每个工具的端到端持续时间
- 调用失败 (计数器):调用失败,按工具和错误代码标记
- 工具状态 (仪表):每个工具的操作状态(启用/降级/不可用/禁用)
- 电路状态 (仪表):每个工具的断路器状态(闭合/打开/半打开)
- 会话生命周期 (计数器):会话创建和终止事件
- 健康转型 (计数器):刀具健康状态更改
分布式跟踪
使用启用跟踪 WithTracing() 选项。为每次调用创建跟踪,并跨越从入口到出口的完整路径。W3C traceparent 和 tracestate 标头在所有到工具后端的出站HTTP调用上传播,从而在整个AI工具堆栈中实现端到端的分布式跟踪。
结构化日志记录
goaktmcp使用可插拔的日志接口。您可以通过以下方式提供自己的日志后端(zap、zerolog、slog、logrus等) WithLogger(logger),或通过声明方式设置日志级别 mcp.Config.LogLevel.
type Logger interface {
Debug(msg string, args ...any)
Info(msg string, args ...any)
Warn(msg string, args ...any)
Error(msg string, args ...any)
}所有日志行都包含结构化的相关字段:租户ID、工具ID、请求ID和跟踪ID。如果您的 Logger 还实现了可选 LeveledLogger 接口(Level() string),适配器将其用于发动机侧日志门控;否则默认为 info.
持久审计追踪
审计日记账以结构化的方式记录每一个重大事件 AuditEvent 并将其写入您的 AuditSink 异步实现。
type AuditSink interface {
Write(event *AuditEvent) error
Close() error
}policy_decision:任何调用都会达到策略评估(允许或拒绝)invocation_start:对工具后端开始执行调用invocation_complete:调用成功完成invocation_failed:调用失败(超时、传输错误等)health_transition:工具的操作状态发生变化circuit_state_change:断路器在状态之间转换
每个事件都携带租户ID、客户端ID、工具ID、请求ID、跟踪ID、结果、错误代码和自由形式的元数据映射。内置水槽包括 MemorySink (用于测试)以及 FileSink (NDJSON行分隔文件)。自定义接收器可以写入任何持久存储:Postgres、Kafka、BigQuery、S3等。
群集模式
goakt mcp支持多节点操作。启用集群模式时,网关节点使用基于八卦的成员资格和GoAkt远程处理形成对等集群,用于分布式参与者通信。
graph TB
LB["Load Balancer\n(nginx / K8s Ingress)"]
subgraph Node0["Node: gateway-0"]
GM0["GatewayManager\ngateway-manager-gateway-0"]
Router0["RouterActor"]
Reg["RegistrarActor\n★ cluster singleton"]
end
subgraph Node1["Node: gateway-1"]
GM1["GatewayManager\ngateway-manager-gateway-1"]
Router1["RouterActor"]
end
subgraph Node2["Node: gateway-2"]
GM2["GatewayManager\ngateway-manager-gateway-2"]
Router2["RouterActor"]
end
subgraph Discovery["Discovery Provider"]
K8s["Kubernetes\npod label selector"]
DNS["DNS-SD\nservice domain"]
Custom["Custom\nimplementation"]
end
LB --> Node0 & Node1 & Node2
Node0 |"Gossip + Remoting"| Node1
Node1 |"Gossip + Remoting"| Node2
Node0 |"Gossip + Remoting"| Node2
Router1 & Router2 -.->|"Remote Ask"| Reg
Discovery --> Node0 & Node1 & Node2运作原理
- 每个节点都运行自己的
GatewayManager带有主机名后缀的名称(gateway-manager-)因此,生成多个节点永远不会引发集群范围内的名称冲突。 - 这
RegistrarActor(工具注册表)作为 集群单例 :整个集群中只存在一个实例,在首先启动的节点上选择。不承载单例的节点通过GoAkt的分布式参与者消息传递到达它。 RouterActor在每个节点上本地运行。它按名称解析单例注册器,并将调用透明地路由到本地或远程监控器。- 工具会议和主管在注册商放置的地方进行。在当前拓扑中,它们是单例节点的本地节点;集群感知路由是一种进化路径。
同行发现
实施 DiscoveryProvider 接口来控制节点如何找到彼此。
type DiscoveryProvider interface {
ID() string
Start(ctx context.Context) error
DiscoverPeers(ctx context.Context) ([]string, error)
Stop(ctx context.Context) error
}DiscoverPeers 返回对等地址列表(主机:发现端口的端口)。goakt mcp定义了 DiscoveryProvider 接口:您提供与您的基础设施(Kubernetes、Consul、DNS-SD、静态列表等)相匹配的实现。这 集群示例 包括一个基于Kubernetes的示例提供者。
传输层安全
集 ClusterConfig.TLS 到一个 RemotingTLSConfig 为所有远程处理和集群通信启用TLS。所有节点必须共享相同的根CA。支持双向TLS(客户端证书验证)。
群集感知关闭
对于Kubernetes中的干净有序关闭,请在删除StatefulSet之前将其扩展到零副本 podManagementPolicy: OrderedReadyKubernetes以相反的顺序终止Pod:每个离开的节点都有活动的对等体来复制参与者状态,从而防止复制错误。这 集群示例Makefile cluster-down target展示了这种模式。
配置参考
构建一个 mcp.Config 并将其传递给 goaktmcp.New。零值字段用安全默认值填充。
运行时
控制actor运行时的超时和探测间隔。
SessionIdleTimeout(默认值5 min):在此段不活动期后,将会话参与者钝化RequestTimeout(默认值30 s):单个参与者Ask(路由器、注册器、会话)的最长时间StartupTimeout(默认值10 s):等待演员系统准备就绪的最长时间HealthProbeInterval(默认值30 s):健康行动者多久对工具主管进行一次调查ShutdownTimeout(默认值30 s):正常关机的最长时间
簇
Enabled:激活集群模式DiscoveryProvider:启用时必需:对等体发现实现DiscoveryPort:发现协议的端口(默认15000)PeersPort:八卦成员列表协议的端口(默认15000)RemotingPort:GoAkt参与者到参与者远程处理的端口(默认15001)RegistrarRole:将单例注册器固定到特定节点的可选集群角色TLS:可选*RemotingTLSConfig用于加密远程处理
遥测
OTLPEndpoint:用于指标和跟踪导出的OTLP HTTP端点(例如。http://otel-collector:4318)
审计
Sink(默认值MemorySink) :mcp.AuditSink实施;使用FileSink或定制水槽用于生产MailboxSize(默认值1024):发送方阻塞前的最大运行中审核事件数(背压)
凭证
Providers(默认值:):订购的切片mcp.CredentialsProvider实现CacheTTL(默认值:):重新获取之前缓存已解析的凭据需要多长时间MaxCacheEntries(默认值:):最大缓存条目数;满时驱逐最近最少使用的
租户
ID:租户标识符:必须匹配IdentityResolver输出Quotas.RequestsPerMinute:每分钟请求速率限制;零=无限制Quotas.ConcurrentSessions:同期会议上限;零=无限制Evaluator:可选mcp.PolicyEvaluator用于自定义授权
工具定义
每个条目 Config.Tools (或致电 RegisterTool)接受以下字段。
ID:唯一工具标识符Transport:stdio,http,或grpcStdio:*StdioTransportConfig:command、args、env、工作目录HTTP:*HTTPTransportConfig:URL和可选TLS配置GRPC:*GRPCTransportConfig:目标、服务、方法、TLS、元数据、描述符集或反射State:初始状态:enabled或disabledRouting:sticky(会话关联性,默认)或least_loadedMaxSessionsPerTool:背压限制;零=无限制RequestTimeout:每次工具调用超时;覆盖运行时默认值IdleTimeout:在此持续时间之后,将空闲会话设为被动StartupTimeout:执行器创建超时Circuit:CircuitConfig:失败阈值、打开持续时间、半打开最大请求数CredentialPolicy:optional(默认)或requiredAuthorizationPolicy:tenant_allowlist将工具限制为特定租户
公共API
这 Gateway struct是唯一的公共入口点。从多个goroutine并发调用所有方法都是安全的。
生命周期
- **
New(cfg, ...opts) (*Gateway, error)** :构建网关;验证配置并应用选项 Start(ctx) error:启动actor系统,生成运行时actor,并引导配置的工具Stop(ctx) error:优雅地耗尽正在进行的调用并关闭actor系统System() ActorSystem:访问底层GoAkt参与者系统(例如用于自定义参与者交互)
工具调用
- **
Invoke(ctx, inv) (*ExecutionResult, error)** :同步工具执行:阻塞直到后端响应或超时 - **
InvokeStream(ctx, inv) (*StreamingResult, error)** :流媒体执行;返回aStreamingResult带着一个Progress频道和aFinal频道
工具管理
ListTools(ctx) ([]Tool, error):返回所有已注册的工具及其当前架构RegisterTool(ctx, tool) error:动态注册新工具或替换现有工具UpdateTool(ctx, tool) error更新工具的元数据;工具必须存在EnableTool(ctx, toolID) error:重新启用以前禁用的工具DisableTool(ctx, toolID) error:在行政上禁用工具;飞行中请求完成,新请求被拒绝RemoveTool(ctx, toolID) error:从注册表中完全删除工具GetToolSchema(ctx, toolID) ([]ToolSchema, error):返回从后端发现的缓存MCP架构
管理和运营API
这些方法提供对正在运行的网关的实时可见性和控制。他们在为交通服务时可以安全地拨打电话,并立即返回。
- **
GetGatewayStatus(ctx) (*GatewayStatus, error)** :总体状态:运行标志、注册工具计数、总活动会话计数 - **
GetToolStatus(ctx, toolID) (*ToolStatus, error)** :每个工具状态:操作状态、电路状态、会话计数、排放标志、缓存模式 ListSessions(ctx) ([]SessionInfo, error):所有工具中的每个活动会话,包括工具ID、租户ID和客户端IDDrainTool(ctx, toolID) error:在现有会话结束时,停止接受工具的新会话;禁用或删除前使用ResetCircuit(ctx, toolID) error:确认后端正常后,手动关闭断路器
入口处理程序
Handler(cfg IngressConfig) (http.Handler, error):MCP流式HTTP处理程序(2025-11-25规范)SSEHandler(cfg IngressConfig) (http.Handler, error):服务器发送事件处理程序(2024-11-05规范)- **
WSHandler(cfg IngressConfig, wsCfg *WSConfig) (http.Handler, error)** :WebSocket处理程序 - **
RegisterGRPCService(srv *grpc.Server, cfg GRPCIngressConfig) error** :注册具有流媒体支持的MCPToolService gRPC服务 - **
GRPCAuthInterceptors(ea *EnterpriseAuthConfig) (unary, stream, error)** :gRPC入口的承载令牌身份验证拦截器
选项
WithLogger(logger Logger):插入自定义日志后端(zap、zerolog、slog、logrus等)WithDebug():启用从actor引擎到stdout的详细调试日志记录WithMetrics():启用OpenTetry指标导出WithTracing():启用OpenTetry跟踪和W3C跟踪上下文传播
关键接口
以下接口是定制goakt-mcp行为的扩展点。所有这些都在 mcp 包裹。
身份解析(HTTP) :每个新的HTTP/SSE/Webocket会话调用一次:
type IdentityResolver interface {
ResolveIdentity(r *http.Request) (TenantID, ClientID, error)
}身份解析(gRPC) :每个gRPC请求调用一次:
type GRPCIdentityResolver interface {
ResolveGRPCIdentity(ctx context.Context) (TenantID, ClientID, error)
}政策评估 :在每次调用时调用,经过内置检查:
type PolicyEvaluator interface {
Evaluate(ctx context.Context, input PolicyInput) *RuntimeError
}凭证解析 :每次调用时由凭证代理调用:
type CredentialsProvider interface {
ID() string
ResolveCredentials(ctx context.Context, tenantID TenantID, toolID ToolID) (*Credentials, error)
}审计持久性 :接收每一个重要事件:
type AuditSink interface {
Write(event *AuditEvent) error
Close() error
}群集对等体发现 :返回八卦和远程处理的实时对等地址:
type DiscoveryProvider interface {
ID() string
Start(ctx context.Context) error
DiscoverPeers(ctx context.Context) ([]string, error)
Stop(ctx context.Context) error
}自定义工具执行器 :自带交通工具:
type ToolExecutor interface {
Execute(ctx context.Context, inv *Invocation) (*ExecutionResult, error)
Close() error
}定制执行器工厂 :为每个会话创建新的执行器:
type ExecutorFactory interface {
Create(ctx context.Context, tool Tool, credentials map[string]string) (ToolExecutor, error)
}示例
所有示例均在 examples/ 目录,可以运行 go run ./examples/.
- 文件系统 :使用stdio文件系统工具的最小网关
- 审核http :具有HTTP出口工具的持久文件审计接收器
- 入口 :具有基于标头的身份解析的MCP流式HTTP入口
- 入口grpc :gRPC入口,具有基于元数据的身份解析、ListTools、CallTool和CallToolStream
- 管理员 :完全管理员API和自定义
PolicyEvaluator - 配额 :每个租户的速率限制和并发执行
- 完整配置 :涵盖每个字段的完整配置参考
- 人工智能中心 :端到端多对端人工智能工具集线器示例:stdio+HTTP出口、可流式HTTP入口、可插拔策略、凭据代理、持久审核、OpenTelemetry和完整管理员API
- 簇 :具有Kubernetes对等体发现、nginx会话关联和Jaeger跟踪的三节点Kubernetes集群
这 人工智能中心 示例是理解所有部分如何端到端地组合在一起的推荐起点。这 簇 示例包括一个完整的 Makefile Kubernetes清单用于部署到本地 亲切的 集群。
贡献
欢迎投稿!请阅读 贡献指南 有关开发工作流、代码标准、测试和拉取请求流程的详细信息。
安装
go get github.com/tochemey/goakt-mcpgoakt mcp需要Go 1.26或更高版本。这 mcp 包包含所有公共域类型。这 goaktmcp 根包公开 Gateway 以及它的选择。
许可证
本项目根据以下条款获得许可 许可证.
