From 67a9343ada1ebc7efc50e688207b733d7f52a6eb Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 29 Sep 2026 14:55:48 +0800 Subject: [PATCH 1/4] feat(kubernetes): select clusters by discovery_args.cluster_ids --- apisix/discovery/kubernetes/init.lua | 78 +++- apisix/schema_def.lua | 9 + apisix/upstream.lua | 26 ++ docs/en/latest/discovery/kubernetes.md | 29 ++ docs/zh/latest/discovery/kubernetes.md | 29 ++ t/kubernetes/discovery/kubernetes5.t | 517 +++++++++++++++++++++++++ 6 files changed, 686 insertions(+), 2 deletions(-) create mode 100644 t/kubernetes/discovery/kubernetes5.t diff --git a/apisix/discovery/kubernetes/init.lua b/apisix/discovery/kubernetes/init.lua index f9207792356f..dccc3aed15ef 100644 --- a/apisix/discovery/kubernetes/init.lua +++ b/apisix/discovery/kubernetes/init.lua @@ -20,7 +20,9 @@ local type = type local ipairs = ipairs local pairs = pairs local string = string +local str_find = string.find local error = error +local tostring = tostring local is_http = ngx.config.subsystem == "http" local process = require("ngx.process") local core = require("apisix.core") @@ -87,7 +89,13 @@ local function single_mode_init(conf) end -local function single_mode_nodes(service_name) +local function single_mode_nodes(service_name, discovery_args) + if discovery_args and discovery_args.cluster_ids then + core.log.error("discovery_args.cluster_ids requires kubernetes discovery ", + "configured with multiple clusters, service: ", service_name) + return nil + end + return k8s_core.resolve_nodes( endpoint_lrucache, service_name, "^(.*):(.*)$", -- namespace/name:port_name @@ -154,7 +162,73 @@ local function multiple_mode_init(confs) end -local function multiple_mode_nodes(service_name) +local function merge_cluster_nodes(endpoint_dicts, endpoint_key, endpoint_port) + local nodes = {} + local seen = {} + for _, endpoint_dict in ipairs(endpoint_dicts) do + local cluster_nodes = k8s_core.create_endpoint_lrucache(endpoint_dict, endpoint_key, + endpoint_port) + for _, node in ipairs(cluster_nodes or {}) do + local addr = node.host .. ":" .. tostring(node.port) + if not seen[addr] then + seen[addr] = true + core.table.insert(nodes, node) + end + end + end + + return nodes +end + + +-- service_name is "namespace/name:port_name", the clusters come from cluster_ids +local function selected_clusters_nodes(service_name, cluster_ids) + if str_find(service_name, "/[^/]*/") then + core.log.error("service_name must not carry a cluster id prefix ", + "when discovery_args.cluster_ids is set: ", service_name) + return nil + end + + local match = ngx.re.match(service_name, "^(.*):(.*)$", "jo") + if not match then + core.log.error("get unexpected upstream service_name: ", service_name) + return nil + end + local endpoint_key, endpoint_port = match[1], match[2] + + local endpoint_dicts = core.table.new(#cluster_ids, 0) + local versions = core.table.new(#cluster_ids, 0) + for _, id in ipairs(cluster_ids) do + local endpoint_dict = ctx[id] + if not endpoint_dict then + core.log.error("kubernetes discovery cluster id not exist: ", id, + ", service: ", service_name) + else + local endpoint_version = endpoint_dict:get(endpoint_key .. "#version") + if endpoint_version then + core.table.insert(endpoint_dicts, endpoint_dict) + core.table.insert(versions, id .. "#" .. endpoint_version) + end + end + end + + if #endpoint_dicts == 0 then + core.log.info("get empty endpoint version from selected clusters for ", service_name) + return nil + end + + return endpoint_lrucache(service_name .. "#" .. core.table.concat(cluster_ids, ","), + core.table.concat(versions, ","), + merge_cluster_nodes, endpoint_dicts, endpoint_key, endpoint_port) +end + + +local function multiple_mode_nodes(service_name, discovery_args) + local cluster_ids = discovery_args and discovery_args.cluster_ids + if cluster_ids then + return selected_clusters_nodes(service_name, cluster_ids) + end + return k8s_core.resolve_nodes( endpoint_lrucache, service_name, "^(.*)/(.*/.*):(.*)$", -- id/namespace/name:port_name diff --git a/apisix/schema_def.lua b/apisix/schema_def.lua index 64f661f82142..8b008fe67b4c 100644 --- a/apisix/schema_def.lua +++ b/apisix/schema_def.lua @@ -570,6 +570,15 @@ local upstream_schema = { description = "group name", type = "string", }, + cluster_ids = { + description = "ids of the kubernetes discovery clusters to get nodes from", + type = "array", + minItems = 1, + uniqueItems = true, + items = { + type = "string", + }, + }, } }, pass_host = { diff --git a/apisix/upstream.lua b/apisix/upstream.lua index 76865fc06495..12ecce0a827f 100644 --- a/apisix/upstream.lua +++ b/apisix/upstream.lua @@ -28,6 +28,7 @@ local ipairs = ipairs local pairs = pairs local pcall = pcall local str_byte = string.byte +local str_find = string.find local ngx_var = ngx.var local is_http = ngx.config.subsystem == "http" local upstreams @@ -613,6 +614,26 @@ end _M.check_warm_up_conf = check_warm_up_conf +-- With `discovery_args.cluster_ids`, kubernetes discovery takes the clusters +-- from that list, so `service_name` is `namespace/name:port_name` and a +-- `cluster_id/namespace/name:port_name` form would be ambiguous. +local function check_discovery_args(conf) + local discovery_args = conf.discovery_args + if conf.discovery_type ~= "kubernetes" or not discovery_args + or not discovery_args.cluster_ids + then + return true + end + + if conf.service_name and str_find(conf.service_name, "/[^/]*/") then + return false, "service_name must not carry a cluster id prefix " .. + "when discovery_args.cluster_ids is set" + end + + return true +end + + local function check_upstream_conf(in_dp, conf) if not in_dp then local ok, err = check_schema(conf) @@ -625,6 +646,11 @@ local function check_upstream_conf(in_dp, conf) return false, err end + local ok, err = check_discovery_args(conf) + if not ok then + return false, err + end + if conf.nodes and not core.table.isarray(conf.nodes) then local port for addr,_ in pairs(conf.nodes) do diff --git a/docs/en/latest/discovery/kubernetes.md b/docs/en/latest/discovery/kubernetes.md index a2f137c39635..399fabf3aa3d 100644 --- a/docs/en/latest/discovery/kubernetes.md +++ b/docs/en/latest/discovery/kubernetes.md @@ -284,6 +284,35 @@ a nodes("release/default/plat-dev:port") call will get follow result: } ``` +### Select Clusters with `cluster_ids` + +In multi-cluster mode, an upstream can get its nodes from several clusters at once by listing their `id`s in `discovery_args.cluster_ids`. The `service_name` then uses the single-cluster pattern _[namespace]/[name]:[portName]_, without the `id` prefix: + +```json +{ + "type": "roundrobin", + "discovery_type": "kubernetes", + "service_name": "default/plat-dev:port", + "discovery_args": { + "cluster_ids": ["release", "staging"] + } +} +``` + +The upstream nodes are the union of the matching endpoints in the listed clusters, and they follow the endpoint changes in those clusters. When the same `host:port` is found in more than one listed cluster, it is used once. + ++ `cluster_ids` is a non-empty array of unique strings. It only applies to multi-cluster mode. In single-cluster mode, an upstream with `cluster_ids` gets no nodes and an error is logged. + ++ Clusters that are not listed, including clusters added to the configuration later, never contribute nodes. + ++ When `cluster_ids` is set, `service_name` must not carry an `id` prefix. The Admin API rejects such a configuration. + ++ An `id` that is not defined in the service discovery configuration is skipped, and an error is logged when the upstream is used. The Admin API does not check `cluster_ids` against the service discovery configuration. + ++ When none of the listed clusters has matching endpoints, the upstream has no valid nodes. Nodes from other clusters are never used instead. + +Without `cluster_ids`, the query interface described above is unchanged. + ## Q&A **Q: Why only support configuration token to access _Kubernetes APIServer_?** diff --git a/docs/zh/latest/discovery/kubernetes.md b/docs/zh/latest/discovery/kubernetes.md index b769400c3d5c..8cf59460135e 100644 --- a/docs/zh/latest/discovery/kubernetes.md +++ b/docs/zh/latest/discovery/kubernetes.md @@ -282,6 +282,35 @@ nodes("release/default/plat-dev:port") 调用会得到如下的返回值: } ``` +### 使用 `cluster_ids` 选择集群 + +在多集群模式下,上游可以通过 `discovery_args.cluster_ids` 列出多个集群的 `id`,同时从这些集群获取节点。此时 `service_name` 使用单集群模式的格式 _[namespace]/[name]:[portName]_,不带 `id` 前缀: + +```json +{ + "type": "roundrobin", + "discovery_type": "kubernetes", + "service_name": "default/plat-dev:port", + "discovery_args": { + "cluster_ids": ["release", "staging"] + } +} +``` + +上游节点是所列集群中匹配 endpoints 的并集,并随这些集群中 endpoints 的变化而更新。同一个 `host:port` 出现在多个所列集群中时只使用一次。 + ++ `cluster_ids` 是非空且元素不重复的字符串数组,仅适用于多集群模式。在单集群模式下,设置了 `cluster_ids` 的上游不会获得任何节点,并会记录错误日志。 + ++ 未列出的集群(包括之后新增到配置中的集群)不会提供节点。 + ++ 设置 `cluster_ids` 时,`service_name` 不能带 `id` 前缀,Admin API 会拒绝这样的配置。 + ++ 服务发现配置中不存在的 `id` 会被跳过,并在使用该上游时记录错误日志。Admin API 不会根据服务发现配置校验 `cluster_ids`。 + ++ 所列集群都没有匹配的 endpoints 时,上游没有可用节点,不会改用其他集群的节点。 + +不设置 `cluster_ids` 时,上述查询接口的行为不变。 + ## Q&A **Q: 为什么只支持配置 token 来访问 Kubernetes APIServer?** diff --git a/t/kubernetes/discovery/kubernetes5.t b/t/kubernetes/discovery/kubernetes5.t new file mode 100644 index 000000000000..941bd60b12de --- /dev/null +++ b/t/kubernetes/discovery/kubernetes5.t @@ -0,0 +1,517 @@ +# +# 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. +# + +BEGIN { + our $token_file = "/tmp/var/run/secrets/kubernetes.io/serviceaccount/token"; + our $token_value = eval {`cat $token_file 2>/dev/null`}; + + our $yaml_config = <<_EOC_; +apisix: + node_listen: 1984 +deployment: + role: data_plane + role_data_plane: + config_provider: yaml +discovery: + kubernetes: + - id: first + service: + host: "127.0.0.1" + port: "6443" + ssl_verify: false + client: + token_file: "/tmp/var/run/secrets/kubernetes.io/serviceaccount/token" + namespace_selector: + equal: ns-a + - id: second + service: + schema: "http" + host: "127.0.0.1" + port: "6445" + client: + token_file: "/tmp/var/run/secrets/kubernetes.io/serviceaccount/token" + +_EOC_ + + # no apiserver listens on 6446, so the tests below own the endpoint dicts + our $offline_yaml_config = <<_EOC_; +apisix: + node_listen: 1984 +deployment: + role: data_plane + role_data_plane: + config_provider: yaml +discovery: + kubernetes: + - id: first + service: + schema: "http" + host: "127.0.0.1" + port: "6446" + client: + token: "fake" + - id: second + service: + schema: "http" + host: "127.0.0.1" + port: "6446" + client: + token: "fake" + +_EOC_ + +} + +use t::APISIX 'no_plan'; + +repeat_each(1); +log_level('warn'); +no_root_location(); +no_shuffle(); +workers(4); + +add_block_preprocessor(sub { + my ($block) = @_; + + my $apisix_yaml = $block->apisix_yaml // <<_EOC_; +routes: [] +#END +_EOC_ + + $block->set_value("apisix_yaml", $apisix_yaml); + + my $main_config = $block->main_config // <<_EOC_; +env KUBERNETES_SERVICE_HOST=127.0.0.1; +env KUBERNETES_SERVICE_PORT=6443; +env KUBERNETES_CLIENT_TOKEN=$::token_value; +env KUBERNETES_CLIENT_TOKEN_FILE=$::token_file; +_EOC_ + + $block->set_value("main_config", $main_config); + + my $config = $block->config // <<_EOC_; + location /queries { + content_by_lua_block { + local core = require("apisix.core") + local d = require("apisix.discovery.kubernetes") + + ngx.sleep(1) + + ngx.req.read_body() + local request_body = ngx.req.get_body_data() + local queries = core.json.decode(request_body) + local response_body = "{" + for _, query in ipairs(queries) do + local args + if query.cluster_ids then + args = {cluster_ids = query.cluster_ids} + end + local nodes = d.nodes(query.service_name, args) + if nodes == nil or #nodes == 0 then + response_body = response_body .. " " .. 0 + else + response_body = response_body .. " " .. #nodes + end + end + ngx.say(response_body .. " }") + } + } + + location /operators { + content_by_lua_block { + local http = require("resty.http") + local core = require("apisix.core") + local ipairs = ipairs + + ngx.req.read_body() + local request_body = ngx.req.get_body_data() + local operators = core.json.decode(request_body) + + for _, op in ipairs(operators) do + local path = "/api/v1/namespaces/" .. op.namespace .. "/endpoints/" .. op.name + local body + if #op.subsets == 0 then + body = '[{"path":"/subsets","op":"replace","value":[]}]' + else + local t = { { op = "replace", path = "/subsets", value = op.subsets } } + body = core.json.encode(t, true) + end + + local httpc = http.new() + local res, err = httpc:request_uri("http://127.0.0.1:6445" .. path, { + method = "PATCH", + headers = { + ["Host"] = "127.0.0.1:6445", + ["Content-Type"] = "application/json-patch+json", + }, + body = body, + }) + if not res then + core.log.error("operator k8s cluster error: ", err) + return 500 + end + if res.status ~= 200 and res.status ~= 201 then + return res.status + end + end + ngx.say("DONE") + } + } + +_EOC_ + + $block->set_value("config", $config); + +}); + +run_tests(); + +__DATA__ + +=== TEST 1: create endpoints +--- yaml_config eval: $::yaml_config +--- request +POST /operators +[ + { + "namespace": "ns-a", + "name": "ep", + "subsets": [ + { + "addresses": [{"ip": "10.0.1.1"}], + "ports": [{"name": "p1", "port": 5001}] + } + ] + }, + { + "namespace": "ns-b", + "name": "ep", + "subsets": [ + { + "addresses": [{"ip": "10.0.2.1"}, {"ip": "10.0.2.2"}], + "ports": [{"name": "p1", "port": 5001}] + } + ] + } +] +--- more_headers +Content-type: application/json +--- response_body +DONE + + + +=== TEST 2: select clusters by cluster_ids +--- yaml_config eval: $::yaml_config +--- request +GET /queries +[ + {"service_name": "ns-a/ep:p1", "cluster_ids": ["first"]}, + {"service_name": "ns-a/ep:p1", "cluster_ids": ["second"]}, + {"service_name": "ns-a/ep:p1", "cluster_ids": ["first", "second"]}, + {"service_name": "ns-b/ep:p1", "cluster_ids": ["first"]}, + {"service_name": "ns-b/ep:p1", "cluster_ids": ["second"]}, + {"service_name": "ns-b/ep:p1", "cluster_ids": ["first", "second"]}, + {"service_name": "ns-b/ep:p1", "cluster_ids": ["third", "second"]}, + {"service_name": "ns-b/ep:p1", "cluster_ids": ["third"]}, + {"service_name": "first/ns-a/ep:p1"}, + {"service_name": "first/ns-b/ep:p1"}, + {"service_name": "second/ns-b/ep:p1"} +] +--- more_headers +Content-type: application/json +--- response_body +{ 1 1 1 0 2 2 2 0 1 0 2 } +--- error_log +kubernetes discovery cluster id not exist: third + + + +=== TEST 3: endpoint changes in the selected clusters keep refreshing +--- yaml_config eval: $::yaml_config +--- request eval +[ + +"GET /queries +[ + {\"service_name\": \"ns-b/ep:p1\", \"cluster_ids\": [\"first\", \"second\"]} +]", + +"POST /operators +[{\"name\":\"ep\",\"namespace\":\"ns-b\",\"subsets\":[{\"addresses\":[{\"ip\":\"10.0.2.1\"},{\"ip\":\"10.0.2.2\"},{\"ip\":\"10.0.2.3\"}],\"ports\":[{\"name\":\"p1\",\"port\":5001}]}]}]", + +"GET /queries +[ + {\"service_name\": \"ns-b/ep:p1\", \"cluster_ids\": [\"first\", \"second\"]} +]", + +"POST /operators +[{\"name\":\"ep\",\"namespace\":\"ns-b\",\"subsets\":[]}]", + +"GET /queries +[ + {\"service_name\": \"ns-b/ep:p1\", \"cluster_ids\": [\"first\", \"second\"]} +]" + +] +--- response_body eval +[ + "{ 2 }\n", + "DONE\n", + "{ 3 }\n", + "DONE\n", + "{ 0 }\n", +] + + + +=== TEST 4: union of the selected clusters +--- yaml_config eval: $::offline_yaml_config +--- config + location /t { + content_by_lua_block { + local core = require("apisix.core") + local function set_endpoints(id, key, endpoints, version) + local dict = ngx.shared["kubernetes-" .. id] + dict:set(key, core.json.encode(endpoints)) + dict:set(key .. "#version", version) + end + local d = require("apisix.discovery.kubernetes") + + local function show(service_name, cluster_ids) + local nodes = d.nodes(service_name, {cluster_ids = cluster_ids}) + if not nodes then + ngx.say("nil") + return + end + local addrs = {} + for _, node in ipairs(nodes) do + table.insert(addrs, node.host .. ":" .. node.port) + end + table.sort(addrs) + ngx.say(table.concat(addrs, " ")) + end + + set_endpoints("first", "ns/svc", {p1 = { + {host = "10.0.0.1", port = 80, weight = 50}, + {host = "10.0.0.9", port = 80, weight = 50}, + }}, "1") + set_endpoints("second", "ns/svc", {p1 = { + {host = "10.0.0.2", port = 80, weight = 50}, + {host = "10.0.0.9", port = 80, weight = 50}, + }}, "1") + set_endpoints("first", "ns/only-first", {p1 = { + {host = "10.0.1.1", port = 80, weight = 50}, + }}, "1") + + show("ns/svc:p1", {"first", "second"}) + show("ns/svc:p1", {"first"}) + show("ns/svc:p1", {"second"}) + show("ns/only-first:p1", {"second"}) + show("ns/only-first:p1", {"second", "first"}) + show("second/ns/svc:p1", {"second"}) + + set_endpoints("second", "ns/svc", {p1 = { + {host = "10.0.0.2", port = 80, weight = 50}, + {host = "10.0.0.3", port = 80, weight = 50}, + }}, "2") + show("ns/svc:p1", {"first", "second"}) + } + } +--- request +GET /t +--- response_body +10.0.0.1:80 10.0.0.2:80 10.0.0.9:80 +10.0.0.1:80 10.0.0.9:80 +10.0.0.2:80 10.0.0.9:80 +nil +10.0.1.1:80 +nil +10.0.0.1:80 10.0.0.2:80 10.0.0.3:80 10.0.0.9:80 +--- error_log +service_name must not carry a cluster id prefix when discovery_args.cluster_ids is set: second/ns/svc:p1 +--- no_error_log +cluster id not exist + + + +=== TEST 5: proxy to the selected clusters only +--- yaml_config eval: $::offline_yaml_config +--- apisix_yaml +routes: + - + uri: /hello + upstream: + service_name: ns/svc:p1 + discovery_type: kubernetes + discovery_args: + cluster_ids: + - second + type: roundrobin + - + uri: /hello1 + upstream: + service_name: ns/svc:p1 + discovery_type: kubernetes + discovery_args: + cluster_ids: + - third + type: roundrobin + - + uri: /hello_chunked + upstream: + service_name: ns/only-second:p1 + discovery_type: kubernetes + discovery_args: + cluster_ids: + - first + type: roundrobin +#END +--- config + location /t { + content_by_lua_block { + local core = require("apisix.core") + local function set_endpoints(id, key, endpoints, version) + local dict = ngx.shared["kubernetes-" .. id] + dict:set(key, core.json.encode(endpoints)) + dict:set(key .. "#version", version) + end + local http = require("resty.http") + + set_endpoints("first", "ns/svc", {p1 = { + {host = "127.0.0.1", port = 1979, weight = 50}, + }}, "1") + set_endpoints("second", "ns/svc", {p1 = { + {host = "127.0.0.1", port = 1980, weight = 50}, + }}, "1") + set_endpoints("second", "ns/only-second", {p1 = { + {host = "127.0.0.1", port = 1980, weight = 50}, + }}, "1") + + local uri = "http://127.0.0.1:" .. ngx.var.server_port + for _, path in ipairs({"/hello", "/hello", "/hello", "/hello1", "/hello_chunked"}) do + local httpc = http.new() + local res, err = httpc:request_uri(uri .. path) + if not res then + ngx.say(err) + return + end + ngx.say(path, " ", res.status) + end + } + } +--- request +GET /t +--- response_body +/hello 200 +/hello 200 +/hello 200 +/hello1 503 +/hello_chunked 503 +--- error_log +kubernetes discovery cluster id not exist: third + + + +=== TEST 6: cluster_ids requires multiple clusters +--- yaml_config +apisix: + node_listen: 1984 +deployment: + role: data_plane + role_data_plane: + config_provider: yaml +discovery: + kubernetes: + service: + schema: "http" + host: "127.0.0.1" + port: "6446" + client: + token: "fake" +--- config + location /t { + content_by_lua_block { + local core = require("apisix.core") + local d = require("apisix.discovery.kubernetes") + + local dict = ngx.shared["kubernetes"] + dict:set("ns/svc", core.json.encode({p1 = { + {host = "10.0.0.1", port = 80, weight = 50}, + }})) + dict:set("ns/svc#version", "1") + + local nodes = d.nodes("ns/svc:p1") + ngx.say(#nodes) + nodes = d.nodes("ns/svc:p1", {cluster_ids = {"first"}}) + ngx.say(tostring(nodes)) + } + } +--- request +GET /t +--- response_body +1 +nil +--- error_log +discovery_args.cluster_ids requires kubernetes discovery configured with multiple clusters + + + +=== TEST 7: validate cluster_ids +--- yaml_config +apisix: + node_listen: 1984 +deployment: + role: data_plane + role_data_plane: + config_provider: yaml +--- config + location /t { + content_by_lua_block { + local upstream = require("apisix.upstream") + local cases = { + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first", "second"}}}, + {service_name = "first/ns/svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first", "first"}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {1}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = "first"}}, + {service_name = "first/ns/svc:p1"}, + {service_name = "first/ns/svc:p1", discovery_type = "nacos", + discovery_args = {cluster_ids = {"first"}}}, + } + for _, case in ipairs(cases) do + case.discovery_type = case.discovery_type or "kubernetes" + case.type = "roundrobin" + local ok, err = upstream.check_upstream_conf(case) + ngx.say(ok and "passed" or err) + end + } + } +--- request +GET /t +--- response_body +passed +passed +service_name must not carry a cluster id prefix when discovery_args.cluster_ids is set +invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: expect array to have at least 1 items +invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: expected unique items but items 1 and 2 are equal +invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: failed to validate item 1: wrong type: expected string, got number +invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: wrong type: expected array, got string +passed +passed From 8cc8b16f78b16fb3b2047615c68f9d41006c46b3 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 29 Sep 2026 16:47:54 +0800 Subject: [PATCH 2/4] feat(kubernetes): validate discovery_args.cluster_ids in the Admin API --- apisix/discovery/kubernetes/init.lua | 73 +++++++++- apisix/upstream.lua | 35 ++--- docs/en/latest/discovery/kubernetes.md | 6 +- docs/zh/latest/discovery/kubernetes.md | 6 +- t/discovery/kubernetes_cluster_ids.t | 176 +++++++++++++++++++++++++ t/kubernetes/discovery/kubernetes5.t | 149 +++++++++++++++------ 6 files changed, 368 insertions(+), 77 deletions(-) create mode 100644 t/discovery/kubernetes_cluster_ids.t diff --git a/apisix/discovery/kubernetes/init.lua b/apisix/discovery/kubernetes/init.lua index dccc3aed15ef..2eabeed6763d 100644 --- a/apisix/discovery/kubernetes/init.lua +++ b/apisix/discovery/kubernetes/init.lua @@ -181,11 +181,20 @@ local function merge_cluster_nodes(endpoint_dicts, endpoint_key, endpoint_port) end --- service_name is "namespace/name:port_name", the clusters come from cluster_ids -local function selected_clusters_nodes(service_name, cluster_ids) +-- with cluster_ids, service_name is "namespace/name:port_name" +local function check_cluster_ids_service_name(service_name) if str_find(service_name, "/[^/]*/") then - core.log.error("service_name must not carry a cluster id prefix ", - "when discovery_args.cluster_ids is set: ", service_name) + return false, "service_name must be namespace/name:port_name when " + .. "discovery_args.cluster_ids is set, got: " .. service_name + end + return true +end + + +local function selected_clusters_nodes(service_name, cluster_ids) + local ok, err = check_cluster_ids_service_name(service_name) + if not ok then + core.log.error(err) return nil end @@ -198,11 +207,12 @@ local function selected_clusters_nodes(service_name, cluster_ids) local endpoint_dicts = core.table.new(#cluster_ids, 0) local versions = core.table.new(#cluster_ids, 0) + local unknown_ids for _, id in ipairs(cluster_ids) do local endpoint_dict = ctx[id] if not endpoint_dict then - core.log.error("kubernetes discovery cluster id not exist: ", id, - ", service: ", service_name) + unknown_ids = unknown_ids or {} + core.table.insert(unknown_ids, id) else local endpoint_version = endpoint_dict:get(endpoint_key .. "#version") if endpoint_version then @@ -212,6 +222,11 @@ local function selected_clusters_nodes(service_name, cluster_ids) end end + if unknown_ids then + core.log.warn("skip unknown kubernetes discovery cluster ids: ", + core.table.concat(unknown_ids, ", "), ", service: ", service_name) + end + if #endpoint_dicts == 0 then core.log.info("get empty endpoint version from selected clusters for ", service_name) return nil @@ -257,6 +272,52 @@ function _M.init_worker() end +function _M.check_discovery_args(discovery_args, service_name, in_dp) + local cluster_ids = discovery_args and discovery_args.cluster_ids + if not cluster_ids then + return true + end + + if service_name then + local ok, err = check_cluster_ids_service_name(service_name) + if not ok then + return false, err + end + end + + -- the data plane resolves the ids at runtime, where an unknown id is skipped + if in_dp then + return true + end + + local discovery_conf = local_conf.discovery.kubernetes + if #discovery_conf == 0 then + return false, "discovery_args.cluster_ids requires kubernetes discovery " + .. "configured with multiple clusters" + end + + local known_ids = {} + for _, conf in ipairs(discovery_conf) do + known_ids[conf.id] = true + end + + local unknown_ids + for _, id in ipairs(cluster_ids) do + if not known_ids[id] then + unknown_ids = unknown_ids or {} + core.table.insert(unknown_ids, id) + end + end + + if unknown_ids then + return false, "unknown kubernetes discovery cluster ids in " + .. "discovery_args.cluster_ids: " .. core.table.concat(unknown_ids, ", ") + end + + return true +end + + function _M.dump_data() local discovery_conf = local_conf.discovery.kubernetes local eps = {} diff --git a/apisix/upstream.lua b/apisix/upstream.lua index 12ecce0a827f..9ccfeaaa66b3 100644 --- a/apisix/upstream.lua +++ b/apisix/upstream.lua @@ -28,7 +28,6 @@ local ipairs = ipairs local pairs = pairs local pcall = pcall local str_byte = string.byte -local str_find = string.find local ngx_var = ngx.var local is_http = ngx.config.subsystem == "http" local upstreams @@ -614,26 +613,6 @@ end _M.check_warm_up_conf = check_warm_up_conf --- With `discovery_args.cluster_ids`, kubernetes discovery takes the clusters --- from that list, so `service_name` is `namespace/name:port_name` and a --- `cluster_id/namespace/name:port_name` form would be ambiguous. -local function check_discovery_args(conf) - local discovery_args = conf.discovery_args - if conf.discovery_type ~= "kubernetes" or not discovery_args - or not discovery_args.cluster_ids - then - return true - end - - if conf.service_name and str_find(conf.service_name, "/[^/]*/") then - return false, "service_name must not carry a cluster id prefix " .. - "when discovery_args.cluster_ids is set" - end - - return true -end - - local function check_upstream_conf(in_dp, conf) if not in_dp then local ok, err = check_schema(conf) @@ -646,11 +625,6 @@ local function check_upstream_conf(in_dp, conf) return false, err end - local ok, err = check_discovery_args(conf) - if not ok then - return false, err - end - if conf.nodes and not core.table.isarray(conf.nodes) then local port for addr,_ in pairs(conf.nodes) do @@ -693,6 +667,15 @@ local function check_upstream_conf(in_dp, conf) end end + -- a discovery module may check discovery_args against its own configuration + local dis = conf.discovery_type and discovery and discovery[conf.discovery_type] + if dis and dis.check_discovery_args then + local ok, err = dis.check_discovery_args(conf.discovery_args, conf.service_name, in_dp) + if not ok then + return false, err + end + end + if is_http then if conf.pass_host == "rewrite" and (conf.upstream_host == nil or conf.upstream_host == "") diff --git a/docs/en/latest/discovery/kubernetes.md b/docs/en/latest/discovery/kubernetes.md index 399fabf3aa3d..bd5da1971e3b 100644 --- a/docs/en/latest/discovery/kubernetes.md +++ b/docs/en/latest/discovery/kubernetes.md @@ -301,13 +301,13 @@ In multi-cluster mode, an upstream can get its nodes from several clusters at on The upstream nodes are the union of the matching endpoints in the listed clusters, and they follow the endpoint changes in those clusters. When the same `host:port` is found in more than one listed cluster, it is used once. -+ `cluster_ids` is a non-empty array of unique strings. It only applies to multi-cluster mode. In single-cluster mode, an upstream with `cluster_ids` gets no nodes and an error is logged. ++ `cluster_ids` is a non-empty array of unique strings. It only applies to multi-cluster mode. The Admin API rejects `cluster_ids` when Kubernetes service discovery is configured in single-cluster mode. If such an upstream reaches the data plane in another way, it gets no nodes and an error is logged. + Clusters that are not listed, including clusters added to the configuration later, never contribute nodes. -+ When `cluster_ids` is set, `service_name` must not carry an `id` prefix. The Admin API rejects such a configuration. ++ When `cluster_ids` is set, `service_name` must not carry an `id` prefix. Such a configuration is rejected, and the error shows the `service_name` value. -+ An `id` that is not defined in the service discovery configuration is skipped, and an error is logged when the upstream is used. The Admin API does not check `cluster_ids` against the service discovery configuration. ++ The Admin API rejects an `id` that is not defined in the Kubernetes service discovery configuration of the APISIX instance that serves the Admin API, and the error lists all unknown `id`s. When a configuration with an unknown `id` still reaches the data plane, for example because it was written before a cluster was removed from the configuration, or because it comes from a standalone configuration file, the unknown `id`s are skipped and a warning that lists them is logged. + When none of the listed clusters has matching endpoints, the upstream has no valid nodes. Nodes from other clusters are never used instead. diff --git a/docs/zh/latest/discovery/kubernetes.md b/docs/zh/latest/discovery/kubernetes.md index 8cf59460135e..11a2216a0607 100644 --- a/docs/zh/latest/discovery/kubernetes.md +++ b/docs/zh/latest/discovery/kubernetes.md @@ -299,13 +299,13 @@ nodes("release/default/plat-dev:port") 调用会得到如下的返回值: 上游节点是所列集群中匹配 endpoints 的并集,并随这些集群中 endpoints 的变化而更新。同一个 `host:port` 出现在多个所列集群中时只使用一次。 -+ `cluster_ids` 是非空且元素不重复的字符串数组,仅适用于多集群模式。在单集群模式下,设置了 `cluster_ids` 的上游不会获得任何节点,并会记录错误日志。 ++ `cluster_ids` 是非空且元素不重复的字符串数组,仅适用于多集群模式。Kubernetes 服务发现为单集群模式时,Admin API 会拒绝 `cluster_ids`。如果这样的上游通过其他方式到达数据面,它不会获得任何节点,并会记录错误日志。 + 未列出的集群(包括之后新增到配置中的集群)不会提供节点。 -+ 设置 `cluster_ids` 时,`service_name` 不能带 `id` 前缀,Admin API 会拒绝这样的配置。 ++ 设置 `cluster_ids` 时,`service_name` 不能带 `id` 前缀。这样的配置会被拒绝,错误信息中会给出 `service_name` 的值。 -+ 服务发现配置中不存在的 `id` 会被跳过,并在使用该上游时记录错误日志。Admin API 不会根据服务发现配置校验 `cluster_ids`。 ++ 如果 `id` 不在提供 Admin API 的 APISIX 实例的 Kubernetes 服务发现配置中,Admin API 会拒绝该配置,错误信息会列出所有未知的 `id`。如果带有未知 `id` 的配置仍然到达了数据面,例如该配置写入于某个集群从配置中移除之前,或者来自 standalone 配置文件,未知的 `id` 会被跳过,并记录一条列出这些 `id` 的警告日志。 + 所列集群都没有匹配的 endpoints 时,上游没有可用节点,不会改用其他集群的节点。 diff --git a/t/discovery/kubernetes_cluster_ids.t b/t/discovery/kubernetes_cluster_ids.t new file mode 100644 index 000000000000..e4cd7648f874 --- /dev/null +++ b/t/discovery/kubernetes_cluster_ids.t @@ -0,0 +1,176 @@ +# +# 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. +# +use t::APISIX 'no_plan'; + +repeat_each(1); +no_long_string(); +no_shuffle(); +no_root_location(); +log_level("warn"); + +add_block_preprocessor(sub { + my ($block) = @_; + + if (!$block->request) { + $block->set_value("request", "GET /t"); + } + + # no apiserver listens on 6446, the tests only check the Admin API + my $extra_yaml_config = $block->extra_yaml_config // <<_EOC_; +discovery: + kubernetes: + - id: first + service: + schema: "http" + host: "127.0.0.1" + port: "6446" + client: + token: "fake" + - id: second + service: + schema: "http" + host: "127.0.0.1" + port: "6446" + client: + token: "fake" +_EOC_ + + $block->set_value("extra_yaml_config", $extra_yaml_config); + + if (!$block->error_log && !$block->no_error_log) { + $block->set_value("no_error_log", "[alert]"); + } +}); + +run_tests; + +__DATA__ + +=== TEST 1: validate cluster_ids in the Admin API +--- config + location /t { + content_by_lua_block { + local t = require("lib.test_admin").test + local cases = { + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"second", "first"}}}, + {service_name = "ns/svc:p1", + discovery_args = {cluster_ids = {"third", "first", "fourth"}}}, + {service_name = "first/ns/svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first", "first"}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {1}}}, + {service_name = "ns/svc:p1", discovery_args = {cluster_ids = "first"}}, + {service_name = "first/ns/svc:p1"}, + } + for _, case in ipairs(cases) do + case.discovery_type = "kubernetes" + case.type = "roundrobin" + local code, body = t('/apisix/admin/upstreams/1', ngx.HTTP_PUT, case) + if code >= 300 then + ngx.say(code, " ", (body:gsub("%s+$", ""))) + else + ngx.say("passed") + end + end + } + } +--- response_body +passed +passed +400 {"error_msg":"unknown kubernetes discovery cluster ids in discovery_args.cluster_ids: third, fourth"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: first/ns/svc:p1"} +400 {"error_msg":"invalid configuration: property \"discovery_args\" validation failed: property \"cluster_ids\" validation failed: expect array to have at least 1 items"} +400 {"error_msg":"invalid configuration: property \"discovery_args\" validation failed: property \"cluster_ids\" validation failed: expected unique items but items 1 and 2 are equal"} +400 {"error_msg":"invalid configuration: property \"discovery_args\" validation failed: property \"cluster_ids\" validation failed: failed to validate item 1: wrong type: expected string, got number"} +400 {"error_msg":"invalid configuration: property \"discovery_args\" validation failed: property \"cluster_ids\" validation failed: wrong type: expected array, got string"} +passed + + + +=== TEST 2: validate cluster_ids in an inline upstream +--- config + location /t { + content_by_lua_block { + local t = require("lib.test_admin").test + for _, uri in ipairs({"/apisix/admin/routes/1", "/apisix/admin/services/1"}) do + local conf = { + upstream = { + type = "roundrobin", + discovery_type = "kubernetes", + service_name = "ns/svc:p1", + discovery_args = {cluster_ids = {"first", "third"}}, + }, + } + if uri == "/apisix/admin/routes/1" then + conf.uri = "/hello" + end + local code, body = t(uri, ngx.HTTP_PUT, conf) + ngx.say(code, " ", (body:gsub("%s+$", ""))) + end + } + } +--- response_body +400 {"error_msg":"unknown kubernetes discovery cluster ids in discovery_args.cluster_ids: third"} +400 {"error_msg":"unknown kubernetes discovery cluster ids in discovery_args.cluster_ids: third"} + + + +=== TEST 3: cluster_ids of another discovery type is not checked by kubernetes +--- config + location /t { + content_by_lua_block { + local t = require("lib.test_admin").test + local code, body = t('/apisix/admin/upstreams/1', ngx.HTTP_PUT, { + type = "roundrobin", + discovery_type = "consul", + service_name = "first/ns/svc:p1", + discovery_args = {cluster_ids = {"third"}}, + }) + ngx.say(code < 300 and "passed" or body) + } + } +--- response_body +passed + + + +=== TEST 4: cluster_ids requires multiple clusters +--- extra_yaml_config +discovery: + kubernetes: + service: + schema: "http" + host: "127.0.0.1" + port: "6446" + client: + token: "fake" +--- config + location /t { + content_by_lua_block { + local t = require("lib.test_admin").test + local code, body = t('/apisix/admin/upstreams/1', ngx.HTTP_PUT, { + type = "roundrobin", + discovery_type = "kubernetes", + service_name = "ns/svc:p1", + discovery_args = {cluster_ids = {"first"}}, + }) + ngx.say(code, " ", (body:gsub("%s+$", ""))) + } + } +--- response_body +400 {"error_msg":"discovery_args.cluster_ids requires kubernetes discovery configured with multiple clusters"} diff --git a/t/kubernetes/discovery/kubernetes5.t b/t/kubernetes/discovery/kubernetes5.t index 941bd60b12de..572494343b48 100644 --- a/t/kubernetes/discovery/kubernetes5.t +++ b/t/kubernetes/discovery/kubernetes5.t @@ -237,7 +237,7 @@ Content-type: application/json --- response_body { 1 1 1 0 2 2 2 0 1 0 2 } --- error_log -kubernetes discovery cluster id not exist: third +skip unknown kubernetes discovery cluster ids: third, service: ns-b/ep:p1 @@ -343,9 +343,9 @@ nil nil 10.0.0.1:80 10.0.0.2:80 10.0.0.3:80 10.0.0.9:80 --- error_log -service_name must not carry a cluster id prefix when discovery_args.cluster_ids is set: second/ns/svc:p1 +service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: second/ns/svc:p1 --- no_error_log -cluster id not exist +skip unknown kubernetes discovery cluster ids @@ -370,6 +370,7 @@ routes: discovery_args: cluster_ids: - third + - fourth type: roundrobin - uri: /hello_chunked @@ -423,7 +424,7 @@ GET /t /hello1 503 /hello_chunked 503 --- error_log -kubernetes discovery cluster id not exist: third +skip unknown kubernetes discovery cluster ids: third, fourth, service: ns/svc:p1 @@ -471,47 +472,117 @@ discovery_args.cluster_ids requires kubernetes discovery configured with multipl -=== TEST 7: validate cluster_ids ---- yaml_config -apisix: - node_listen: 1984 -deployment: - role: data_plane - role_data_plane: - config_provider: yaml +=== TEST 7: the data plane rejects a cluster id prefix but not an unknown id +--- yaml_config eval: $::offline_yaml_config +--- apisix_yaml +upstreams: + - + id: 1 + service_name: second/ns/svc:p1 + discovery_type: kubernetes + discovery_args: + cluster_ids: + - second + type: roundrobin + - + id: 2 + service_name: ns/svc:p1 + discovery_type: kubernetes + discovery_args: + cluster_ids: + - second + - third + type: roundrobin +routes: + - + uri: /hello + upstream_id: 1 + - + uri: /hello1 + upstream_id: 2 +#END --- config location /t { content_by_lua_block { - local upstream = require("apisix.upstream") - local cases = { - {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first"}}}, - {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first", "second"}}}, - {service_name = "first/ns/svc:p1", discovery_args = {cluster_ids = {"first"}}}, - {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {}}}, - {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {"first", "first"}}}, - {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {1}}}, - {service_name = "ns/svc:p1", discovery_args = {cluster_ids = "first"}}, - {service_name = "first/ns/svc:p1"}, - {service_name = "first/ns/svc:p1", discovery_type = "nacos", - discovery_args = {cluster_ids = {"first"}}}, - } - for _, case in ipairs(cases) do - case.discovery_type = case.discovery_type or "kubernetes" - case.type = "roundrobin" - local ok, err = upstream.check_upstream_conf(case) - ngx.say(ok and "passed" or err) + local core = require("apisix.core") + local http = require("resty.http") + + local dict = ngx.shared["kubernetes-second"] + dict:set("ns/svc", core.json.encode({p1 = { + {host = "127.0.0.1", port = 1980, weight = 50}, + }})) + dict:set("ns/svc#version", "1") + + local uri = "http://127.0.0.1:" .. ngx.var.server_port + for _, path in ipairs({"/hello", "/hello1"}) do + local httpc = http.new() + local res, err = httpc:request_uri(uri .. path) + if not res then + ngx.say(err) + return + end + ngx.say(path, " ", res.status) end } } --- request GET /t --- response_body -passed -passed -service_name must not carry a cluster id prefix when discovery_args.cluster_ids is set -invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: expect array to have at least 1 items -invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: expected unique items but items 1 and 2 are equal -invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: failed to validate item 1: wrong type: expected string, got number -invalid configuration: property "discovery_args" validation failed: property "cluster_ids" validation failed: wrong type: expected array, got string -passed -passed +/hello 502 +/hello1 200 +--- error_log +service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: second/ns/svc:p1 +skip unknown kubernetes discovery cluster ids: third, service: ns/svc:p1 + + + +=== TEST 8: stream route with cluster_ids +--- yaml_config eval: $::offline_yaml_config +--- apisix_yaml +stream_routes: + - + id: 1 + server_addr: 127.0.0.1 + server_port: 1985 + upstream: + service_name: ns/svc:p1 + discovery_type: kubernetes + discovery_args: + cluster_ids: + - second + type: roundrobin +#END +--- stream_extra_init_worker_by_lua + local core = require("apisix.core") + local d = require("apisix.discovery.kubernetes") + + local function set_endpoints(id, key, endpoints, version) + local dict = ngx.shared["kubernetes-" .. id .. "-stream"] + dict:set(key, core.json.encode(endpoints)) + dict:set(key .. "#version", version) + end + + set_endpoints("first", "ns/svc", {p1 = { + {host = "127.0.0.1", port = 1979, weight = 50}, + {host = "127.0.0.1", port = 1995, weight = 50}, + }}, "1") + set_endpoints("second", "ns/svc", {p1 = { + {host = "127.0.0.1", port = 1995, weight = 50}, + {host = "127.0.0.2", port = 1995, weight = 50}, + }}, "1") + + local nodes = d.nodes("ns/svc:p1", {cluster_ids = {"first", "second"}}) + local addrs = {} + for _, node in ipairs(nodes) do + table.insert(addrs, node.host .. ":" .. node.port) + end + table.sort(addrs) + core.log.warn("stream nodes of first and second: ", table.concat(addrs, " ")) +--- stream_request +m +--- stream_response +hello world +--- error_log +stream nodes of first and second: 127.0.0.1:1979 127.0.0.1:1995 127.0.0.2:1995 +--- no_error_log +skip unknown kubernetes discovery cluster ids From 200b5c91fe327f30151811c743b19a511a429099 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 29 Sep 2026 16:57:15 +0800 Subject: [PATCH 3/4] fix(kubernetes): validate the full service_name shape with cluster_ids --- apisix/discovery/kubernetes/init.lua | 29 ++++++++++++-------------- docs/en/latest/discovery/kubernetes.md | 2 +- docs/zh/latest/discovery/kubernetes.md | 2 +- t/discovery/kubernetes_cluster_ids.t | 10 +++++++++ 4 files changed, 25 insertions(+), 18 deletions(-) diff --git a/apisix/discovery/kubernetes/init.lua b/apisix/discovery/kubernetes/init.lua index 2eabeed6763d..0ad19fe514e6 100644 --- a/apisix/discovery/kubernetes/init.lua +++ b/apisix/discovery/kubernetes/init.lua @@ -20,7 +20,6 @@ local type = type local ipairs = ipairs local pairs = pairs local string = string -local str_find = string.find local error = error local tostring = tostring local is_http = ngx.config.subsystem == "http" @@ -182,25 +181,23 @@ end -- with cluster_ids, service_name is "namespace/name:port_name" -local function check_cluster_ids_service_name(service_name) - if str_find(service_name, "/[^/]*/") then - return false, "service_name must be namespace/name:port_name when " - .. "discovery_args.cluster_ids is set, got: " .. service_name +local cluster_ids_service_name_pattern = [[^([^/]+/[^/:]+):(.+)$]] + + +local function parse_cluster_ids_service_name(service_name) + local match = ngx.re.match(service_name, cluster_ids_service_name_pattern, "jo") + if not match then + return nil, "service_name must be namespace/name:port_name when " + .. "discovery_args.cluster_ids is set, got: " .. service_name end - return true + return match end local function selected_clusters_nodes(service_name, cluster_ids) - local ok, err = check_cluster_ids_service_name(service_name) - if not ok then - core.log.error(err) - return nil - end - - local match = ngx.re.match(service_name, "^(.*):(.*)$", "jo") + local match, err = parse_cluster_ids_service_name(service_name) if not match then - core.log.error("get unexpected upstream service_name: ", service_name) + core.log.error(err) return nil end local endpoint_key, endpoint_port = match[1], match[2] @@ -279,8 +276,8 @@ function _M.check_discovery_args(discovery_args, service_name, in_dp) end if service_name then - local ok, err = check_cluster_ids_service_name(service_name) - if not ok then + local match, err = parse_cluster_ids_service_name(service_name) + if not match then return false, err end end diff --git a/docs/en/latest/discovery/kubernetes.md b/docs/en/latest/discovery/kubernetes.md index bd5da1971e3b..b00e3750fcee 100644 --- a/docs/en/latest/discovery/kubernetes.md +++ b/docs/en/latest/discovery/kubernetes.md @@ -305,7 +305,7 @@ The upstream nodes are the union of the matching endpoints in the listed cluster + Clusters that are not listed, including clusters added to the configuration later, never contribute nodes. -+ When `cluster_ids` is set, `service_name` must not carry an `id` prefix. Such a configuration is rejected, and the error shows the `service_name` value. ++ When `cluster_ids` is set, `service_name` must match _[namespace]/[name]:[portName]_, so it cannot carry an `id` prefix. Any other value is rejected, and the error shows the `service_name` value. + The Admin API rejects an `id` that is not defined in the Kubernetes service discovery configuration of the APISIX instance that serves the Admin API, and the error lists all unknown `id`s. When a configuration with an unknown `id` still reaches the data plane, for example because it was written before a cluster was removed from the configuration, or because it comes from a standalone configuration file, the unknown `id`s are skipped and a warning that lists them is logged. diff --git a/docs/zh/latest/discovery/kubernetes.md b/docs/zh/latest/discovery/kubernetes.md index 11a2216a0607..6ec9993ed34e 100644 --- a/docs/zh/latest/discovery/kubernetes.md +++ b/docs/zh/latest/discovery/kubernetes.md @@ -303,7 +303,7 @@ nodes("release/default/plat-dev:port") 调用会得到如下的返回值: + 未列出的集群(包括之后新增到配置中的集群)不会提供节点。 -+ 设置 `cluster_ids` 时,`service_name` 不能带 `id` 前缀。这样的配置会被拒绝,错误信息中会给出 `service_name` 的值。 ++ 设置 `cluster_ids` 时,`service_name` 必须满足格式 _[namespace]/[name]:[portName]_,因此不能带 `id` 前缀。其他格式的值会被拒绝,错误信息中会给出 `service_name` 的值。 + 如果 `id` 不在提供 Admin API 的 APISIX 实例的 Kubernetes 服务发现配置中,Admin API 会拒绝该配置,错误信息会列出所有未知的 `id`。如果带有未知 `id` 的配置仍然到达了数据面,例如该配置写入于某个集群从配置中移除之前,或者来自 standalone 配置文件,未知的 `id` 会被跳过,并记录一条列出这些 `id` 的警告日志。 diff --git a/t/discovery/kubernetes_cluster_ids.t b/t/discovery/kubernetes_cluster_ids.t index e4cd7648f874..5122887c58c0 100644 --- a/t/discovery/kubernetes_cluster_ids.t +++ b/t/discovery/kubernetes_cluster_ids.t @@ -76,6 +76,11 @@ __DATA__ {service_name = "ns/svc:p1", discovery_args = {cluster_ids = {1}}}, {service_name = "ns/svc:p1", discovery_args = {cluster_ids = "first"}}, {service_name = "first/ns/svc:p1"}, + {service_name = "svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "/svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc:", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc", discovery_args = {cluster_ids = {"first"}}}, } for _, case in ipairs(cases) do case.discovery_type = "kubernetes" @@ -99,6 +104,11 @@ passed 400 {"error_msg":"invalid configuration: property \"discovery_args\" validation failed: property \"cluster_ids\" validation failed: failed to validate item 1: wrong type: expected string, got number"} 400 {"error_msg":"invalid configuration: property \"discovery_args\" validation failed: property \"cluster_ids\" validation failed: wrong type: expected array, got string"} passed +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: svc:p1"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: /svc:p1"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/:p1"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/svc:"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/svc"} From b2d7ddbc87a02afebf8cf69c5149b16a7cf36385 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 29 Sep 2026 17:13:56 +0800 Subject: [PATCH 4/4] fix(kubernetes): reject ':' and '/' in each cluster_ids service_name segment --- apisix/discovery/kubernetes/init.lua | 2 +- t/discovery/kubernetes_cluster_ids.t | 4 ++++ 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/apisix/discovery/kubernetes/init.lua b/apisix/discovery/kubernetes/init.lua index 0ad19fe514e6..11c8b063073c 100644 --- a/apisix/discovery/kubernetes/init.lua +++ b/apisix/discovery/kubernetes/init.lua @@ -181,7 +181,7 @@ end -- with cluster_ids, service_name is "namespace/name:port_name" -local cluster_ids_service_name_pattern = [[^([^/]+/[^/:]+):(.+)$]] +local cluster_ids_service_name_pattern = [[^([^/:]+/[^/:]+):([^/:]+)$]] local function parse_cluster_ids_service_name(service_name) diff --git a/t/discovery/kubernetes_cluster_ids.t b/t/discovery/kubernetes_cluster_ids.t index 5122887c58c0..9606d51b85c0 100644 --- a/t/discovery/kubernetes_cluster_ids.t +++ b/t/discovery/kubernetes_cluster_ids.t @@ -81,6 +81,8 @@ __DATA__ {service_name = "ns/:p1", discovery_args = {cluster_ids = {"first"}}}, {service_name = "ns/svc:", discovery_args = {cluster_ids = {"first"}}}, {service_name = "ns/svc", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns:extra/svc:p1", discovery_args = {cluster_ids = {"first"}}}, + {service_name = "ns/svc:p1:extra", discovery_args = {cluster_ids = {"first"}}}, } for _, case in ipairs(cases) do case.discovery_type = "kubernetes" @@ -109,6 +111,8 @@ passed 400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/:p1"} 400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/svc:"} 400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/svc"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns:extra/svc:p1"} +400 {"error_msg":"service_name must be namespace/name:port_name when discovery_args.cluster_ids is set, got: ns/svc:p1:extra"}