diff --git a/apisix/discovery/kubernetes/init.lua b/apisix/discovery/kubernetes/init.lua index f9207792356f..11c8b063073c 100644 --- a/apisix/discovery/kubernetes/init.lua +++ b/apisix/discovery/kubernetes/init.lua @@ -21,6 +21,7 @@ local ipairs = ipairs local pairs = pairs local string = string 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 +88,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 +161,86 @@ 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 + + +-- with cluster_ids, service_name is "namespace/name:port_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 match +end + + +local function selected_clusters_nodes(service_name, cluster_ids) + local match, err = parse_cluster_ids_service_name(service_name) + if not match then + core.log.error(err) + 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) + local unknown_ids + for _, id in ipairs(cluster_ids) do + local endpoint_dict = ctx[id] + if not endpoint_dict then + 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 + core.table.insert(endpoint_dicts, endpoint_dict) + core.table.insert(versions, id .. "#" .. endpoint_version) + end + 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 + 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 @@ -183,6 +269,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 match, err = parse_cluster_ids_service_name(service_name) + if not match 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/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..9ccfeaaa66b3 100644 --- a/apisix/upstream.lua +++ b/apisix/upstream.lua @@ -667,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 a2f137c39635..b00e3750fcee 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. 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 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. + ++ 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..6ec9993ed34e 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` 是非空且元素不重复的字符串数组,仅适用于多集群模式。Kubernetes 服务发现为单集群模式时,Admin API 会拒绝 `cluster_ids`。如果这样的上游通过其他方式到达数据面,它不会获得任何节点,并会记录错误日志。 + ++ 未列出的集群(包括之后新增到配置中的集群)不会提供节点。 + ++ 设置 `cluster_ids` 时,`service_name` 必须满足格式 _[namespace]/[name]:[portName]_,因此不能带 `id` 前缀。其他格式的值会被拒绝,错误信息中会给出 `service_name` 的值。 + ++ 如果 `id` 不在提供 Admin API 的 APISIX 实例的 Kubernetes 服务发现配置中,Admin API 会拒绝该配置,错误信息会列出所有未知的 `id`。如果带有未知 `id` 的配置仍然到达了数据面,例如该配置写入于某个集群从配置中移除之前,或者来自 standalone 配置文件,未知的 `id` 会被跳过,并记录一条列出这些 `id` 的警告日志。 + ++ 所列集群都没有匹配的 endpoints 时,上游没有可用节点,不会改用其他集群的节点。 + +不设置 `cluster_ids` 时,上述查询接口的行为不变。 + ## Q&A **Q: 为什么只支持配置 token 来访问 Kubernetes APIServer?** diff --git a/t/discovery/kubernetes_cluster_ids.t b/t/discovery/kubernetes_cluster_ids.t new file mode 100644 index 000000000000..9606d51b85c0 --- /dev/null +++ b/t/discovery/kubernetes_cluster_ids.t @@ -0,0 +1,190 @@ +# +# 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"}, + {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"}}}, + {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" + 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 +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"} +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"} + + + +=== 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 new file mode 100644 index 000000000000..572494343b48 --- /dev/null +++ b/t/kubernetes/discovery/kubernetes5.t @@ -0,0 +1,588 @@ +# +# 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 +skip unknown kubernetes discovery cluster ids: third, service: ns-b/ep:p1 + + + +=== 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 be namespace/name:port_name when discovery_args.cluster_ids is set, got: second/ns/svc:p1 +--- no_error_log +skip unknown kubernetes discovery cluster ids + + + +=== 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 + - fourth + 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 +skip unknown kubernetes discovery cluster ids: third, fourth, service: ns/svc:p1 + + + +=== 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: 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 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 +/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