Skip to content
Open
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
136 changes: 134 additions & 2 deletions apisix/discovery/kubernetes/init.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 = {}
Expand Down
9 changes: 9 additions & 0 deletions apisix/schema_def.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {
Expand Down
9 changes: 9 additions & 0 deletions apisix/upstream.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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 == "")
Expand Down
29 changes: 29 additions & 0 deletions docs/en/latest/discovery/kubernetes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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_?**
Expand Down
29 changes: 29 additions & 0 deletions docs/zh/latest/discovery/kubernetes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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?**
Expand Down
Loading
Loading