Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -99,3 +99,5 @@ dump.rdb
/data/
# telemetry dumped by the otel collector during tests
ci/pod/otelcol-contrib/data-otlp.json
# bytecode from importing the .py test drivers
t/plugin/__pycache__/
8 changes: 8 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -419,6 +419,14 @@ install: runtime
$(ENV_INSTALL) apisix/plugins/mcp/broker/*.lua $(ENV_INST_LUADIR)/apisix/plugins/mcp/broker
$(ENV_INSTALL) apisix/plugins/mcp/transport/*.lua $(ENV_INST_LUADIR)/apisix/plugins/mcp/transport

$(ENV_INSTALL) -d $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp/openapi
$(ENV_INSTALL) -d $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp/tools
$(ENV_INSTALL) -d $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp/transport
$(ENV_INSTALL) apisix/plugins/openapi-to-mcp/*.lua $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp
$(ENV_INSTALL) apisix/plugins/openapi-to-mcp/openapi/*.lua $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp/openapi
$(ENV_INSTALL) apisix/plugins/openapi-to-mcp/tools/*.lua $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp/tools
$(ENV_INSTALL) apisix/plugins/openapi-to-mcp/transport/*.lua $(ENV_INST_LUADIR)/apisix/plugins/openapi-to-mcp/transport

$(ENV_INSTALL) -d $(ENV_INST_LUADIR)/apisix/plugins/jwt-auth
$(ENV_INSTALL) apisix/plugins/jwt-auth/*.lua $(ENV_INST_LUADIR)/apisix/plugins/jwt-auth

Expand Down
1 change: 1 addition & 0 deletions apisix/cli/config.lua
Original file line number Diff line number Diff line change
Expand Up @@ -269,6 +269,7 @@ local _M = {
"traffic-split",
"redirect",
"response-rewrite",
"openapi-to-mcp",
"oas-validator",
"mcp-bridge",
"degraphql",
Expand Down
2 changes: 1 addition & 1 deletion apisix/cli/ngx_tpl.lua
Original file line number Diff line number Diff line change
Expand Up @@ -507,7 +507,7 @@ http {
lua_shared_dict ext-plugin {* http.lua_shared_dict["ext-plugin"] *}; # cache for ext-plugin
{% end %}

{% if enabled_plugins["mcp-bridge"] then %}
{% if enabled_plugins["mcp-bridge"] or enabled_plugins["openapi-to-mcp"] then %}
lua_shared_dict mcp-session {* http.lua_shared_dict["mcp-session"] *}; # cache for mcp-session
{% end %}

Expand Down
165 changes: 165 additions & 0 deletions apisix/plugins/openapi-to-mcp.lua
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
--
-- Licensed to the Apache Software Foundation (ASF) under one or more
-- contributor license agreements. See the NOTICE file distributed with
-- this work for additional information regarding copyright ownership.
-- The ASF licenses this file to You under the Apache License, Version 2.0
-- (the "License"); you may not use this file except in compliance with
-- the License. You may obtain a copy of the License at
--
-- http://www.apache.org/licenses/LICENSE-2.0
--
-- Unless required by applicable law or agreed to in writing, software
-- distributed under the License is distributed on an "AS IS" BASIS,
-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-- See the License for the specific language governing permissions and
-- limitations under the License.
--
local core = require("apisix.core")
local streamable_http = require("apisix.plugins.openapi-to-mcp.transport.streamable_http")
local mcp_sse = require("apisix.plugins.openapi-to-mcp.transport.sse")
local ngx = ngx
local pairs = pairs

local schema = {
type = "object",
properties = {
transport = {
description = "The transport mechanisms for client-server communication",
type = "string",
default = "sse",
enum = {"sse", "streamable_http"},
},
openapi_url = {
description = "URL of the OpenAPI specification document",
type = "string",
minLength = 1,
},
base_url = {
description = "Base URL of the external service",
type = "string",
minLength = 1,
},
headers = {
description = "Headers to include in requests to the external service",
type = "object",
minProperties = 0,
patternProperties = {
["^[^:]+$"] = {
oneOf = {
{ type = "string" }
}
}
},
},
flatten_parameters = {
description = "Whether to flatten query and path parameters " ..
"in the tool inputSchema. When false (default), " ..
"parameters are nested under queryParameters or pathParameters. " ..
"When true, parameters are placed directly in properties.",
type = "boolean",
default = false,
},
},
required = { "openapi_url", "base_url" },
}

local plugin_name = "openapi-to-mcp"

local _M = {
version = 0.1,
priority = 540,
name = plugin_name,
schema = schema,
}


function _M.check_schema(conf)
return core.schema.check(schema, conf)
end


-- Resolve base_url and headers against request variables.
local function resolve_conf(conf, ctx)
local base_url, err = core.utils.resolve_var(conf.base_url, ctx.var)
if err then
core.log.error("failed to resolve variable for base_url: ",
conf.base_url, ", error: ", err)
base_url = conf.base_url
end

local headers = {}
for key, value in pairs(conf.headers or {}) do
local resolved_value, herr = core.utils.resolve_var(value, ctx.var)
if herr then
core.log.error("failed to resolve variable for header, key: ", key,
", error: ", herr)
resolved_value = value
end
headers[key] = resolved_value
end

return base_url, headers
end


function _M.access(conf, ctx)
if conf.transport == "streamable_http" then
local base_url, headers = resolve_conf(conf, ctx)

-- Defer the answer to before_proxy. Exiting here would skip every
-- plugin with a lower priority that still has to run in access, such
-- as an authorization check on the tool being called.
ctx.mcp_inprocess_opts = {
conf = conf,
base_url = base_url,
headers = headers,
transport = "streamable_http",
}
-- The answer is produced in before_proxy, so the request never reaches
-- an upstream; handle_upstream() runs before_proxy and returns.
ctx.bypass_nginx_upstream = true
return
end

if conf.transport ~= "sse" then
core.log.error("Invalid MCP transport: ", conf.transport)
return 500, { message = "Invalid MCP transport"}
end

local base_url, headers = resolve_conf(conf, ctx)

-- The client is told to POST its messages to the path the route matched.
local message_path = ctx.curr_req_matched and ctx.curr_req_matched._path

ngx.ctx.disable_proxy_buffering = true
ctx.mcp_inprocess_opts = {
conf = conf,
base_url = base_url,
headers = headers,
transport = "sse",
message_path = message_path,
}
ctx.bypass_nginx_upstream = true
end


function _M.before_proxy(conf, ctx)
local opts = ctx.mcp_inprocess_opts
if not opts then
return
end

-- Every body the transports produce is JSON, and core.response.exit() sets
-- no content type of its own, so without this they would all go out as
-- text/plain. Set once here so no rejection path can miss it; the two that
-- stream override it with text/event-stream on their way out.
core.response.set_header("Content-Type", "application/json")

if opts.transport == "sse" then
return mcp_sse.handle(ctx, opts)
end
return streamable_http.handle(ctx, opts)
end


return _M
74 changes: 74 additions & 0 deletions apisix/plugins/openapi-to-mcp/cache.lua
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
--
-- Licensed to the Apache Software Foundation (ASF) under one or more
-- contributor license agreements. See the NOTICE file distributed with
-- this work for additional information regarding copyright ownership.
-- The ASF licenses this file to You under the Apache License, Version 2.0
-- (the "License"); you may not use this file except in compliance with
-- the License. You may obtain a copy of the License at
--
-- http://www.apache.org/licenses/LICENSE-2.0
--
-- Unless required by applicable law or agreed to in writing, software
-- distributed under the License is distributed on an "AS IS" BASIS,
-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-- See the License for the specific language governing permissions and
-- limitations under the License.
--
local core = require("apisix.core")
local loader = require("apisix.plugins.openapi-to-mcp.openapi.loader")
local ref = require("apisix.plugins.openapi-to-mcp.openapi.ref")
local generator = require("apisix.plugins.openapi-to-mcp.tools.generator")
local tostring = tostring

local _M = {}

-- A generated tool list is kept for an hour, for up to 100 documents.
local SPEC_TTL = 3600
local SPEC_COUNT = 100

-- A failed fetch is cached only briefly. Without neg_ttl core.lrucache caches
-- nothing on failure, which would let an unreachable spec host be re-dialed on
-- every single request; a long negative TTL would instead keep the route broken
-- long after the host recovers.
local NEG_TTL = 5
local NEG_COUNT = 32

local CACHE_VERSION = "1"

-- invalid_stale: without it core.lrucache hands an expired entry back and
-- re-arms its TTL whenever the version still matches, and the version here
-- never changes, so a document updated at the same URL would never be fetched
-- again.
local lru = core.lrucache.new({

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking: cache expiry behavior

With a constant CACHE_VERSION and invalid_stale unset, core.lrucache revives expired entries whose version still matches instead of invoking build_tools again. In a targeted probe using the existing cache wrapper with a mocked clock and loader, calls after 3,601 and 7,202 seconds still returned the original tools, with only one fetch. A document updated at the same URL can therefore stay stale beyond the advertised one-hour TTL, until eviction or worker restart.

This does not block merging the PR. Please double-check the intended refresh policy and decide whether to fix this behavior, for example by invalidating expired entries and adding a same-URL refresh regression test.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed and fixed in a6b4674. core.lrucache re-arms an expired entry whose version still matches unless invalid_stale is set, and CACHE_VERSION never changes, so a document updated at the same URL was never fetched again. The cache now sets invalid_stale = true. openapi-to-mcp-cache.t TEST 6 loads the cache with a one-second TTL, serves a document that changes on every fetch, and checks the tool list is rebuilt once after expiry (it fails without the fix).

ttl = SPEC_TTL,
count = SPEC_COUNT,
invalid_stale = true,
neg_ttl = NEG_TTL,
neg_count = NEG_COUNT,
})


local function build_tools(openapi_url, flatten_parameters)
local spec, path_order, err = loader.fetch(openapi_url)
if not spec then
return nil, err
end

local resolved = ref.resolve(spec)
return generator.generate(resolved, path_order, {
flatten_parameters = flatten_parameters,
})
end


-- Returns the tool list for a plugin conf, building it on first use.
-- base_url and headers do not take part in the key: they affect how a tool is
-- invoked, never how it is generated.
function _M.get_tools(conf)
local flatten_parameters = conf.flatten_parameters == true
local key = conf.openapi_url .. "#" .. tostring(flatten_parameters)
return lru(key, CACHE_VERSION, build_tools, conf.openapi_url, flatten_parameters)
end


return _M
96 changes: 96 additions & 0 deletions apisix/plugins/openapi-to-mcp/json_pretty.lua
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
--
-- Licensed to the Apache Software Foundation (ASF) under one or more
-- contributor license agreements. See the NOTICE file distributed with
-- this work for additional information regarding copyright ownership.
-- The ASF licenses this file to You under the Apache License, Version 2.0
-- (the "License"); you may not use this file except in compliance with
-- the License. You may obtain a copy of the License at
--
-- http://www.apache.org/licenses/LICENSE-2.0
--
-- Unless required by applicable law or agreed to in writing, software
-- distributed under the License is distributed on an "AS IS" BASIS,
-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-- See the License for the specific language governing permissions and
-- limitations under the License.
--
local core = require("apisix.core")
local str_sub = string.sub
local str_rep = string.rep
local table_concat = table.concat

local _M = {}

local INDENT_UNIT = " "


-- cjson has no pretty printer, so re-flow its compact output instead of
-- re-implementing value encoding and string escaping. Matches
-- JSON.stringify(value, null, 2): two-space indent, ": " after keys, and
-- empty containers kept on one line.
--
-- Object key order still comes from the Lua table, so the result is not
-- byte-identical to JSON.stringify for nested upstream payloads. That is a
-- known, semantically irrelevant difference.
function _M.encode(value)
local compact, err = core.json.encode(value)
if not compact then
return nil, err
end

local out = {}
local indent = 0
local in_string = false
local escaped = false
local index = 1
local length = #compact

while index <= length do
local char = str_sub(compact, index, index)

if in_string then
out[#out + 1] = char
if escaped then
escaped = false
elseif char == "\\" then
escaped = true
elseif char == '"' then
in_string = false
end

elseif char == '"' then
in_string = true
out[#out + 1] = char

elseif char == "{" or char == "[" then
local next_char = str_sub(compact, index + 1, index + 1)
if (char == "{" and next_char == "}") or (char == "[" and next_char == "]") then
out[#out + 1] = char .. next_char
index = index + 1
else
indent = indent + 1
out[#out + 1] = char .. "\n" .. str_rep(INDENT_UNIT, indent)
end

elseif char == "}" or char == "]" then
indent = indent - 1
out[#out + 1] = "\n" .. str_rep(INDENT_UNIT, indent) .. char

elseif char == "," then
out[#out + 1] = ",\n" .. str_rep(INDENT_UNIT, indent)

elseif char == ":" then
out[#out + 1] = ": "

else
out[#out + 1] = char
end

index = index + 1
end

return table_concat(out)
end


return _M
Loading
Loading