This is an automated email from the ASF dual-hosted git repository.
AlinsRan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/apisix.git
The following commit(s) were added to refs/heads/master by this push:
new c525ae5c8f feat(upstream): slow start for newly observed upstream
nodes (#13941)
c525ae5c8f is described below
commit c525ae5c8f01defeaaf0b10c979bc642e5c82dfc
Author: AlinsRan <[email protected]>
AuthorDate: Thu Sep 17 09:01:32 2026 +0800
feat(upstream): slow start for newly observed upstream nodes (#13941)
---
apisix/balancer.lua | 48 ++-
apisix/cli/config.lua | 1 +
apisix/cli/ngx_tpl.lua | 1 +
apisix/init.lua | 1 +
apisix/plugins/traffic-split.lua | 10 +
apisix/schema_def.lua | 43 ++
apisix/slow_start.lua | 531 +++++++++++++++++++++++++
apisix/upstream.lua | 57 +++
conf/config.yaml.example | 1 +
docs/en/latest/admin-api.md | 14 +
docs/zh/latest/admin-api.md | 11 +
t/APISIX.pm | 1 +
t/admin/upstream-slow-start.t | 622 +++++++++++++++++++++++++++++
t/node/upstream-slow-start.t | 833 +++++++++++++++++++++++++++++++++++++++
14 files changed, 2157 insertions(+), 17 deletions(-)
diff --git a/apisix/balancer.lua b/apisix/balancer.lua
index 35f015da1b..0258d4bb54 100644
--- a/apisix/balancer.lua
+++ b/apisix/balancer.lua
@@ -20,6 +20,7 @@ local core = require("apisix.core")
local priority_balancer = require("apisix.balancer.priority")
local apisix_upstream = require("apisix.upstream")
local healthcheck_manager = require("apisix.healthcheck_manager")
+local slow_start = require("apisix.slow_start")
local ipairs = ipairs
local is_http = ngx.config.subsystem == "http"
local enable_keepalive = balancer.enable_keepalive and is_http
@@ -50,7 +51,7 @@ local _M = {
}
-local function transform_node(new_nodes, node)
+local function transform_node(new_nodes, node, weight)
if not new_nodes._priority_index then
new_nodes._priority_index = {}
end
@@ -60,16 +61,17 @@ local function transform_node(new_nodes, node)
core.table.insert(new_nodes._priority_index, node.priority)
end
- new_nodes[node.priority][node.host .. ":" .. node.port] = node.weight
+ new_nodes[node.priority][node.host .. ":" .. node.port] = weight or
node.weight
return new_nodes
end
-local function fetch_all_nodes(upstream)
- local nodes = upstream.nodes
+-- `weights` carries the slow start weight of every node, indexed like `nodes`;
+-- without it each node keeps its configured weight
+local function transform_nodes(nodes, weights)
local new_nodes = core.table.new(0, #nodes)
- for _, node in ipairs(nodes) do
- new_nodes = transform_node(new_nodes, node)
+ for i, node in ipairs(nodes) do
+ new_nodes = transform_node(new_nodes, node, weights and weights[i])
end
return new_nodes
end
@@ -107,26 +109,27 @@ local function create_health_status(upstream, checker)
end
--- Build the picker node set from the healthy subset, reusing
create_health_status
--- so the per-node health lookup lives in exactly one place.
-local function fetch_health_nodes(upstream, checker)
+-- The nodes that actually reach the picker, reusing create_health_status so
the
+-- per-node health lookup lives in exactly one place. When every node is
unhealthy
+-- the whole set is kept, which is the existing fail-open behaviour.
+local function fetch_eligible_nodes(upstream, checker)
if not checker then
- return fetch_all_nodes(upstream)
+ return upstream.nodes
end
local health_status = create_health_status(upstream, checker)
if health_status.all_unhealthy then
- return fetch_all_nodes(upstream)
+ return upstream.nodes
end
- local up_nodes = core.table.new(0, #upstream.nodes)
+ local nodes = core.table.new(#upstream.nodes, 0)
for _, node in ipairs(upstream.nodes) do
if health_status.status[node.host .. ":" .. node.port] then
- up_nodes = transform_node(up_nodes, node)
+ core.table.insert(nodes, node)
end
end
- return up_nodes
+ return nodes
end
@@ -164,9 +167,12 @@ local function create_server_picker(upstream, checker)
local up_nodes
if upstream.type == "chash" then
- up_nodes = fetch_all_nodes(upstream)
+ up_nodes = transform_nodes(upstream.nodes)
else
- up_nodes = fetch_health_nodes(upstream, checker)
+ -- slow start runs on the eligible set, so a node only starts its
ramp
+ -- once it can actually be picked
+ local nodes = fetch_eligible_nodes(upstream, checker)
+ up_nodes = transform_nodes(nodes,
slow_start.effective_weights(upstream, nodes))
end
if #up_nodes._priority_index > 1 then
@@ -256,7 +262,11 @@ local function pick_server(route, ctx)
-- balancer here would leave it blind to everything routed before the
second
-- node showed up, which is the state a k8s deployment or a discovery
service
-- starts from. See #12217
- if nodes_count == 1 and up_conf.type ~= "least_conn" then
+ --
+ -- Slow start is in the same position: the node set of a single node
upstream
+ -- is what the second node is later compared against, and only the picker
build
+ -- records it. The node still takes every request either way.
+ if nodes_count == 1 and up_conf.type ~= "least_conn" and not
up_conf.warm_up_conf then
local node = up_conf.nodes[1]
ctx.balancer_ip = node.host
ctx.balancer_port = node.port
@@ -307,6 +317,10 @@ local function pick_server(route, ctx)
version = version .. "#" .. checker.status_ver
end
+ if up_conf.warm_up_conf then
+ version = version .. (slow_start.version_suffix(up_conf) or "")
+ end
+
-- the same picker will be used in the whole request, especially during
the retry
local server_picker = ctx.server_picker
if not server_picker then
diff --git a/apisix/cli/config.lua b/apisix/cli/config.lua
index cfb757bb4d..2063eb8d71 100644
--- a/apisix/cli/config.lua
+++ b/apisix/cli/config.lua
@@ -168,6 +168,7 @@ local _M = {
["plugin-limit-conn"] = "10m",
["worker-events"] = "10m",
["lrucache-lock"] = "10m",
+ ["upstream-slow-start"] = "10m",
["balancer-ewma"] = "10m",
["balancer-ewma-locks"] = "10m",
["balancer-ewma-last-touched-at"] = "10m",
diff --git a/apisix/cli/ngx_tpl.lua b/apisix/cli/ngx_tpl.lua
index f0a8475cdf..8d6947fea0 100644
--- a/apisix/cli/ngx_tpl.lua
+++ b/apisix/cli/ngx_tpl.lua
@@ -408,6 +408,7 @@ http {
lua_shared_dict internal-status {* http.lua_shared_dict["internal-status"]
*};
lua_shared_dict worker-events {* http.lua_shared_dict["worker-events"] *};
lua_shared_dict lrucache-lock {* http.lua_shared_dict["lrucache-lock"] *};
+ lua_shared_dict upstream-slow-start {*
http.lua_shared_dict["upstream-slow-start"] *};
lua_shared_dict balancer-ewma {* http.lua_shared_dict["balancer-ewma"] *};
lua_shared_dict balancer-ewma-locks {*
http.lua_shared_dict["balancer-ewma-locks"] *};
lua_shared_dict balancer-ewma-last-touched-at {*
http.lua_shared_dict["balancer-ewma-last-touched-at"] *};
diff --git a/apisix/init.lua b/apisix/init.lua
index 20d0caa1f1..b4ef63195e 100644
--- a/apisix/init.lua
+++ b/apisix/init.lua
@@ -94,6 +94,7 @@ function _M.http_init(args)
core.resolver.init_resolver(args)
core.id.init()
core.env.init()
+ require("apisix.slow_start").init()
local process = require("ngx.process")
local ok, err = process.enable_privileged_agent()
diff --git a/apisix/plugins/traffic-split.lua b/apisix/plugins/traffic-split.lua
index da6014e8b5..35243f502c 100644
--- a/apisix/plugins/traffic-split.lua
+++ b/apisix/plugins/traffic-split.lua
@@ -117,6 +117,16 @@ function _M.check_schema(conf)
end
end
end
+
+ -- the upstreams this plugin builds are rebuilt per request and
carry
+ -- no stable scope, so slow start has nowhere to keep the
lifecycle of
+ -- their nodes
+ for _, wupstream in ipairs(rule.weighted_upstreams or {}) do
+ if wupstream.upstream and wupstream.upstream.warm_up_conf then
+ return false, "warm_up_conf is not supported by the
upstream of " ..
+ "the traffic-split plugin"
+ end
+ end
end
end
diff --git a/apisix/schema_def.lua b/apisix/schema_def.lua
index 959cc941b0..0ff40e5646 100644
--- a/apisix/schema_def.lua
+++ b/apisix/schema_def.lua
@@ -414,6 +414,48 @@ local private_key_schema = {
}
+local warm_up_conf_schema = {
+ description = "slow start: ramp a newly observed node up to its configured
weight",
+ type = "object",
+ properties = {
+ slow_start_time_seconds = {
+ description = "seconds a new node takes to reach its full weight",
+ type = "integer",
+ minimum = 1,
+ },
+ min_weight_percent = {
+ description = "lowest effective weight, as a percentage of the
original weight",
+ type = "integer",
+ minimum = 1,
+ maximum = 100,
+ },
+ interval = {
+ description = "seconds between two effective weight refreshes",
+ type = "integer",
+ minimum = 1,
+ default = 1,
+ },
+ aggression = {
+ description = "shape of the ramp: 1 is linear, above 1 ramps up
faster " ..
+ "at the beginning, below 1 slower",
+ type = "number",
+ minimum = 0.01,
+ default = 1,
+ },
+ startup_grace_period_seconds = {
+ description = "seconds after the data plane started during which a
node " ..
+ "observed for the first time is considered mature",
+ type = "integer",
+ minimum = 0,
+ default = 0,
+ },
+ },
+ required = {"slow_start_time_seconds", "min_weight_percent"},
+ additionalProperties = false,
+}
+_M.warm_up_conf = warm_up_conf_schema
+
+
local upstream_schema = {
type = "object",
properties = {
@@ -427,6 +469,7 @@ local upstream_schema = {
-- properties
nodes = nodes_schema,
+ warm_up_conf = warm_up_conf_schema,
retries = {
type = "integer",
minimum = 0,
diff --git a/apisix/slow_start.lua b/apisix/slow_start.lua
new file mode 100644
index 0000000000..1ca6e8f1f3
--- /dev/null
+++ b/apisix/slow_start.lua
@@ -0,0 +1,531 @@
+--
+-- 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.
+--
+-- Slow start (`upstream.warm_up_conf`): a node that the data plane observes
for
+-- the first time takes a reduced share of the traffic and ramps back to its
+-- configured weight over `slow_start_time_seconds`.
+--
+-- Which node is new is decided here, by comparing the node set of the current
+-- picker build against the set the previous build recorded. No node lifecycle
+-- timestamp is read from the configuration, the Admin API or service
discovery:
+-- every start point is generated locally with `ngx.now()` and lives only in
the
+-- `upstream-slow-start` shared dict, so it is never written back to etcd.
+--
+-- Shared dict layout, all keys prefixed by the upstream scope (the resource
key
+-- of the standalone upstream, or of the route/service that embeds it):
+--
+-- <scope>|# node ids the configuration held at the previous
reconcile
+-- <scope>|! aggregate ramp deadline, read on the request hot
path
+-- <scope>|<id> ramp start: > 0 ramping, 0 mature, -1 not picked
yet
+-- <scope>|<id>|i when a reconcile first saw it configured but
unpickable
+-- <scope>|<id>@<start> the successor a worker elected for that ramp start
+--
+-- Workers reconcile in parallel, and the shared dict has no compare-and-set,
so
+-- a value is never derived from a read and then written back blindly: a worker
+-- that read a ramp start just before another one changed it would put the old
+-- one back. Instead a node's ramp start is only ever created with `add`, only
+-- ever replaced through an election on `add` keyed by the value being
replaced,
+-- and otherwise kept alive with `expire`, which leaves the value alone. Every
+-- worker that decides on a change to the same ramp start therefore ends up
+-- writing the same successor.
+--
+-- Presence and eligibility are tracked separately. Presence comes from the
+-- configuration, so every worker sees the same set; eligibility is this
worker's
+-- health view and only decides when a ramp starts and when it is interrupted.
+-- Tombstoning off the eligible set would let one worker's transient health
+-- opinion drop a node the other workers are still serving, and the state of a
+-- node that never left could then expire and be ramped again from scratch.
+--
+-- A node the configuration still holds keeps its entry alive. One that leaves
the
+-- configuration keeps it for a further `slow_start_time_seconds` (the
tombstone
+-- window) so a short absence resumes the ramp instead of restarting it. Coming
+-- back later, or being out of the picker for longer than a full window,
starts a
+-- new lifecycle.
+local core = require("apisix.core")
+local ipairs = ipairs
+local pairs = pairs
+local tostring = tostring
+local tonumber = tonumber
+local os_time = os.time
+local ngx_now = ngx.now
+local math_floor = math.floor
+local math_max = math.max
+local math_min = math.min
+local str_gmatch = string.gmatch
+local type = type
+
+local shdict = ngx.shared["upstream-slow-start"]
+
+local INSTANCE_STARTED_KEY = "@instance_started_at"
+local SNAPSHOT_SUFFIX = "|#"
+local DEADLINE_SUFFIX = "|!"
+local INELIGIBLE_SUFFIX = "|i"
+
+-- A node that is present but needs no ramp is stored with this ramp start, so
+-- that it can still be tombstoned when it disappears: without an entry it
would
+-- look new again on its way back.
+local MATURE = 0
+
+-- And one the configuration holds but that has never reached a picker - an
+-- active health check has not cleared it yet - is stored with this one. It is
+-- what separates "known and already mature" from "known but has not had its
+-- chance to ramp", which otherwise both look like a node with no state.
+local PENDING = -1
+
+-- Lifetime of the state of a node that is still eligible. Refreshed on every
+-- reconcile, which a busy upstream runs at least as often as the picker LRU
+-- expires (300s), so it only ever elapses for an upstream that stopped
receiving
+-- traffic altogether. State lost that way makes the whole node set a fresh
mature
+-- baseline again, which is the safe direction, and it is also what keeps the
+-- state of a deleted upstream from living forever.
+local STATE_TTL = 86400
+
+-- Lifetime of an election on a ramp start, see `replace_start`.
+local ELECTION_TTL = 10
+
+local _M = {}
+
+
+local reported_scopes = {}
+
+
+local function report_once(scope, ...)
+ if reported_scopes[scope] then
+ return
+ end
+ reported_scopes[scope] = true
+ core.log.error(...)
+end
+
+
+-- The instance start point behind `startup_grace_period_seconds`. Written from
+-- the master process with `add`, so a HUP reload - which keeps the shared
dict -
+-- continues to use the start point of the original start.
+function _M.init()
+ if not shdict then
+ return
+ end
+
+ local ok, err = shdict:add(INSTANCE_STARTED_KEY, os_time())
+ if not ok and err ~= "exists" then
+ core.log.error("failed to record the slow start instance start time:
", err)
+ end
+end
+
+
+local function scope_key(up_conf)
+ return up_conf.resource_key
+end
+
+
+-- The lifecycle identity of a node. A domain node is identified by the
hostname
+-- it was configured with, not by the address it currently resolves to, so a
DNS
+-- rotation does not restart its ramp.
+local function node_id(node)
+ return (node.domain or node.host) .. ":" .. tostring(node.port)
+end
+
+
+local function effective_weight(original_weight, first_seen_at, conf, now)
+ if original_weight <= 0 then
+ -- a node configured out of the rotation stays out of it
+ return original_weight
+ end
+
+ if not first_seen_at or first_seen_at <= MATURE then
+ return original_weight
+ end
+
+ local window = conf.slow_start_time_seconds
+ local elapsed = now - first_seen_at
+ if elapsed < 0 then
+ -- the local clock moved backwards; start the window over rather than
+ -- letting a negative elapsed produce a complex power below
+ elapsed = 0
+ end
+ if elapsed >= window then
+ return original_weight
+ end
+
+ -- `max(elapsed, 1)` keeps the first bucket at one second of progress, so a
+ -- node never ramps from exactly zero
+ local time_factor = math_max(elapsed, 1) / window
+ local ratio = math_max(conf.min_weight_percent / 100,
+ time_factor ^ (1 / (conf.aggression or 1)))
+ local weight = math_floor(original_weight * ratio)
+ if weight < 1 then
+ -- integer weights: anything above zero has to keep at least one unit,
+ -- otherwise the node gets no traffic at all and never warms up
+ weight = 1
+ end
+
+ return weight
+end
+_M.effective_weight = effective_weight
+
+
+local function within_startup_grace(conf, now)
+ local grace = conf.startup_grace_period_seconds or 0
+ if grace <= 0 then
+ return false
+ end
+
+ local started_at = shdict:get(INSTANCE_STARTED_KEY)
+ if not started_at then
+ return false
+ end
+
+ return now < started_at + grace
+end
+
+
+local function decode_snapshot(snapshot)
+ local known = {}
+ if not snapshot then
+ return nil
+ end
+
+ for id in str_gmatch(snapshot, "[^,]+") do
+ known[id] = true
+ end
+ return known
+end
+
+
+local function store(key, value, ttl)
+ local ok, err, forcible = shdict:set(key, value, ttl)
+ if not ok then
+ core.log.error("failed to store the slow start state of ", key, ": ",
err)
+ elseif forcible then
+ core.log.warn("the upstream-slow-start shared dict is full, storing ",
key,
+ " evicted another entry; nodes whose state is lost keep
",
+ "their configured weight")
+ end
+ return ok
+end
+
+
+local function keep_alive(key, ttl)
+ local ok, err = shdict:expire(key, ttl)
+ if not ok and err ~= "not found" then
+ core.log.error("failed to refresh the slow start state of ", key, ":
", err)
+ end
+end
+
+
+-- Create the ramp start of a node that has none, or adopt the one another
worker
+-- created first. Returns the ramp start everyone now shares, and whether this
+-- worker is the one that set it.
+local function create_start(key, candidate)
+ local ok, err, forcible = shdict:add(key, candidate, STATE_TTL)
+ if ok then
+ if forcible then
+ core.log.warn("the upstream-slow-start shared dict is full,
storing ", key,
+ " evicted another entry")
+ end
+ return candidate, true
+ end
+
+ if err ~= "exists" then
+ core.log.error("failed to store the slow start state of ", key, ": ",
err)
+ return MATURE, false
+ end
+
+ return tonumber(shdict:get(key)) or candidate, false
+end
+
+
+-- Replace the ramp start `from` with `to`, unless another worker already
elected
+-- a successor for `from`, in which case that one is used. The election
outlives
+-- the few microseconds a worker spends between reading `from` and getting here
+-- by a wide margin, and is short enough that `from` cannot come round again
for
+-- the same node before it expires.
+local function replace_start(key, from, to, window)
+ local election = key .. "@" .. from
+ local won, err = shdict:add(election, to, math_min(window, ELECTION_TTL))
+ if not won then
+ if err == "exists" then
+ to = tonumber(shdict:get(election)) or to
+ else
+ core.log.error("failed to elect the slow start state of ", key, ":
", err)
+ end
+ end
+
+ store(key, to, STATE_TTL)
+ return to, won
+end
+
+
+-- Compare the node set of this picker build against the one the previous build
+-- recorded, and return the ramp start point of every eligible node, indexed
like
+-- `nodes`. `present_nodes` is the configured set, identical in every worker;
+-- `nodes` is the subset this worker can actually pick from.
+local function reconcile(conf, present_nodes, nodes, scope, now)
+ local window = conf.slow_start_time_seconds
+
+ local present_ids = core.table.new(#present_nodes, 0)
+ local present = core.table.new(0, #present_nodes)
+ for i, node in ipairs(present_nodes) do
+ present_ids[i] = node_id(node)
+ present[present_ids[i]] = true
+ end
+
+ local ids = core.table.new(#nodes, 0)
+ local eligible = core.table.new(0, #nodes)
+ for i, node in ipairs(nodes) do
+ ids[i] = node_id(node)
+ eligible[ids[i]] = true
+ end
+
+ local snapshot_key = scope .. SNAPSHOT_SUFFIX
+ local known = decode_snapshot(shdict:get(snapshot_key))
+
+ -- The first reconcile of a scope is a baseline: the nodes an upstream is
+ -- bootstrapped with, and the nodes it already had when `warm_up_conf` was
+ -- turned on, are mature. So are nodes observed inside the startup grace
+ -- period, which absorbs the ordering differences of a cold restart.
+ local baseline = (known == nil) or within_startup_grace(conf, now)
+
+ local first_seen = core.table.new(#nodes, 0)
+ local deadline = 0
+
+ local function log_began(id, weight)
+ core.log.info("slow start began for node ", id, " of upstream ", scope,
+ ", weight ", weight, ", window ", window, "s, from ",
now)
+ end
+
+ for i, id in ipairs(ids) do
+ local key = scope .. "|" .. id
+ local ineligible_key = key .. INELIGIBLE_SUFFIX
+ local start_at = tonumber(shdict:get(key))
+ local ineligible_since = tonumber(shdict:get(ineligible_key))
+ local changed
+
+ if not start_at then
+ if baseline or (known and known[id]) then
+ -- either the bootstrap set, or a node whose state the shared
dict
+ -- evicted while it stayed configured: both keep the full
weight
+ -- rather than ramping a node that has been serving all along
+ start_at = create_start(key, MATURE)
+ else
+ start_at, changed = create_start(key, now)
+ if changed then
+ log_began(id, nodes[i].weight)
+ end
+ end
+
+ elseif start_at == PENDING then
+ -- the first picker it can actually be part of: this is where its
+ -- window starts, not when the configuration first mentioned it
+ start_at, changed = replace_start(key, PENDING, now, window)
+ if changed then
+ log_began(id, nodes[i].weight)
+ end
+
+ elseif ineligible_since and now - ineligible_since > window then
+ -- observed out of the picker for longer than a full window:
whatever
+ -- answers on this address now is not the process that was ramping
+ local out_for = now - ineligible_since
+ start_at, changed = replace_start(key, start_at, now, window)
+ if changed then
+ core.log.info("slow start restarted for node ", id, " of
upstream ",
+ scope, " after ", out_for, "s out of the picker")
+ end
+
+ elseif start_at > MATURE and now - start_at >= window then
+ start_at, changed = replace_start(key, start_at, MATURE, window)
+ if changed then
+ core.log.info("slow start finished for node ", id, " of
upstream ",
+ scope)
+ end
+
+ else
+ -- refresh the lifetime, which also drops the tombstone window if
the
+ -- node came back after having left the configuration
+ keep_alive(key, STATE_TTL)
+ end
+
+ if ineligible_since then
+ -- back in the picker
+ shdict:delete(ineligible_key)
+ end
+
+ first_seen[i] = start_at
+ if start_at > MATURE then
+ deadline = math_max(deadline, start_at + window)
+ end
+ end
+
+ for _, id in ipairs(present_ids) do
+ if not eligible[id] then
+ local key = scope .. "|" .. id
+ local start_at = tonumber(shdict:get(key))
+ if start_at then
+ -- still configured, just not pickable here: keep the state
alive
+ -- and mark when it dropped out, keeping the earliest mark. The
+ -- ramp clock itself runs on, so a short outage resumes where
it
+ -- left off
+ keep_alive(key, STATE_TTL)
+ local ineligible_key = key .. INELIGIBLE_SUFFIX
+ if not shdict:add(ineligible_key, now, STATE_TTL) then
+ keep_alive(ineligible_key, STATE_TTL)
+ end
+ elseif baseline or (known and known[id]) then
+ -- part of the bootstrap set even though a health check has not
+ -- cleared it yet, or state the shared dict evicted: every node
+ -- the previous reconcile saw was given state then
+ start_at = create_start(key, MATURE)
+ else
+ -- configured but never picked: its window starts when it first
+ -- becomes usable
+ start_at = create_start(key, PENDING)
+ end
+
+ -- a ramp this worker cannot pick from is still a ramp other
workers
+ -- may be serving, and the deadline is what keeps their picker keys
+ -- moving; leaving it out would let this worker publish an early
one
+ if start_at > MATURE then
+ deadline = math_max(deadline, start_at + window)
+ end
+ end
+ end
+
+ -- A node that left the configuration keeps its state for one slow start
+ -- window and is then forgotten, so that the address coming back later - a
+ -- different process behind the same host and port - is warmed up again.
+ if known then
+ for id in pairs(known) do
+ if not present[id] then
+ local key = scope .. "|" .. id
+ if shdict:get(key) then
+ core.log.info("node ", id, " of upstream ", scope,
+ " left the upstream, keeping its slow start
state for ",
+ window, "s")
+ keep_alive(key, window)
+ keep_alive(key .. INELIGIBLE_SUFFIX, window)
+ end
+ end
+ end
+ end
+
+ core.table.sort(present_ids)
+ store(snapshot_key, core.table.concat(present_ids, ","), STATE_TTL)
+ -- published for the hot path: zero means every node is mature
+ store(scope .. DEADLINE_SUFFIX, deadline, STATE_TTL)
+
+ return first_seen
+end
+
+
+local function usable(up_conf, nodes)
+ local scope = scope_key(up_conf)
+ if not scope then
+ core.log.error("slow start needs an upstream with a resource key, ",
+ "ignoring warm_up_conf")
+ return nil
+ end
+
+ if not shdict then
+ -- slow start only ramps HTTP upstreams, and the stream subsystem has
no
+ -- such shared dict. Like every other upstream field that does not
apply
+ -- there, `warm_up_conf` is quietly ignored and the configured weights
are
+ -- used, rather than failing the connection or logging on every build
+ return nil
+ end
+
+ if up_conf.type ~= "roundrobin" then
+ report_once(scope, "slow start only supports roundrobin, ignoring ",
+ "warm_up_conf of upstream ", scope)
+ return nil
+ end
+
+ if nodes then
+ local priority = nodes[1] and nodes[1].priority
+ for _, node in ipairs(nodes) do
+ if node.priority ~= priority then
+ report_once(scope, "slow start does not support an upstream
with ",
+ "mixed node priorities, ignoring warm_up_conf of
upstream ",
+ scope)
+ return nil
+ end
+ end
+ end
+
+ return scope
+end
+
+
+-- Effective weight of every node of this picker build, indexed like `nodes`,
or
+-- nil when the upstream does not use slow start. Runs once per picker build,
not
+-- per request.
+function _M.effective_weights(up_conf, nodes)
+ local conf = up_conf.warm_up_conf
+ if type(conf) ~= "table" then
+ return nil
+ end
+
+ local scope = usable(up_conf, nodes)
+ if not scope then
+ -- settle the picker cache key: without a deadline `version_suffix`
keeps
+ -- appending a fresh time bucket, so an upstream that can never ramp
would
+ -- rebuild its picker every `interval` for nothing
+ if shdict and up_conf.resource_key then
+ store(up_conf.resource_key .. DEADLINE_SUFFIX, 0, STATE_TTL)
+ end
+ return nil
+ end
+
+ local now = ngx_now()
+ local first_seen = reconcile(conf, up_conf.nodes, nodes, scope, now)
+
+ local weights = core.table.new(#nodes, 0)
+ for i, node in ipairs(nodes) do
+ weights[i] = effective_weight(node.weight, first_seen[i], conf, now)
+ end
+
+ return weights
+end
+
+
+-- Suffix of the picker cache key. While any node ramps, the key carries the
+-- current `interval` bucket so that the picker is rebuilt once per bucket per
+-- worker; once every node is mature the key settles on a stable suffix, and
the
+-- rebuild it causes is the one that restores the full weights.
+--
+-- A rebuild restarts the round robin cursor, which favours the heaviest node
for
+-- the first few picks. Over a bucket that carries real traffic this averages
out,
+-- but on an upstream that sees only a handful of requests per `interval` the
+-- ramping node can end up with even less traffic than its weight asks for.
+function _M.version_suffix(up_conf)
+ local conf = up_conf.warm_up_conf
+ if type(conf) ~= "table" or not shdict or not up_conf.resource_key
+ or up_conf.type ~= "roundrobin" then
+ return nil
+ end
+
+ local now = ngx_now()
+ local deadline = shdict:get(up_conf.resource_key .. DEADLINE_SUFFIX)
+ -- an unknown deadline means no picker has been built for this upstream
yet;
+ -- the build that this bucket triggers is the one that publishes it
+ if deadline and now >= deadline then
+ return "#wm"
+ end
+
+ return "#w" .. math_floor(now / (conf.interval or 1))
+end
+
+
+return _M
diff --git a/apisix/upstream.lua b/apisix/upstream.lua
index 99e1e857d3..acb6ab3663 100644
--- a/apisix/upstream.lua
+++ b/apisix/upstream.lua
@@ -559,6 +559,58 @@ local function get_chash_key_schema(hash_on)
end
+-- Constraints of `warm_up_conf` within one upstream that JSON schema cannot
+-- express. A ramp needs a single roundrobin tier to work in, so a combination
it
+-- could never act on is rejected at the Admin API rather than accepted and
+-- ignored. Where the upstream is used is a different question: like every
other
+-- field that only applies to HTTP, `warm_up_conf` is simply ignored on the
+-- stream path.
+--
+-- This runs on the configuration entry points only, never on the data plane
+-- checker: a configuration that reaches a running gateway some other way - it
is
+-- written to etcd or to a standalone config file directly, or it is embedded
in a
+-- route or a service, neither of which runs this - must not take the whole
+-- upstream out of service over a field
+-- that only accelerates a ramp. `slow_start.usable()` logs and keeps proxying
+-- with the configured weights there.
+local function check_warm_up_conf(conf)
+ local warm_up_conf = conf.warm_up_conf
+ if not warm_up_conf then
+ return true
+ end
+
+ if (conf.type or "roundrobin") ~= "roundrobin" then
+ return false, "warm_up_conf is only supported by the roundrobin
upstream type"
+ end
+
+ local interval = warm_up_conf.interval or 1
+ if interval > warm_up_conf.slow_start_time_seconds then
+ return false, "warm_up_conf.interval can't be greater than " ..
+ "warm_up_conf.slow_start_time_seconds"
+ end
+
+ -- APISIX drains the highest priority tier before it uses the next one, so
a
+ -- weight ramp inside one tier can't hold traffic back from a new node in a
+ -- tier above the mature ones
+ local nodes = conf.nodes
+ if nodes and core.table.isarray(nodes) then
+ local priority
+ for i, node in ipairs(nodes) do
+ local node_priority = node.priority or 0
+ if i == 1 then
+ priority = node_priority
+ elseif node_priority ~= priority then
+ return false, "warm_up_conf doesn't support an upstream with "
..
+ "nodes of different priorities"
+ end
+ end
+ end
+
+ return true
+end
+_M.check_warm_up_conf = check_warm_up_conf
+
+
local function check_upstream_conf(in_dp, conf)
if not in_dp then
local ok, err = check_schema(conf)
@@ -566,6 +618,11 @@ local function check_upstream_conf(in_dp, conf)
return false, "invalid configuration: " .. err
end
+ local ok, err = check_warm_up_conf(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/conf/config.yaml.example b/conf/config.yaml.example
index b9f10261b0..86510feb5c 100644
--- a/conf/config.yaml.example
+++ b/conf/config.yaml.example
@@ -355,6 +355,7 @@ nginx_config: # Config for render the
template to generate n
plugin-limit-conn: 10m
worker-events: 10m
lrucache-lock: 10m
+ upstream-slow-start: 10m # Slow start state of upstream nodes
(`warm_up_conf`)
balancer-ewma: 10m
balancer-ewma-locks: 10m
balancer-ewma-last-touched-at: 10m
diff --git a/docs/en/latest/admin-api.md b/docs/en/latest/admin-api.md
index abd9965d37..77d8d76969 100644
--- a/docs/en/latest/admin-api.md
+++ b/docs/en/latest/admin-api.md
@@ -1024,6 +1024,11 @@ In addition to the equalization algorithm selections,
Upstream also supports pas
| keepalive_pool.size | False
| Auxiliary | Sets `keepalive` directive
dynamically.
[...]
| keepalive_pool.idle_timeout | False
| Auxiliary | Sets `keepalive_timeout`
directive dynamically.
[...]
| keepalive_pool.requests | False
| Auxiliary | Sets `keepalive_requests`
directive dynamically.
[...]
+| warm_up_conf.slow_start_time_seconds | True, when `warm_up_conf` is set |
Integer | Slow start window in seconds. A node the gateway observes for the
first time takes a reduced share of the traffic and ramps back to its
configured weight over this window. Must be at least 1. | 300 |
+| warm_up_conf.min_weight_percent | True, when `warm_up_conf` is set | Integer
| Lowest effective weight during the ramp, as a percentage of the configured
weight, from 1 to 100. | 1 |
+| warm_up_conf.interval | False | Integer | Seconds between two effective
weight refreshes. Defaults to `1`, and cannot be greater than
`slow_start_time_seconds`. | 1 |
+| warm_up_conf.aggression | False | Number | Shape of the ramp. `1` (default)
is linear, above `1` ramps up faster at the beginning, below `1` slower. At
least `0.01`. | 1 |
+| warm_up_conf.startup_grace_period_seconds | False | Integer | Seconds after
the gateway starts during which a node observed for the first time is treated
as already warmed up, so that a restart does not ramp the whole node set again.
Defaults to `0`. | 180 |
An Upstream can be one of the following `types`:
@@ -1067,6 +1072,15 @@ To verify the certificate presented by the Upstream, set
`tls.verify` to `true`.
To allow Upstream to have a separate connection pool, use `keepalive_pool`. It
can be configured by modifying its child fields.
+`warm_up_conf` enables slow start for the nodes of a `roundrobin` Upstream.
Whether a node is new is decided by the gateway itself, from the node set it
observes, and the start of each ramp is recorded locally in the
`upstream-slow-start` shared dict:
+
+- The node set an Upstream has when the gateway first builds a load balancer
for it is treated as warmed up. So is the node set it already has when
`warm_up_conf` is turned on.
+- A node added afterwards ramps from `min_weight_percent` back to its
configured weight over `slow_start_time_seconds`. A node held out of the load
balancer by a health check starts its ramp when it first becomes available.
+- A node that leaves the Upstream and comes back within
`slow_start_time_seconds` resumes its ramp. One that comes back later, or that
a health check kept out for longer than that, ramps again from the start.
+- Every APISIX instance ramps independently, from the moment it observed the
node.
+
+`warm_up_conf` is only supported by `roundrobin` Upstreams whose nodes share a
single priority, and it is rejected in the Upstreams of the `traffic-split`
Plugin, which are rebuilt per request. Like the other Upstream fields that only
apply to HTTP, it is ignored when the Upstream is used by a stream route. A
ramp only shifts traffic between nodes: a single-node Upstream, or one whose
nodes are all new, keeps sending every request to them.
+
Example Configuration:
```shell
diff --git a/docs/zh/latest/admin-api.md b/docs/zh/latest/admin-api.md
index aae200b88f..ec6aa9c8a7 100644
--- a/docs/zh/latest/admin-api.md
+++ b/docs/zh/latest/admin-api.md
@@ -1032,6 +1032,11 @@ APISIX 的 Upstream 除了基本的负载均衡算法选择外,还支持对上
|keepalive_pool.size | 否 | 辅助 |
动态设置 `keepalive` 指令,详细信息请参考下文。 |
|keepalive_pool.idle_timeout | 否
| 辅助 | 动态设置 `keepalive_timeout` 指令,详细信息请参考下文。 |
|keepalive_pool.requests | 否 | 辅助
| 动态设置 `keepalive_requests` 指令,详细信息请参考下文。 |
+|warm_up_conf.slow_start_time_seconds | 设置 `warm_up_conf` 时必填 | 整型 |
慢启动窗口,单位为秒。网关首次观察到的节点先承接较少的流量,并在该窗口内逐步恢复到配置的权重。最小值为 1。 |
+|warm_up_conf.min_weight_percent | 设置 `warm_up_conf` 时必填 | 整型 |
爬坡期间有效权重的下限,以配置权重的百分比表示,取值范围 1 到 100。 |
+|warm_up_conf.interval | 否 | 整型 | 两次刷新有效权重之间的间隔,单位为秒。默认为 `1`,不能大于
`slow_start_time_seconds`。 |
+|warm_up_conf.aggression | 否 | 数值 | 爬坡曲线。`1`(默认)为线性,大于 `1` 前期增长更快,小于 `1`
前期增长更慢。最小值为 `0.01`。 |
+|warm_up_conf.startup_grace_period_seconds | 否 | 整型 |
网关启动后的宽限时间,单位为秒。在此期间首次观察到的节点直接视为已完成预热,避免重启后整组节点重新爬坡。默认为 `0`。 |
`type` 详细信息如下:
@@ -1070,6 +1075,12 @@ APISIX 的 Upstream 除了基本的负载均衡算法选择外,还支持对上
```
- `keepalive_pool` 允许 Upstream 有自己单独的连接池。它下属的字段,比如 `requests`,可以用于配置上游连接保持的参数。
+- `warm_up_conf` 为 `roundrobin` 类型的 Upstream
开启节点慢启动。节点是否为新节点由网关根据自己观察到的节点集合判断,每次爬坡的起点记录在本地的 `upstream-slow-start` 共享字典中:
+ - 网关首次为 Upstream 构建负载均衡器时已有的节点视为已完成预热,开启 `warm_up_conf` 时 Upstream 已有的节点同样如此。
+ - 之后新增的节点从 `min_weight_percent` 开始,在 `slow_start_time_seconds`
内逐步恢复到配置的权重。被健康检查挡在负载均衡器之外的节点,从它首次可用时开始爬坡。
+ - 节点离开 Upstream 后在 `slow_start_time_seconds`
内回来,会继续原来的爬坡;更晚回来,或被健康检查排除超过该时长,会重新开始爬坡。
+ - 每个 APISIX 实例独立计时,起点为该实例观察到节点的时刻。
+ - `warm_up_conf` 仅支持节点优先级一致的 `roundrobin` 类型 Upstream,`traffic-split` 插件中的
Upstream 每请求重建,不允许配置。与其他只对 HTTP 生效的 Upstream 字段一样,该字段在 stream route 使用的
Upstream 上会被忽略。爬坡只在节点之间调整流量:单节点 Upstream,或所有节点都是新节点时,请求仍会全部发往这些节点。
Upstream 对象 JSON 配置示例:
diff --git a/t/APISIX.pm b/t/APISIX.pm
index 76dc94530a..dca50da5cb 100644
--- a/t/APISIX.pm
+++ b/t/APISIX.pm
@@ -660,6 +660,7 @@ _EOC_
lua_shared_dict internal-status 10m;
lua_shared_dict worker-events 10m;
lua_shared_dict lrucache-lock 10m;
+ lua_shared_dict upstream-slow-start 10m;
lua_shared_dict balancer-ewma 1m;
lua_shared_dict balancer-ewma-locks 1m;
lua_shared_dict balancer-ewma-last-touched-at 1m;
diff --git a/t/admin/upstream-slow-start.t b/t/admin/upstream-slow-start.t
new file mode 100644
index 0000000000..58dd861e86
--- /dev/null
+++ b/t/admin/upstream-slow-start.t
@@ -0,0 +1,622 @@
+#
+# 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_root_location();
+no_shuffle();
+log_level("info");
+
+add_block_preprocessor(sub {
+ my ($block) = @_;
+
+ if (!$block->request) {
+ $block->set_value("request", "GET /t");
+ }
+
+ if (!$block->no_error_log && !$block->error_log) {
+ $block->set_value("no_error_log", "[error]\n[alert]");
+ }
+});
+
+run_tests;
+
+__DATA__
+
+=== TEST 1: set upstream with warm_up_conf
+--- 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",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 300,
+ "min_weight_percent": 1
+ }
+ }]]
+ )
+
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+
+
+
+=== TEST 2: defaults are filled in by the schema
+--- config
+ location /t {
+ content_by_lua_block {
+ local core = require("apisix.core")
+ local conf = {
+ type = "roundrobin",
+ nodes = {{host = "127.0.0.1", port = 1980, weight = 100}},
+ warm_up_conf = {
+ slow_start_time_seconds = 300,
+ min_weight_percent = 1,
+ },
+ }
+
+ local ok, err = core.schema.check(core.schema.upstream, conf)
+ if not ok then
+ ngx.say("failed: ", err)
+ return
+ end
+
+ ngx.say("interval: ", conf.warm_up_conf.interval,
+ ", aggression: ", conf.warm_up_conf.aggression,
+ ", startup_grace_period_seconds: ",
+ conf.warm_up_conf.startup_grace_period_seconds)
+ }
+ }
+--- response_body
+interval: 1, aggression: 1, startup_grace_period_seconds: 0
+
+
+
+=== TEST 3: all fields
+--- 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",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 300,
+ "min_weight_percent": 20,
+ "interval": 5,
+ "aggression": 2.5,
+ "startup_grace_period_seconds": 180
+ }
+ }]]
+ )
+
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+
+
+
+=== TEST 4: slow_start_time_seconds is required
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "min_weight_percent": 1
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/property .*slow_start_time_seconds.* is required/
+
+
+
+=== TEST 5: min_weight_percent is required
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/property .*min_weight_percent.* is required/
+
+
+
+=== TEST 6: min_weight_percent is a percentage, not a ratio
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ for _, percent in ipairs({0, 101}) do
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": ]] .. percent .. [[
+ }
+ }]]
+ )
+ if code < 300 then
+ ngx.say("unexpectedly accepted min_weight_percent ",
percent)
+ return
+ end
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+
+
+
+=== TEST 7: reject slow_start_time_seconds below 1 and aggression below 0.01
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local cases = {
+ [["slow_start_time_seconds": 0, "min_weight_percent": 1]],
+ [["slow_start_time_seconds": 10, "min_weight_percent": 1,
"aggression": 0]],
+ [["slow_start_time_seconds": 10, "min_weight_percent": 1,
"interval": 0]],
+ [["slow_start_time_seconds": 10, "min_weight_percent": 1,
"default_weight": 1]],
+ [["slow_start_time_seconds": 10, "min_weight_percent": 1,
+ "startup_grace_period_seconds": -1]],
+ [["slow_start_time_seconds": 10, "min_weight_percent": 1,
"unknown": 1]],
+ }
+
+ for i, case in ipairs(cases) do
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {]] .. case .. [[}
+ }]]
+ )
+ if code < 300 then
+ ngx.say("unexpectedly accepted case ", i)
+ return
+ end
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+
+
+
+=== TEST 8: reject warm_up_conf on a non roundrobin upstream
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "chash",
+ "key": "remote_addr",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/warm_up_conf is only supported by the roundrobin upstream type/
+
+
+
+=== TEST 9: reject an interval longer than the slow start window
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1,
+ "interval": 11
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/warm_up_conf.interval can't be greater than
warm_up_conf.slow_start_time_seconds/
+
+
+
+=== TEST 10: reject nodes with different priorities
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100,
"priority": 0},
+ {"host": "127.0.0.1", "port": 1981, "weight": 100,
"priority": -1}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/warm_up_conf doesn't support an upstream with nodes of different priorities/
+
+
+
+=== TEST 11: accept warm_up_conf on a route embedded upstream
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/routes/1',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/hello",
+ "upstream": {
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }
+ }]]
+ )
+
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+
+
+
+=== TEST 12: accept warm_up_conf on a service embedded upstream
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/services/1',
+ ngx.HTTP_PUT,
+ [[{
+ "upstream": {
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }
+ }]]
+ )
+
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+
+
+
+=== TEST 13: reject warm_up_conf on a route embedded upstream that is not
roundrobin
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/routes/2',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/hello",
+ "upstream": {
+ "type": "least_conn",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/warm_up_conf is only supported by the roundrobin upstream type/
+
+
+
+=== TEST 14: reject warm_up_conf in a traffic-split upstream
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/routes/3',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/hello",
+ "plugins": {
+ "traffic-split": {
+ "rules": [{
+ "weighted_upstreams": [{
+ "upstream": {
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port":
1981, "weight": 100}
+ ],
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ },
+ "weight": 1
+ }]
+ }]
+ }
+ },
+ "upstream": {
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100}
+ ]
+ }
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/warm_up_conf is not supported by the upstream of the traffic-split plugin/
+
+
+
+=== TEST 15: declarative validation accepts warm_up_conf
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/configs/validate',
+ ngx.HTTP_POST,
+ [[{
+ "upstreams": [
+ {
+ "id": "u1",
+ "type": "roundrobin",
+ "nodes": {"127.0.0.1:1980": 1},
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }
+ ]
+ }]]
+ )
+
+ ngx.status = code
+ ngx.say(body)
+ }
+ }
+--- error_code: 200
+--- response_body
+passed
+
+
+
+=== TEST 16: declarative validation still rejects an unusable warm_up_conf
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/configs/validate',
+ ngx.HTTP_POST,
+ [[{
+ "upstreams": [
+ {
+ "id": "u1",
+ "type": "chash",
+ "key": "remote_addr",
+ "nodes": {"127.0.0.1:1980": 1},
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1
+ }
+ }
+ ]
+ }]]
+ )
+
+ ngx.status = code
+ ngx.print(body)
+ }
+ }
+--- error_code: 400
+--- response_body eval
+qr/warm_up_conf is only supported by the roundrobin upstream type/
+
+
+
+=== TEST 17: the data plane keeps an upstream it cannot ramp
+--- config
+ location /t {
+ content_by_lua_block {
+ local core = require("apisix.core")
+ local apisix_upstream = require("apisix.upstream")
+
+ -- the Admin API rejects this, but a control plane or an embedded
+ -- upstream reaches the data plane without passing through it; the
+ -- upstream still has to load, or every route using it returns 503
+ local conf = {
+ type = "chash",
+ key = "remote_addr",
+ nodes = {{host = "127.0.0.1", port = 1980, weight = 100,
priority = 0}},
+ warm_up_conf = {
+ slow_start_time_seconds = 10,
+ min_weight_percent = 1,
+ },
+ }
+
+ local ok, err = apisix_upstream.check_upstream_conf(conf)
+ ngx.say("admin: ", tostring(ok), " ", tostring(err))
+
+ local dp_conf = core.table.deepcopy(conf)
+ local dp_ok, dp_err = core.schema.check(core.schema.upstream,
dp_conf)
+ ngx.say("data plane schema: ", tostring(dp_ok), " ",
tostring(dp_err))
+ }
+ }
+--- response_body
+admin: false warm_up_conf is only supported by the roundrobin upstream type
+data plane schema: true nil
+
+
+
+=== TEST 18: clean up
+--- 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',
+ '/apisix/admin/upstreams/1?force=true'}) do
+ local code, body = t(uri, ngx.HTTP_DELETE)
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(uri, ": ", body)
+ return
+ end
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
diff --git a/t/node/upstream-slow-start.t b/t/node/upstream-slow-start.t
new file mode 100644
index 0000000000..4fa385b473
--- /dev/null
+++ b/t/node/upstream-slow-start.t
@@ -0,0 +1,833 @@
+#
+# 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_root_location();
+no_shuffle();
+log_level("info");
+workers(1);
+
+add_block_preprocessor(sub {
+ my ($block) = @_;
+
+ if (!$block->request) {
+ $block->set_value("request", "GET /t");
+ }
+
+ if (!$block->no_error_log && !$block->error_log) {
+ $block->set_value("no_error_log", "[error]\n[alert]");
+ }
+});
+
+run_tests;
+
+__DATA__
+
+=== TEST 1: the effective weight follows the ramp
+--- config
+ location /t {
+ content_by_lua_block {
+ local slow_start = require("apisix.slow_start")
+ local conf = {
+ slow_start_time_seconds = 300,
+ min_weight_percent = 1,
+ interval = 1,
+ aggression = 1,
+ }
+ local now = 10000
+
+ -- the first bucket counts as one second of progress, not zero:
over a
+ -- 10s window that is a tenth of the weight, well above the 1%
floor
+ ngx.say(slow_start.effective_weight(100, now, {
+ slow_start_time_seconds = 10,
+ min_weight_percent = 1,
+ aggression = 1,
+ }, now))
+ ngx.say(slow_start.effective_weight(100, now, conf, now))
+ ngx.say(slow_start.effective_weight(100, now, conf, now + 30))
+ ngx.say(slow_start.effective_weight(100, now, conf, now + 150))
+ ngx.say(slow_start.effective_weight(100, now, conf, now + 300))
+ ngx.say(slow_start.effective_weight(100, now, conf, now + 3000))
+ -- a node configured out of the rotation stays out of it
+ ngx.say(slow_start.effective_weight(0, now, conf, now + 30))
+ -- anything above zero keeps at least one weight unit
+ ngx.say(slow_start.effective_weight(1, now, conf, now + 1))
+ -- a node with no state, and a node marked mature, use the full
weight
+ ngx.say(slow_start.effective_weight(100, nil, conf, now + 1))
+ ngx.say(slow_start.effective_weight(100, 0, conf, now + 1))
+ }
+ }
+--- response_body
+10
+1
+10
+50
+100
+100
+0
+1
+100
+100
+
+
+
+=== TEST 2: min_weight_percent is the floor, aggression the shape
+--- config
+ location /t {
+ content_by_lua_block {
+ local slow_start = require("apisix.slow_start")
+ local now = 10000
+ local function weight(percent, aggression, elapsed)
+ return slow_start.effective_weight(100, now, {
+ slow_start_time_seconds = 300,
+ min_weight_percent = percent,
+ aggression = aggression,
+ }, now + elapsed)
+ end
+
+ ngx.say(weight(20, 1, 1))
+ ngx.say(weight(1, 1, 30))
+ ngx.say(weight(1, 2, 30) > weight(1, 1, 30))
+ ngx.say(weight(1, 0.5, 30) < weight(1, 1, 30))
+ }
+ }
+--- response_body
+20
+10
+true
+true
+
+
+
+=== TEST 3: a local clock jump does not produce a weight outside the range
+--- config
+ location /t {
+ content_by_lua_block {
+ local slow_start = require("apisix.slow_start")
+ local conf = {
+ slow_start_time_seconds = 300,
+ min_weight_percent = 1,
+ aggression = 1,
+ }
+
+ -- backwards
+ ngx.say(slow_start.effective_weight(100, 10000, conf, 9000))
+ -- and forwards, clamped to the full window
+ ngx.say(slow_start.effective_weight(100, 10000, conf, 99999))
+ }
+ }
+--- response_body
+1
+100
+
+
+
+=== TEST 4: a node added to a running upstream ramps up, the existing one does
not
+--- timeout: 30
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local http = require("resty.http")
+ local uri = "http://127.0.0.1:" .. ngx.var.server_port ..
"/server_port"
+
+ local function set_upstream(nodes, desc)
+ return t('/apisix/admin/upstreams/1',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "desc": "]] .. desc .. [[",
+ "nodes": ]] .. nodes .. [[,
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1,
+ "interval": 1
+ }
+ }]]
+ )
+ end
+
+ local function split(times)
+ local ports = {}
+ for _ = 1, times do
+ local httpc = http.new()
+ local res, err = httpc:request_uri(uri)
+ if not res then
+ return nil, err
+ end
+ ports[res.body] = (ports[res.body] or 0) + 1
+ end
+ return ports
+ end
+
+ local one = [=[[{"host": "127.0.0.1", "port": 1980, "weight":
100}]]=]
+ local two = [=[[{"host": "127.0.0.1", "port": 1980, "weight": 100},
+ {"host": "127.0.0.1", "port": 1981, "weight":
100}]]=]
+
+ local code, body = set_upstream(one, "one node")
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ code, body = t('/apisix/admin/routes/1',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/server_port",
+ "upstream_id": "1"
+ }]]
+ )
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ -- the first picker build records the bootstrapped node as the
baseline
+ ngx.sleep(0.5)
+ split(1)
+
+ code, body = set_upstream(two, "two nodes")
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.sleep(0.5)
+
+ local ports, err = split(20)
+ if not ports then
+ ngx.say(err)
+ return
+ end
+ local mature, warming = ports["1980"] or 0, ports["1981"] or 0
+ ngx.log(ngx.WARN, "ramping split: 1980=", mature, " 1981=",
warming)
+ if not (mature >= 17 and warming <= 3) then
+ ngx.say("failed while ramping: 1980=", mature, " 1981=",
warming)
+ return
+ end
+
+ -- once the window has passed both nodes are back to their weights
+ ngx.sleep(11)
+ ports, err = split(20)
+ if not ports then
+ ngx.say(err)
+ return
+ end
+ if ports["1980"] ~= 10 or ports["1981"] ~= 10 then
+ ngx.say("failed after the window: 1980=",
tostring(ports["1980"]),
+ " 1981=", tostring(ports["1981"]))
+ return
+ end
+
+ -- an unrelated field changing must not start the ramp over
+ code, body = set_upstream(two, "a new description")
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.sleep(0.5)
+
+ ports, err = split(20)
+ if not ports then
+ ngx.say(err)
+ return
+ end
+ if ports["1980"] ~= 10 or ports["1981"] ~= 10 then
+ ngx.say("failed after an unrelated change: 1980=",
+ tostring(ports["1980"]), " 1981=",
tostring(ports["1981"]))
+ return
+ end
+
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed
+--- error_log
+slow start began for node 127.0.0.1:1981
+slow start finished for node 127.0.0.1:1981
+--- no_error_log eval
+[qr/\[error\]/, qr/\[alert\]/, qr/slow start began for node 127\.0\.0\.1:1980/]
+
+
+
+=== TEST 5: a node that comes back inside the tombstone window keeps its
progress
+--- timeout: 30
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local http = require("resty.http")
+ local uri = "http://127.0.0.1:" .. ngx.var.server_port ..
"/slow-start-tombstone"
+
+ local function set_upstream(nodes)
+ return t('/apisix/admin/upstreams/2',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": ]] .. nodes .. [[,
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1,
+ "interval": 10
+ }
+ }]]
+ )
+ end
+
+ local function split(times)
+ local ports = {}
+ for _ = 1, times do
+ local httpc = http.new()
+ local res, err = httpc:request_uri(uri)
+ if not res then
+ return nil, err
+ end
+ ports[res.body] = (ports[res.body] or 0) + 1
+ end
+ return ports
+ end
+
+ local one = [=[[{"host": "127.0.0.1", "port": 1980, "weight":
100}]]=]
+ local two = [=[[{"host": "127.0.0.1", "port": 1980, "weight": 100},
+ {"host": "127.0.0.1", "port": 1981, "weight":
100}]]=]
+
+ local code, body = set_upstream(one)
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ code, body = t('/apisix/admin/routes/2',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/slow-start-tombstone",
+ "plugins": {
+ "proxy-rewrite": {"uri": "/server_port"}
+ },
+ "upstream_id": "2"
+ }]]
+ )
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ ngx.sleep(0.5)
+ split(1)
+
+ -- 1981 joins and starts its ramp
+ set_upstream(two)
+ ngx.sleep(0.5)
+ split(1)
+
+ -- it leaves, and comes back well inside the 10s tombstone window
+ set_upstream(one)
+ ngx.sleep(0.5)
+ split(1)
+ ngx.sleep(4)
+ set_upstream(two)
+ ngx.sleep(0.5)
+
+ local ports, err = split(100)
+ if not ports then
+ ngx.say(err)
+ return
+ end
+
+ -- the ramp continued while it was away, so 1981 is about half way
+ -- through the window: a real share of the traffic, still below
1980.
+ -- Restarting the ramp would have put it back at a tenth of the
weight
+ local mature, warming = ports["1980"] or 0, ports["1981"] or 0
+ ngx.log(ngx.WARN, "tombstone split: 1980=", mature, " 1981=",
warming)
+ if warming >= 15 and warming < mature then
+ ngx.say("passed")
+ else
+ ngx.say("failed: 1980=", mature, " 1981=", warming)
+ end
+ }
+ }
+--- response_body
+passed
+--- grep_error_log eval
+qr/slow start began for node \S+/
+--- grep_error_log_out
+slow start began for node 127.0.0.1:1981
+
+
+
+=== TEST 6: without warm_up_conf the weights are used as configured
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local code, body = t('/apisix/admin/upstreams/3',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100},
+ {"host": "127.0.0.1", "port": 1981, "weight": 100}
+ ]
+ }]]
+ )
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ code, body = t('/apisix/admin/routes/3',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/slow-start-off",
+ "plugins": {
+ "proxy-rewrite": {"uri": "/server_port"}
+ },
+ "upstream_id": "3"
+ }]]
+ )
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ ngx.sleep(0.5)
+
+ local http = require("resty.http")
+ local uri = "http://127.0.0.1:" .. ngx.var.server_port ..
"/slow-start-off"
+ local ports = {}
+ for _ = 1, 20 do
+ local httpc = http.new()
+ local res, err = httpc:request_uri(uri)
+ if not res then
+ ngx.say(err)
+ return
+ end
+ ports[res.body] = (ports[res.body] or 0) + 1
+ end
+
+ if ports["1980"] == 10 and ports["1981"] == 10 then
+ ngx.say("passed")
+ else
+ ngx.say("failed: 1980=", tostring(ports["1980"]),
+ " 1981=", tostring(ports["1981"]))
+ end
+ }
+ }
+--- response_body
+passed
+--- no_error_log
+slow start began for node
+
+
+
+=== TEST 7: a route embedded upstream is its own lifecycle scope
+--- timeout: 15
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local http = require("resty.http")
+ local uri = "http://127.0.0.1:" .. ngx.var.server_port ..
"/slow-start-embedded"
+
+ local function set_route(nodes)
+ return t('/apisix/admin/routes/4',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/slow-start-embedded",
+ "plugins": {
+ "proxy-rewrite": {"uri": "/server_port"}
+ },
+ "upstream": {
+ "type": "roundrobin",
+ "nodes": ]] .. nodes .. [[,
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1,
+ "interval": 1
+ }
+ }
+ }]]
+ )
+ end
+
+ local one = [=[[{"host": "127.0.0.1", "port": 1980, "weight":
100}]]=]
+ local two = [=[[{"host": "127.0.0.1", "port": 1980, "weight": 100},
+ {"host": "127.0.0.1", "port": 1981, "weight":
100}]]=]
+
+ local code, body = set_route(one)
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ ngx.sleep(0.5)
+ local httpc = http.new()
+ httpc:request_uri(uri)
+
+ code, body = set_route(two)
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+ ngx.sleep(0.5)
+
+ local ports = {}
+ for _ = 1, 20 do
+ httpc = http.new()
+ local res, err = httpc:request_uri(uri)
+ if not res then
+ ngx.say(err)
+ return
+ end
+ ports[res.body] = (ports[res.body] or 0) + 1
+ end
+
+ local mature, warming = ports["1980"] or 0, ports["1981"] or 0
+ if mature >= 17 and warming <= 3 then
+ ngx.say("passed")
+ else
+ ngx.say("failed: 1980=", mature, " 1981=", warming)
+ end
+ }
+ }
+--- response_body
+passed
+--- error_log eval
+qr{of upstream \S*/routes/4}
+
+
+
+=== TEST 8: the picker cache key changes once per interval, then settles
+--- config
+ location /t {
+ content_by_lua_block {
+ local core = require("apisix.core")
+ local slow_start = require("apisix.slow_start")
+ local dict = ngx.shared["upstream-slow-start"]
+ local up_conf = {
+ resource_key = "/upstreams/version-suffix",
+ type = "roundrobin",
+ warm_up_conf = {
+ slow_start_time_seconds = 10,
+ min_weight_percent = 1,
+ interval = 2,
+ },
+ }
+
+ -- no deadline published yet: the build this triggers is the one
that
+ -- publishes it, so the key already carries a bucket
+ local first = slow_start.version_suffix(up_conf)
+ ngx.say("bucketed: ", first ~= nil and first:match("^#w%d+$") ~=
nil)
+
+ -- a second request inside the same bucket reuses the cached picker
+ ngx.say("stable within the interval: ",
+ slow_start.version_suffix(up_conf) == first)
+
+ local function publish_deadline(deadline)
+ local ok, err = dict:set(up_conf.resource_key .. "|!",
deadline)
+ if not ok then
+ error("failed to publish the deadline: " .. err)
+ end
+ end
+
+ -- while a node ramps, crossing the bucket rebuilds it exactly once
+ publish_deadline(ngx.now() + 10)
+ ngx.sleep(2.1)
+ local next_bucket = slow_start.version_suffix(up_conf)
+ ngx.say("rebuilt after the interval: ", next_bucket ~= first)
+ ngx.say("stable again: ", slow_start.version_suffix(up_conf) ==
next_bucket)
+
+ -- once every node is mature the key stops moving altogether
+ publish_deadline(0)
+ ngx.say("settled: ", slow_start.version_suffix(up_conf))
+ ngx.sleep(2.1)
+ ngx.say("still settled: ", slow_start.version_suffix(up_conf))
+
+ dict:delete(up_conf.resource_key .. "|!")
+ }
+ }
+--- timeout: 15
+--- response_body
+bucketed: true
+stable within the interval: true
+rebuilt after the interval: true
+stable again: true
+settled: #wm
+still settled: #wm
+
+
+
+=== TEST 9: an unhealthy node keeps its lifecycle, it is not treated as removed
+--- timeout: 20
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ local http = require("resty.http")
+
+ -- 1979 has nothing listening, so the active check takes it out of
the
+ -- picker while the configuration still holds it
+ local code, body = t('/apisix/admin/upstreams/5',
+ ngx.HTTP_PUT,
+ [[{
+ "type": "roundrobin",
+ "nodes": [
+ {"host": "127.0.0.1", "port": 1980, "weight": 100},
+ {"host": "127.0.0.1", "port": 1979, "weight": 100}
+ ],
+ "checks": {
+ "active": {
+ "http_path": "/status",
+ "healthy": {"interval": 1, "successes": 1},
+ "unhealthy": {"interval": 1, "tcp_failures": 1}
+ }
+ },
+ "warm_up_conf": {
+ "slow_start_time_seconds": 10,
+ "min_weight_percent": 1,
+ "interval": 1
+ }
+ }]]
+ )
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ code, body = t('/apisix/admin/routes/5',
+ ngx.HTTP_PUT,
+ [[{
+ "uri": "/slow-start-unhealthy",
+ "plugins": {
+ "proxy-rewrite": {"uri": "/server_port"}
+ },
+ "upstream_id": "5"
+ }]]
+ )
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(body)
+ return
+ end
+
+ local uri = "http://127.0.0.1:" .. ngx.var.server_port ..
"/slow-start-unhealthy"
+ local function hit(times)
+ local ports = {}
+ for _ = 1, times do
+ local httpc = http.new()
+ local res, err = httpc:request_uri(uri)
+ if not res then
+ return nil, err
+ end
+ ports[res.body] = (ports[res.body] or 0) + 1
+ end
+ return ports
+ end
+
+ -- requests before the checker has taken 1979 out may land on it
and
+ -- come back as a 502; only a request that cannot be made at all
fails
+ local ports, err
+ ngx.sleep(0.5)
+ ports, err = hit(2)
+ if not ports then
+ ngx.say("request failed: ", err)
+ return
+ end
+ -- let the active checker settle and keep rebuilding the picker
+ ngx.sleep(3)
+ ports, err = hit(5)
+ if not ports then
+ ngx.say("request failed: ", err)
+ return
+ end
+ ngx.sleep(2)
+ ports, err = hit(10)
+ if not ports then
+ ngx.say("request failed: ", err)
+ return
+ end
+
+ ngx.say("1980: ", ports["1980"] or 0)
+ }
+ }
+--- response_body
+1980: 10
+--- no_error_log
+left the upstream
+
+
+
+=== TEST 10: a node kept out of the picker starts its window when it first
gets in
+--- config
+ location /t {
+ content_by_lua_block {
+ local slow_start = require("apisix.slow_start")
+ local n1 = {host = "10.0.0.1", port = 8080, weight = 100, priority
= 0}
+ local n2 = {host = "10.0.0.2", port = 8080, weight = 100, priority
= 0}
+ local up_conf = {
+ resource_key = "/upstreams/pending",
+ type = "roundrobin",
+ warm_up_conf = {
+ slow_start_time_seconds = 100,
+ min_weight_percent = 1,
+ interval = 1,
+ aggression = 1,
+ },
+ }
+
+ -- the upstream is bootstrapped with one node
+ up_conf.nodes = {n1}
+ local weights = slow_start.effective_weights(up_conf, {n1})
+ ngx.say("baseline: ", weights[1])
+
+ -- a second node is configured, but a health check keeps it out of
the
+ -- picker for a while. It must not age into maturity while
sidelined
+ up_conf.nodes = {n1, n2}
+ weights = slow_start.effective_weights(up_conf, {n1})
+ ngx.say("while sidelined: ", weights[1], " ", tostring(weights[2]))
+
+ -- the first picker it reaches is where its window starts
+ weights = slow_start.effective_weights(up_conf, {n1, n2})
+ ngx.say("first picker: ", weights[1], " ", weights[2])
+ }
+ }
+--- response_body
+baseline: 100
+while sidelined: 100 nil
+first picker: 100 1
+--- error_log
+slow start began for node 10.0.0.2:8080
+
+
+
+=== TEST 11: a worker never writes back a ramp start another worker just
changed
+--- config
+ location /t {
+ content_by_lua_block {
+ local slow_start = require("apisix.slow_start")
+ local dict = ngx.shared["upstream-slow-start"]
+ local mt = getmetatable(dict)
+ local orig_get = mt.get
+
+ -- Reconciles in different workers run truly in parallel. Replay
the
+ -- losing interleaving deterministically: this worker reads a
node's
+ -- state, and another worker writes before this one does
+ local function race(target, other_worker)
+ local fired = false
+ mt.get = function(self, key, ...)
+ local value, flags = orig_get(self, key, ...)
+ if self == dict and key == target and not fired then
+ fired = true
+ other_worker()
+ end
+ return value, flags
+ end
+ end
+
+ local function run(fn)
+ local ok, err = pcall(fn)
+ mt.get = orig_get
+ if not ok then
+ error(err)
+ end
+ end
+
+ local conf = {
+ slow_start_time_seconds = 100,
+ min_weight_percent = 1,
+ interval = 1,
+ aggression = 1,
+ }
+ local n1 = {host = "10.0.1.1", port = 80, weight = 100, priority =
0}
+ local n2 = {host = "10.0.1.2", port = 80, weight = 100, priority =
0}
+ local scope = "/upstreams/race"
+ local key = scope .. "|10.0.1.2:80"
+ local up_conf = {resource_key = scope, type = "roundrobin",
warm_up_conf = conf}
+
+ up_conf.nodes = {n1}
+ slow_start.effective_weights(up_conf, {n1})
+ up_conf.nodes = {n1, n2}
+ slow_start.effective_weights(up_conf, {n1})
+ ngx.say("pending: ", dict:get(key))
+
+ -- both workers see n2 enter the picker; the other one elects first
+ local elected = ngx.now() - 50
+ local weights
+ run(function()
+ race(key, function()
+ assert(dict:add(key .. "@-1", elected, 10))
+ assert(dict:set(key, elected))
+ end)
+ weights = slow_start.effective_weights(up_conf, {n1, n2})
+ end)
+ ngx.say("adopted the elected start: ", dict:get(key) == elected,
+ ", weight ", weights[2])
+
+ -- the other worker finishes the ramp while this one, having read
the
+ -- old start, only means to refresh it
+ run(function()
+ race(key, function()
+ assert(dict:add(key .. "@" .. elected, 0, 10))
+ assert(dict:set(key, 0))
+ end)
+ slow_start.effective_weights(up_conf, {n1, n2})
+ end)
+ ngx.say("refresh left the change alone: ", dict:get(key))
+ }
+ }
+--- response_body
+pending: -1
+adopted the elected start: true, weight 50
+refresh left the change alone: 0
+
+
+
+=== TEST 12: clean up
+--- config
+ location /t {
+ content_by_lua_block {
+ local t = require("lib.test_admin").test
+ for _, uri in ipairs({'/apisix/admin/routes/1',
'/apisix/admin/routes/2',
+ '/apisix/admin/routes/3',
'/apisix/admin/routes/4',
+ '/apisix/admin/routes/5',
+ '/apisix/admin/upstreams/1?force=true',
'/apisix/admin/upstreams/2?force=true',
+ '/apisix/admin/upstreams/3?force=true',
'/apisix/admin/upstreams/5?force=true'}) do
+ local code, body = t(uri, ngx.HTTP_DELETE)
+ if code >= 300 then
+ ngx.status = code
+ ngx.say(uri, ": ", body)
+ return
+ end
+ end
+ ngx.say("passed")
+ }
+ }
+--- response_body
+passed