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 6c6877521 fix(balancer): keep least_conn load state across upstream 
scaling (#13666)
6c6877521 is described below

commit 6c6877521807aaeba906fa038cddfba4ec0d38ee
Author: AlinsRan <[email protected]>
AuthorDate: Fri Jul 17 15:13:26 2026 +0800

    fix(balancer): keep least_conn load state across upstream scaling (#13666)
---
 apisix/balancer.lua              |  23 +-
 apisix/balancer/chash.lua        |  11 +-
 apisix/balancer/ewma.lua         |  11 +-
 apisix/balancer/least_conn.lua   | 192 ++++++++++---
 apisix/balancer/priority.lua     |   4 +-
 apisix/balancer/roundrobin.lua   |  11 +-
 apisix/init.lua                  |   4 +
 apisix/upstream.lua              |  28 +-
 t/lib/server.lua                 |  24 ++
 t/node/least_conn.t              |  37 +++
 t/node/least_conn3.t             | 575 +++++++++++++++++++++++++++++++++++++++
 t/node/upstream-discovery-ipv6.t | 192 +++++++++++++
 t/stream-node/least_conn.t       |  74 +++++
 13 files changed, 1140 insertions(+), 46 deletions(-)

diff --git a/apisix/balancer.lua b/apisix/balancer.lua
index 1647bab2f..ff2f8fb8d 100644
--- a/apisix/balancer.lua
+++ b/apisix/balancer.lua
@@ -174,9 +174,9 @@ local function create_server_picker(upstream, checker)
             return server_picker
         end
 
-        core.log.info("upstream nodes: ",
-                      
core.json.delay_encode(up_nodes[up_nodes._priority_index[1]]))
-        local server_picker = 
picker.new(up_nodes[up_nodes._priority_index[1]], upstream)
+        local priority = up_nodes._priority_index[1]
+        core.log.info("upstream nodes: ", 
core.json.delay_encode(up_nodes[priority]))
+        local server_picker = picker.new(up_nodes[priority], upstream, 
priority)
         server_picker.addr_to_domain = addr_to_domain
         return server_picker
     end
@@ -248,7 +248,13 @@ local function pick_server(route, ctx)
     local up_conf = ctx.upstream_conf
 
     local nodes_count = #up_conf.nodes
-    if nodes_count == 1 then
+    -- least_conn counts the in-flight connections of every request it routes, 
so it
+    -- has to see them even while the upstream has a single node: those 
connections
+    -- are what tells a later scale out that the node is not empty. Skipping 
the
+    -- 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
         local node = up_conf.nodes[1]
         ctx.balancer_ip = node.host
         ctx.balancer_port = node.port
@@ -268,7 +274,10 @@ local function pick_server(route, ctx)
     ctx.balancer_try_count = (ctx.balancer_try_count or 0) + 1
     if ctx.balancer_try_count > 1 then
         if ctx.server_picker and ctx.server_picker.after_balance then
-            ctx.server_picker.after_balance(ctx, true)
+            -- remembering the server as tried is what keeps the next pick off 
it, so
+            -- only do it when there is another one to move to. With a single 
node the
+            -- retry has to land on it again, the way the fast path below 
always did
+            ctx.server_picker.after_balance(ctx, nodes_count > 1)
         end
 
         if checker then
@@ -331,6 +340,9 @@ local function pick_server(route, ctx)
         return nil, "failed to find valid upstream server, all upstream 
servers tried"
     end
     ctx.balancer_server = server
+    -- from here on the request holds a server, so the log phase must be able 
to
+    -- release it even if we bail out below
+    ctx.server_picker = server_picker
 
     local domain = server_picker.addr_to_domain[server]
     local res, err = lrucache_addr(server, nil, parse_addr, server)
@@ -347,7 +359,6 @@ local function pick_server(route, ctx)
     if is_http and ctx.var then
         ctx.var.upstream_unresolved_host = ctx.upstream_unresolved_host
     end
-    ctx.server_picker = server_picker
     res.upstream_host = parse_server_for_upstream_host(res, 
ctx.upstream_scheme)
 
     return res
diff --git a/apisix/balancer/chash.lua b/apisix/balancer/chash.lua
index f0e971a36..a8ab632fb 100644
--- a/apisix/balancer/chash.lua
+++ b/apisix/balancer/chash.lua
@@ -123,6 +123,11 @@ function _M.new(up_nodes, upstream)
             return servers[id]
         end,
         after_balance = function (ctx, before_retry)
+            -- the release is what makes the request stop holding the server, 
so drop
+            -- the reference here instead of leaving every caller to remember
+            local server = ctx.balancer_server
+            ctx.balancer_server = nil
+
             if not before_retry then
                 if ctx.balancer_tried_servers then
                     core.tablepool.release("balancer_tried_servers", 
ctx.balancer_tried_servers)
@@ -132,11 +137,15 @@ function _M.new(up_nodes, upstream)
                 return nil
             end
 
+            if not server then
+                return nil
+            end
+
             if not ctx.balancer_tried_servers then
                 ctx.balancer_tried_servers = 
core.tablepool.fetch("balancer_tried_servers", 0, 2)
             end
 
-            ctx.balancer_tried_servers[ctx.balancer_server] = true
+            ctx.balancer_tried_servers[server] = true
             ctx.balancer_tried_servers_count = 
(ctx.balancer_tried_servers_count or 0) + 1
         end,
         before_retry_next_priority = function (ctx)
diff --git a/apisix/balancer/ewma.lua b/apisix/balancer/ewma.lua
index 0834397b9..93707c400 100644
--- a/apisix/balancer/ewma.lua
+++ b/apisix/balancer/ewma.lua
@@ -187,12 +187,21 @@ local function _ewma_find(ctx, up_nodes)
 end
 
 local function _ewma_after_balance(ctx, before_retry)
+    -- the release is what makes the request stop holding the server, so drop 
the
+    -- reference here instead of leaving every caller to remember
+    local server = ctx.balancer_server
+    ctx.balancer_server = nil
+
     if before_retry then
+        if not server then
+            return nil
+        end
+
         if not ctx.balancer_tried_servers then
             ctx.balancer_tried_servers = 
core.tablepool.fetch("balancer_tried_servers", 0, 2)
         end
 
-        ctx.balancer_tried_servers[ctx.balancer_server] = true
+        ctx.balancer_tried_servers[server] = true
         ctx.balancer_tried_servers_count = (ctx.balancer_tried_servers_count 
or 0) + 1
 
         return nil
diff --git a/apisix/balancer/least_conn.lua b/apisix/balancer/least_conn.lua
index 8923d1781..3aa008ffc 100644
--- a/apisix/balancer/least_conn.lua
+++ b/apisix/balancer/least_conn.lua
@@ -24,66 +24,184 @@ local pairs = pairs
 local _M = {}
 
 
+-- Per-worker balancing state, shared by every picker built for the same 
upstream.
+--
+-- A picker is cached by the upstream version, so it is thrown away whenever 
the
+-- upstream changes: scaling, a config update, or a health status flip. The
+-- requests that are still in flight keep the picker they were routed with, and
+-- release their server on that picker in the log phase. If the heap lived 
inside
+-- the picker, the rebuilt one would never see those releases: the connections
+-- established before the rebuild would be forgotten on creation and then, once
+-- they closed, decremented on a heap nobody reads anymore. Long-lived
+-- connections (WebSocket) would keep the load skewed on the original nodes and
+-- least_conn would degrade to round-robin. See #12217.
+--
+-- Keeping the heap and the in-flight counts here, keyed by something stable
+-- across scaling, gives every generation of pickers a single view of the load.
+--
+-- The values are weak, which is exactly the lifetime this state needs. A 
picker
+-- holds its state, an in-flight request holds the picker it was routed with, 
and
+-- the picker cache holds the current one - so a state survives for as long as
+-- anything can still release a connection into it. Once nothing references it
+-- there is no connection left to count, and dropping it costs nothing. That 
also
+-- means there is no size to tune and no eviction that could quietly forget a 
busy
+-- upstream, which is what an LRU would do here: it would rank states by how
+-- recently they were rebuilt, and a stable upstream holding many long-lived
+-- connections is precisely the one that is never rebuilt.
+local states = setmetatable({}, {__mode = "v"})
+
+
 local function least_score(a, b)
     return a.score < b.score
 end
 
 
-function _M.new(up_nodes, upstream)
-    local servers_heap = binaryHeap.minUnique(least_score)
+local function new_state()
+    return {
+        heap = binaryHeap.minUnique(least_score),
+        -- server -> in-flight connections, only holds positive counts
+        conns = {},
+        -- server -> true, mirrors the payloads currently in the heap
+        members = {},
+    }
+end
+
+
+local function update_score(state, server)
+    local info = state.heap:valueByPayload(server)
+    -- the server may have left the upstream while it still held connections
+    if not info then
+        return
+    end
+
+    info.score = (1 + (state.conns[server] or 0)) * info.effect_weight
+    state.heap:update(server, info)
+end
+
+
+-- Align the long-lived heap with the current node set, keeping the in-flight
+-- counts of the nodes that survive. A node that is added back later (scaled in
+-- again, or reported healthy again) gets its score restored from `conns`.
+local function sync_nodes(state, up_nodes)
+    local heap = state.heap
+
+    for server in pairs(state.members) do
+        if not up_nodes[server] then
+            heap:remove(server)
+            state.members[server] = nil
+        end
+    end
+
     for server, weight in pairs(up_nodes) do
-        local score = 1 / weight
-        -- Note: the argument order of insert is different from others
-        servers_heap:insert({
-            server = server,
-            effect_weight = 1 / weight,
-            score = score,
-        }, server)
+        local effect_weight = 1 / weight
+        local info = heap:valueByPayload(server)
+        if info then
+            info.effect_weight = effect_weight
+        else
+            -- Note: the argument order of insert is different from others
+            heap:insert({
+                server = server,
+                effect_weight = effect_weight,
+                score = effect_weight,
+            }, server)
+            state.members[server] = true
+        end
+        -- one place decides what a score is worth
+        update_score(state, server)
     end
+end
+
+
+function _M.new(up_nodes, upstream, priority)
+    -- resource_key identifies the upstream and is stable across node scaling, 
unlike
+    -- the picker version which changes whenever the nodes change. Do not fall 
back to
+    -- resource_id: it is a bare id, so a route and an upstream sharing one 
would land
+    -- on the same heap and evict each other's nodes
+    local up_key = upstream.resource_key
+    local state
+    if up_key and priority then
+        -- each priority level owns a disjoint node set, so it needs its own 
heap.
+        -- Do not default the priority: a caller that does not name one has a 
node
+        -- set we cannot place, and folding it into level 0 would let 
sync_nodes
+        -- evict that level's nodes from the heap it shares
+        local state_key = up_key .. "#" .. priority
+        state = states[state_key]
+        if not state then
+            state = new_state()
+            states[state_key] = state
+        end
+    else
+        -- no stable identity, fall back to a state private to this picker
+        state = new_state()
+    end
+
+    sync_nodes(state, up_nodes)
+
+    local servers_heap = state.heap
+    local conns = state.conns
 
     return {
         upstream = upstream,
         get = function (ctx)
+            local tried = ctx.balancer_tried_servers
             local server, info, err
-            if ctx.balancer_tried_servers then
-                local tried_server_list = {}
-                while true do
-                    server, info = servers_heap:peek()
-                    -- we need to let the retry > #nodes so this branch can be 
hit and
-                    -- the request will retry next priority of nodes
-                    if server == nil then
-                        err = "all upstream servers tried"
-                        break
-                    end
-
-                    if not ctx.balancer_tried_servers[server] then
-                        break
-                    end
-
-                    servers_heap:pop()
-                    core.table.insert(tried_server_list, info)
+            local skipped
+
+            while true do
+                server, info = servers_heap:peek()
+                -- we need to let the retry > #nodes so this branch can be hit 
and
+                -- the request will retry next priority of nodes
+                if server == nil then
+                    err = "all upstream servers tried"
+                    break
                 end
 
-                for _, info in ipairs(tried_server_list) do
-                    servers_heap:insert(info, info.server)
+                -- the heap is shared with the pickers built for later 
versions of
+                -- the upstream, so it can hold nodes this request's conf does 
not
+                -- know about. Only hand out the ones it does
+                if up_nodes[server] and not (tried and tried[server]) then
+                    break
+                end
+
+                servers_heap:pop()
+                if not skipped then
+                    skipped = {}
+                end
+                core.table.insert(skipped, info)
+            end
+
+            if skipped then
+                for _, skipped_info in ipairs(skipped) do
+                    servers_heap:insert(skipped_info, skipped_info.server)
                 end
-            else
-                server, info = servers_heap:peek()
             end
 
             if not server then
                 return nil, err
             end
 
-            info.score = info.score + info.effect_weight
-            servers_heap:update(server, info)
+            conns[server] = (conns[server] or 0) + 1
+            update_score(state, server)
             return server
         end,
         after_balance = function (ctx, before_retry)
+            -- the release is what makes the request stop holding the server, 
so drop
+            -- the reference here instead of leaving every caller to remember. 
A caller
+            -- that goes on to retry gets a fresh one from the next pick
             local server = ctx.balancer_server
-            local info = servers_heap:valueByPayload(server)
-            info.score = info.score - info.effect_weight
-            servers_heap:update(server, info)
+            ctx.balancer_server = nil
+
+            if server then
+                local count = (conns[server] or 0) - 1
+                if count < 0 then
+                    -- a release with no matching pick. The store below floors 
the
+                    -- count either way, so the score cannot be corrupted - 
but the
+                    -- accounting is wrong and it should not pass in silence
+                    core.log.error("released a connection never picked on ", 
server)
+                end
+                conns[server] = count > 0 and count or nil
+                update_score(state, server)
+            end
 
             if not before_retry then
                 if ctx.balancer_tried_servers then
@@ -94,6 +212,10 @@ function _M.new(up_nodes, upstream)
                 return nil
             end
 
+            if not server then
+                return nil
+            end
+
             if not ctx.balancer_tried_servers then
                 ctx.balancer_tried_servers = 
core.tablepool.fetch("balancer_tried_servers", 0, 2)
             end
diff --git a/apisix/balancer/priority.lua b/apisix/balancer/priority.lua
index af5d60cbb..40e65f3bc 100644
--- a/apisix/balancer/priority.lua
+++ b/apisix/balancer/priority.lua
@@ -33,7 +33,9 @@ function _M.new(up_nodes, upstream, picker_mod)
 
     local pickers = core.table.new(#priority_index, 0)
     for i, priority in ipairs(priority_index) do
-        local picker, err = picker_mod.new(up_nodes[priority], upstream)
+        -- the priority is part of the picker's identity: node sets of 
different
+        -- priorities are disjoint and must not share balancing state
+        local picker, err = picker_mod.new(up_nodes[priority], upstream, 
priority)
         if not picker then
             return nil, "failed to create picker with priority " .. priority 
.. ": " .. err
         end
diff --git a/apisix/balancer/roundrobin.lua b/apisix/balancer/roundrobin.lua
index 7090f5261..ca97e9916 100644
--- a/apisix/balancer/roundrobin.lua
+++ b/apisix/balancer/roundrobin.lua
@@ -58,6 +58,11 @@ function _M.new(up_nodes, upstream)
             return server
         end,
         after_balance = function (ctx, before_retry)
+            -- the release is what makes the request stop holding the server, 
so drop
+            -- the reference here instead of leaving every caller to remember
+            local server = ctx.balancer_server
+            ctx.balancer_server = nil
+
             if not before_retry then
                 if ctx.balancer_tried_servers then
                     core.tablepool.release("balancer_tried_servers", 
ctx.balancer_tried_servers)
@@ -67,11 +72,15 @@ function _M.new(up_nodes, upstream)
                 return nil
             end
 
+            if not server then
+                return nil
+            end
+
             if not ctx.balancer_tried_servers then
                 ctx.balancer_tried_servers = 
core.tablepool.fetch("balancer_tried_servers", 0, 2)
             end
 
-            ctx.balancer_tried_servers[ctx.balancer_server] = true
+            ctx.balancer_tried_servers[server] = true
             ctx.balancer_tried_servers_count = 
(ctx.balancer_tried_servers_count or 0) + 1
         end,
         before_retry_next_priority = function (ctx)
diff --git a/apisix/init.lua b/apisix/init.lua
index 162affd3a..3c0ff689f 100644
--- a/apisix/init.lua
+++ b/apisix/init.lua
@@ -1541,6 +1541,10 @@ function _M.stream_log_phase()
 
     healthcheck_passive(api_ctx)
 
+    if api_ctx.server_picker and api_ctx.server_picker.after_balance then
+        api_ctx.server_picker.after_balance(api_ctx, false)
+    end
+
     core.ctx.release_vars(api_ctx)
     if api_ctx.plugins then
         core.tablepool.release("plugins", api_ctx.plugins)
diff --git a/apisix/upstream.lua b/apisix/upstream.lua
index ef89f2a22..232b0eb0d 100644
--- a/apisix/upstream.lua
+++ b/apisix/upstream.lua
@@ -113,6 +113,22 @@ local scheme_to_port = {
 _M.scheme_to_port = scheme_to_port
 
 
+-- A bare (unbracketed) IPv6 host makes the "host:port" key the balancer and 
the
+-- health checker build ambiguous (parse_addr reads the whole thing as an 
address
+-- with no port). The Admin API rejects such hosts and check_upstream_conf 
brackets
+-- them for configured upstreams, but service discovery returns nodes that 
reach
+-- here without going through either.
+--
+-- Only a node that carries its own port is bracketed. A host given without 
one is
+-- itself ambiguous - the map key "::1:1980" is a valid IPv6 literal as much 
as it
+-- is host ::1 port 1980 - so it is left to the existing malformed-config 
handling.
+-- Discovery endpoints always carry a port, so nothing real is missed.
+local function needs_ipv6_bracket(node)
+    return node.port and core.utils.parse_ipv6(node.host)
+                     and str_byte(node.host, 1) ~= str_byte("[")
+end
+
+
 local function fill_node_info(up_conf, scheme, is_stream)
     local nodes = up_conf.nodes
     if up_conf.nodes_ref == nodes then
@@ -135,6 +151,10 @@ local function fill_node_info(up_conf, scheme, is_stream)
         if not n.priority then
             need_filled = true
         end
+
+        if needs_ipv6_bracket(n) then
+            need_filled = true
+        end
     end
 
     if not need_filled then
@@ -148,10 +168,12 @@ local function fill_node_info(up_conf, scheme, is_stream)
     -- keep the original nodes for slow path in `compare_upstream_node()`,
     -- can't use `core.table.deepcopy()` for whole `nodes` array here,
     -- because `compare_upstream_node()` compare `metadata` of node by address.
+    -- The original (bare) host is preserved there, so bracketing below does 
not
+    -- make discovery re-fetches look like a node change.
     up_conf.original_nodes = core.table.new(#nodes, 0)
     for i, n in ipairs(nodes) do
         up_conf.original_nodes[i] = core.table.clone(n)
-        if not n.port or not n.priority then
+        if not n.port or not n.priority or needs_ipv6_bracket(n) then
             nodes[i] = core.table.clone(n)
 
             if not is_stream and not n.port then
@@ -162,6 +184,10 @@ local function fill_node_info(up_conf, scheme, is_stream)
             if not n.priority then
                 nodes[i].priority = 0
             end
+
+            if needs_ipv6_bracket(n) then
+                nodes[i].host = "[" .. n.host .. "]"
+            end
         end
     end
 
diff --git a/t/lib/server.lua b/t/lib/server.lua
index d933478cd..71b34912a 100644
--- a/t/lib/server.lua
+++ b/t/lib/server.lua
@@ -386,6 +386,30 @@ end
 _M.websocket_handshake_route = _M.websocket_handshake
 
 
+-- keep the session open until the peer goes away, so that the request stays in
+-- flight in the balancer the way a real WebSocket session does. An idle 
timeout is
+-- the normal state of such a session, not an error: keep waiting, and only 
give up
+-- once the peer closes or the connection breaks
+function _M.websocket_hold()
+    local websocket = require "resty.websocket.server"
+    local wb, err = websocket:new({timeout = 30000})
+    if not wb then
+        ngx.log(ngx.ERR, "failed to new websocket: ", err)
+        return ngx.exit(400)
+    end
+
+    while true do
+        local _, typ, err = wb:recv_frame()
+        if typ == "close" then
+            return
+        end
+        if not typ and not string.find(err or "", "timeout", 1, true) then
+            return
+        end
+    end
+end
+
+
 function _M.api_breaker()
     ngx.exit(tonumber(ngx.var.arg_code))
 end
diff --git a/t/node/least_conn.t b/t/node/least_conn.t
index 174252fd7..cf79b2674 100644
--- a/t/node/least_conn.t
+++ b/t/node/least_conn.t
@@ -149,3 +149,40 @@ qr/proxy request to \S+ while connecting to upstream/
 --- grep_error_log_out
 proxy request to 127.0.0.1:1999 while connecting to upstream
 proxy request to 0.0.0.0:1999 while connecting to upstream
+
+
+
+=== TEST 5: more retries than nodes, the request ends up holding no server
+--- apisix_yaml
+upstreams:
+  - id: 1
+    type: least_conn
+    retries: 3
+    nodes:
+        "127.0.0.1:1999": 2
+        "0.0.0.0:1999": 1
+--- error_code: 502
+--- error_log
+failed to find valid upstream server, all upstream servers tried
+--- no_error_log
+table index is nil
+
+
+
+=== TEST 6: a single node is retried, it is the only one there is
+--- apisix_yaml
+upstreams:
+  - id: 1
+    type: least_conn
+    retries: 2
+    nodes:
+        "127.0.0.1:1999": 1
+--- error_code: 502
+--- error_log
+connect() failed
+--- grep_error_log eval
+qr/proxy request to \S+ while connecting to upstream/
+--- grep_error_log_out
+proxy request to 127.0.0.1:1999 while connecting to upstream
+proxy request to 127.0.0.1:1999 while connecting to upstream
+proxy request to 127.0.0.1:1999 while connecting to upstream
diff --git a/t/node/least_conn3.t b/t/node/least_conn3.t
new file mode 100644
index 000000000..723108f22
--- /dev/null
+++ b/t/node/least_conn3.t
@@ -0,0 +1,575 @@
+#
+# 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(2);
+log_level('info');
+no_root_location();
+no_shuffle();
+
+run_tests();
+
+__DATA__
+
+=== TEST 1: keep in-flight conn count across balancer recreation on scaling
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            -- resource_key is stable across scaling, so both pickers share the
+            -- same connection count table
+            local up = {resource_key = "/upstreams/lc-scale"}
+
+            -- 2 nodes serving long-lived connections
+            local nodes = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1}
+            local p1 = least_conn.new(nodes, up, 0)
+
+            -- establish 4 in-flight connections (get without after_balance)
+            local ctx = {}
+            local held = {}
+            for _ = 1, 4 do
+                held[#held + 1] = p1.get(ctx)
+            end
+
+            -- scale out: add a third node, the picker is recreated
+            local scaled = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1,
+                            ["127.0.0.1:1982"] = 1}
+            local p2 = least_conn.new(scaled, up, 0)
+
+            -- the freshly added node has no connection, so it must be picked 
first
+            for _ = 1, 2 do
+                local s = p2.get(ctx)
+                held[#held + 1] = s
+                ngx.say(s)
+            end
+
+            -- release everything so repeated runs start from a clean state
+            for _, s in ipairs(held) do
+                ctx.balancer_server = s
+                p2.after_balance(ctx, false)
+            end
+        }
+    }
+--- request
+GET /t
+--- response_body
+127.0.0.1:1982
+127.0.0.1:1982
+
+
+
+=== TEST 2: a drained node returns to the pool across balancer recreation
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            local up = {resource_key = "/upstreams/lc-drain"}
+
+            local nodes = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1}
+            local p1 = least_conn.new(nodes, up, 0)
+
+            -- hold 4 in-flight connections, 2 land on each node
+            local ctx = {}
+            local held = {}
+            for _ = 1, 4 do
+                held[#held + 1] = p1.get(ctx)
+            end
+
+            -- drain every connection on 1980
+            for _, s in ipairs(held) do
+                if s == "127.0.0.1:1980" then
+                    ctx.balancer_server = s
+                    p1.after_balance(ctx, false)
+                end
+            end
+
+            -- picker recreated: 1980 is back to baseline and must be preferred
+            -- over 1981 which is still holding connections
+            local p2 = least_conn.new(nodes, up, 0)
+            local s = p2.get(ctx)
+            ngx.say(s)
+
+            -- release the rest so repeated runs start from a clean state
+            ctx.balancer_server = s
+            p2.after_balance(ctx, false)
+            for _, h in ipairs(held) do
+                if h == "127.0.0.1:1981" then
+                    ctx.balancer_server = h
+                    p2.after_balance(ctx, false)
+                end
+            end
+        }
+    }
+--- request
+GET /t
+--- response_body
+127.0.0.1:1980
+
+
+
+=== TEST 3: scale down drops the removed node, remaining nodes balance
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            local up = {resource_key = "/upstreams/lc-scale-down"}
+
+            local nodes = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1}
+            local p1 = least_conn.new(nodes, up, 0)
+
+            local ctx = {}
+            -- fully complete two requests, one per node
+            for _ = 1, 2 do
+                local s = p1.get(ctx)
+                ctx.balancer_server = s
+                p1.after_balance(ctx, false)
+            end
+
+            -- scale down to a single remaining node, picker recreated
+            local scaled = {["127.0.0.1:1981"] = 1}
+            local p2 = least_conn.new(scaled, up, 0)
+
+            local s = p2.get(ctx)
+            ctx.balancer_server = s
+            p2.after_balance(ctx, false)
+            ngx.say(s)
+        }
+    }
+--- request
+GET /t
+--- response_body
+127.0.0.1:1981
+
+
+
+=== TEST 4: connections released on the old picker are seen by the new one
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            local up = {resource_key = "/upstreams/lc-drain-old-picker"}
+
+            local nodes = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1}
+            local p1 = least_conn.new(nodes, up, 0)
+
+            -- 4 long-lived connections, 2 on each node
+            local ctx = {}
+            local held = {}
+            for _ = 1, 4 do
+                held[#held + 1] = p1.get(ctx)
+            end
+
+            -- scale out: new requests are routed with a freshly built picker
+            local scaled = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1,
+                            ["127.0.0.1:1982"] = 1}
+            local p2 = least_conn.new(scaled, up, 0)
+
+            -- the long-lived connections close. They are still bound to the 
picker
+            -- they were routed with (ctx.server_picker), so they are released 
on p1
+            for _, s in ipairs(held) do
+                ctx.balancer_server = s
+                p1.after_balance(ctx, false)
+            end
+
+            -- every node is empty again, so the next requests must spread 
over all
+            -- of them instead of piling up on the node that was added last
+            local picked = {}
+            for _ = 1, 3 do
+                picked[#picked + 1] = p2.get(ctx)
+            end
+            table.sort(picked)
+            ngx.say(table.concat(picked, " "))
+
+            for _, s in ipairs(picked) do
+                ctx.balancer_server = s
+                p2.after_balance(ctx, false)
+            end
+        }
+    }
+--- request
+GET /t
+--- response_body
+127.0.0.1:1980 127.0.0.1:1981 127.0.0.1:1982
+
+
+
+=== TEST 5: node sets of different priorities keep their own state
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            local up = {resource_key = "/upstreams/lc-priority"}
+
+            local high = least_conn.new({["127.0.0.1:1980"] = 1}, up, 1)
+            local low = least_conn.new({["127.0.0.1:1981"] = 1}, up, 0)
+
+            local ctx = {}
+            for _, picker in ipairs({high, low}) do
+                local s = picker.get(ctx)
+                ngx.say(s)
+                ctx.balancer_server = s
+                picker.after_balance(ctx, false)
+            end
+        }
+    }
+--- request
+GET /t
+--- response_body
+127.0.0.1:1980
+127.0.0.1:1981
+
+
+
+=== TEST 6: a request that runs out of servers does not release twice
+--- config
+    location /t {
+        content_by_lua_block {
+            local balancer = require("apisix.balancer")
+            local least_conn = require("apisix.balancer.least_conn")
+
+            local up_conf = {
+                type = "least_conn",
+                resource_key = "/upstreams/lc-double-release",
+                nodes = {
+                    {host = "127.0.0.1", port = 1980, weight = 10, priority = 
0},
+                    {host = "127.0.0.1", port = 1981, weight = 1, priority = 
0},
+                },
+            }
+            local nodes = {["127.0.0.1:1980"] = 10, ["127.0.0.1:1981"] = 1}
+
+            -- three long-lived connections land on 1981 while it is the only 
node
+            local seeded = {}
+            local seed = least_conn.new({["127.0.0.1:1981"] = 1}, up_conf, 0)
+            for i = 1, 3 do
+                seeded[i] = {}
+                seeded[i].balancer_server = seed.get(seeded[i])
+            end
+
+            -- a request that fails on every node, driven through the real 
balancer
+            -- a fresh version, as a real scale out would produce, so the 
balancer
+            -- builds a picker for the two nodes over the state seeded above
+            local ctx = {upstream_conf = up_conf, upstream_version = 
tostring(ngx.now()),
+                         upstream_key = "lc-double-release", var = {}}
+            assert(balancer.pick_server(nil, ctx), "first try")
+            assert(balancer.pick_server(nil, ctx), "retry")
+            local server, err = balancer.pick_server(nil, ctx)
+            ngx.say(server == nil and err or "expected to run out of servers")
+
+            -- the request holds nothing now, so the log phase must release 
nothing
+            ngx.say("holds a server: ", ctx.balancer_server ~= nil)
+            ctx.server_picker.after_balance(ctx, false)
+
+            -- 1981 still holds the three connections, so 1980 (ten times the 
weight)
+            -- must win every pick below. It only ties once 1981 is thought to 
be
+            -- lighter than it is
+            local picked = {}
+            local stolen = false
+            local p = least_conn.new(nodes, up_conf, 0)
+            for i = 1, 35 do
+                picked[i] = {}
+                picked[i].balancer_server = p.get(picked[i])
+                if picked[i].balancer_server == "127.0.0.1:1981" then
+                    stolen = true
+                end
+            end
+            ngx.say(stolen and "stolen" or "kept")
+
+            for _, c in ipairs(picked) do
+                p.after_balance(c, false)
+            end
+            for _, c in ipairs(seeded) do
+                p.after_balance(c, false)
+            end
+        }
+    }
+--- request
+GET /t
+--- response_body
+failed to find valid upstream server, all upstream servers tried
+holds a server: false
+kept
+
+
+
+=== TEST 7: scaling out an upstream with in-flight requests prefers the new 
node
+--- config
+    location /t {
+        content_by_lua_block {
+            local t = require("lib.test_admin").test
+            local http = require "resty.http"
+
+            -- the requests have to outlast the whole test, and the default 
upstream
+            -- read timeout (6s) would cut them short well before that
+            local function set_upstream(nodes)
+                local code, body = t('/apisix/admin/upstreams/1', ngx.HTTP_PUT,
+                     [[{"type": "least_conn",
+                        "timeout": {"connect": 60, "send": 60, "read": 60},
+                        "nodes": ]] .. nodes .. [[}]])
+                assert(code < 300, body)
+            end
+
+            set_upstream([[{"127.0.0.1:1980": 1, "0.0.0.0:1980": 1}]])
+            local code, body = t('/apisix/admin/routes/1', ngx.HTTP_PUT,
+                 [[{"uri": "/mysleep", "upstream_id": "1"}]])
+            assert(code < 300, body)
+
+            local uri = "http://127.0.0.1:"; .. ngx.var.server_port .. 
"/mysleep?seconds="
+            -- let the route and the upstream reach the router before hitting 
them
+            ngx.sleep(1)
+
+            -- four in-flight requests, two on each of the original nodes
+            local threads = {}
+            for i = 1, 4 do
+                threads[i] = assert(ngx.thread.spawn(function ()
+                    http.new():request_uri(uri .. "60")
+                end))
+            end
+            ngx.sleep(1)
+
+            -- scale out while they are still in flight
+            set_upstream([[{"127.0.0.1:1980": 1, "0.0.0.0:1980": 1, 
"127.0.0.2:1980": 1}]])
+            ngx.sleep(1)
+
+            -- the new node holds no connection, so it must take the next ones
+            for _ = 1, 2 do
+                assert(http.new():request_uri(uri .. "0.1"))
+            end
+
+            -- drop the clients. This only kills the client coroutines: the 
requests
+            -- they made stay parked on the upstream, so their counts are 
released
+            -- whenever those finish, not here. That is fine - a repeated run 
rebuilds
+            -- the upstream from one node, which drops the second one from the 
heap
+            -- and brings it back empty
+            for _, th in ipairs(threads) do
+                ngx.thread.kill(th)
+            end
+        }
+    }
+--- request
+GET /t
+--- timeout: 10
+--- grep_error_log eval
+qr/proxy request to \S+/
+--- grep_error_log_out eval
+qr/\A(?:proxy request to (?:127\.0\.0\.1|0\.0\.0\.0):1980\n){4}(?:proxy 
request to 127\.0\.0\.2:1980\n){2}\z/
+
+
+
+=== TEST 8: an upstream with no stable key keeps its state private to the 
picker
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            -- conf_server and ai-proxy-multi build pickers from an upstream 
table
+            -- that carries no resource key, so there is nothing to share 
state on
+            local nodes = {["127.0.0.1:1980"] = 1, ["127.0.0.1:1981"] = 1}
+
+            local p1 = least_conn.new(nodes, {})
+            local ctx1 = {}
+            ctx1.balancer_server = p1.get(ctx1)
+
+            -- a second picker starts empty: it must not see the connection 
held on
+            -- the first one, so it picks the very same node
+            local p2 = least_conn.new(nodes, {})
+            local ctx2 = {}
+            ctx2.balancer_server = p2.get(ctx2)
+            ngx.say(ctx2.balancer_server == ctx1.balancer_server and "private" 
or "shared")
+
+            p1.after_balance(ctx1, false)
+            p2.after_balance(ctx2, false)
+        }
+    }
+--- request
+GET /t
+--- response_body
+private
+
+
+
+=== TEST 9: scaling out with live WebSocket sessions prefers the new node
+--- config
+    location /t {
+        content_by_lua_block {
+            local t = require("lib.test_admin").test
+            local ws_client = require("resty.websocket.client")
+
+            -- an upstream of its own: the balancing state is shared by every 
picker
+            -- built for one upstream, so a test that asserts on it cannot 
reuse the
+            -- upstream another test has been routing connections to
+            local function set_upstream(nodes)
+                local code, body = t('/apisix/admin/upstreams/2', ngx.HTTP_PUT,
+                     [[{"type": "least_conn",
+                        "timeout": {"connect": 60, "send": 60, "read": 60},
+                        "nodes": ]] .. nodes .. [[}]])
+                assert(code < 300, body)
+            end
+
+            set_upstream([[{"127.0.0.1:1980": 1, "0.0.0.0:1980": 1}]])
+            local code, body = t('/apisix/admin/routes/2', ngx.HTTP_PUT,
+                 [[{"uri": "/websocket_hold", "enable_websocket": true, 
"upstream_id": "2"}]])
+            assert(code < 300, body)
+
+            local uri = "ws://127.0.0.1:" .. ngx.var.server_port .. 
"/websocket_hold"
+            ngx.sleep(1)
+
+            -- a session stays in flight until it is closed, which is what 
makes the
+            -- load stick to the nodes that were there before the scale out
+            local function hold()
+                local wb = ws_client:new({timeout = 60000})
+                assert(wb:connect(uri))
+                ngx.sleep(60)
+            end
+
+            local threads = {}
+            for i = 1, 4 do
+                threads[i] = assert(ngx.thread.spawn(hold))
+            end
+            ngx.sleep(1)
+
+            -- scale out while the four sessions are still connected
+            set_upstream([[{"127.0.0.1:1980": 1, "0.0.0.0:1980": 1, 
"127.0.0.2:1980": 1}]])
+            ngx.sleep(1)
+
+            -- the new node carries no session, so it must take the next ones
+            for i = 5, 6 do
+                threads[i] = assert(ngx.thread.spawn(hold))
+            end
+            ngx.sleep(1)
+
+            -- drop the clients. This only kills the client coroutines: the 
requests
+            -- they made stay parked on the upstream, so their counts are 
released
+            -- whenever those finish, not here. That is fine - a repeated run 
rebuilds
+            -- the upstream from one node, which drops the second one from the 
heap
+            -- and brings it back empty
+            for _, th in ipairs(threads) do
+                ngx.thread.kill(th)
+            end
+        }
+    }
+--- request
+GET /t
+--- timeout: 10
+--- grep_error_log eval
+qr/proxy request to \S+/
+--- grep_error_log_out eval
+qr/\A(?:proxy request to (?:127\.0\.0\.1|0\.0\.0\.0):1980\n){4}(?:proxy 
request to 127\.0\.0\.2:1980\n){2}\z/
+
+
+
+=== TEST 10: scaling a single-node upstream out counts the sessions it already 
has
+--- config
+    location /t {
+        content_by_lua_block {
+            local t = require("lib.test_admin").test
+            local http = require "resty.http"
+
+            -- the sessions have to outlast the whole test, and the default 
upstream
+            -- read timeout (6s) would cut them short well before that
+            local function set_upstream(nodes)
+                local code, body = t('/apisix/admin/upstreams/3', ngx.HTTP_PUT,
+                     [[{"type": "least_conn",
+                        "timeout": {"connect": 60, "send": 60, "read": 60},
+                        "nodes": ]] .. nodes .. [[}]])
+                assert(code < 300, body)
+            end
+
+            -- a single node is the state a k8s deployment or a discovery 
service
+            -- starts from, and the requests routed while it was alone still 
have to
+            -- be counted, or the scale out cannot see that it is loaded
+            set_upstream([[{"127.0.0.1:1980": 1}]])
+            -- reuse route 1: two routes cannot both own /mysleep
+            local code, body = t('/apisix/admin/routes/1', ngx.HTTP_PUT,
+                 [[{"uri": "/mysleep", "upstream_id": "3"}]])
+            assert(code < 300, body)
+
+            local uri = "http://127.0.0.1:"; .. ngx.var.server_port .. 
"/mysleep?seconds="
+            ngx.sleep(1)
+
+            local threads = {}
+            for i = 1, 4 do
+                threads[i] = assert(ngx.thread.spawn(function ()
+                    http.new():request_uri(uri .. "60")
+                end))
+            end
+            ngx.sleep(1)
+
+            set_upstream([[{"127.0.0.1:1980": 1, "127.0.0.2:1980": 1}]])
+            ngx.sleep(1)
+
+            -- the lone node holds four requests, the new one holds none
+            for _ = 1, 4 do
+                assert(http.new():request_uri(uri .. "0.1"))
+            end
+
+            -- abort the sessions rather than outwait them: the log phase 
still runs,
+            -- so the counts are released before the next run of this block
+            for _, th in ipairs(threads) do
+                ngx.thread.kill(th)
+            end
+        }
+    }
+--- request
+GET /t
+--- timeout: 10
+--- grep_error_log eval
+qr/proxy request to \S+/
+--- grep_error_log_out eval
+qr/\A(?:proxy request to 127\.0\.0\.1:1980\n){4}(?:proxy request to 
127\.0\.0\.2:1980\n){4}\z/
+
+
+
+=== TEST 11: a weight change is applied without losing the connection count
+--- config
+    location /t {
+        content_by_lua_block {
+            local least_conn = require("apisix.balancer.least_conn")
+            local up = {resource_key = "/upstreams/lc-reweight"}
+
+            local p1 = least_conn.new({["127.0.0.1:1980"] = 1, 
["127.0.0.1:1981"] = 1}, up, 0)
+
+            -- one connection on each node: equal weight, equal score
+            local held = {}
+            for i = 1, 2 do
+                held[i] = {}
+                held[i].balancer_server = p1.get(held[i])
+            end
+
+            -- 1980 is given ten times the weight. It carries the same load as 
1981,
+            -- so it is now the lighter of the two and must win the next picks
+            local p2 = least_conn.new({["127.0.0.1:1980"] = 10, 
["127.0.0.1:1981"] = 1}, up, 0)
+
+            local picked = {}
+            for i = 1, 3 do
+                picked[i] = {}
+                picked[i].balancer_server = p2.get(picked[i])
+                ngx.say(picked[i].balancer_server)
+            end
+
+            for _, c in ipairs(picked) do
+                p2.after_balance(c, false)
+            end
+            for _, c in ipairs(held) do
+                p2.after_balance(c, false)
+            end
+        }
+    }
+--- request
+GET /t
+--- response_body
+127.0.0.1:1980
+127.0.0.1:1980
+127.0.0.1:1980
diff --git a/t/node/upstream-discovery-ipv6.t b/t/node/upstream-discovery-ipv6.t
new file mode 100644
index 000000000..438615a7a
--- /dev/null
+++ b/t/node/upstream-discovery-ipv6.t
@@ -0,0 +1,192 @@
+#
+# 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);
+log_level('info');
+no_root_location();
+no_shuffle();
+
+add_block_preprocessor(sub {
+    my ($block) = @_;
+
+    if (!$block->yaml_config) {
+        $block->set_value("yaml_config", <<_EOC_);
+apisix:
+    node_listen: 1984
+deployment:
+    role: data_plane
+    role_data_plane:
+        config_provider: yaml
+_EOC_
+    }
+
+    # a discovery node reaches the balancer without going through the schema 
check
+    # or the bracketing the Admin API applies, so a bare IPv6 host arrives 
here raw
+    $block->set_value("listen_ipv6", 1);
+
+    if ($block->apisix_yaml) {
+        $block->set_value("apisix_yaml", $block->apisix_yaml . <<_EOC_);
+upstreams:
+    - service_name: mock
+      discovery_type: mock
+      type: least_conn
+      id: 1
+#END
+_EOC_
+    }
+
+    if (!$block->request) {
+        $block->set_value("request", "GET /t");
+    }
+});
+
+run_tests();
+
+__DATA__
+
+=== TEST 1: single bare-IPv6 node from discovery, least_conn
+--- apisix_yaml
+routes:
+  - uris:
+      - /hello
+    upstream_id: 1
+--- config
+    location /t {
+        content_by_lua_block {
+            local discovery = require("apisix.discovery.init").discovery
+            discovery.mock = {
+                nodes = function()
+                    return {
+                        {host = "::1", port = 1980, weight = 1},
+                    }
+                end
+            }
+            local http = require "resty.http"
+            local uri = "http://127.0.0.1:"; .. ngx.var.server_port .. "/hello"
+            local res = assert(http.new():request_uri(uri, {keepalive = 
false}))
+            ngx.say(res.status)
+            ngx.print(res.body)
+        }
+    }
+--- response_body
+200
+hello world
+--- no_error_log
+attempt to concatenate
+
+
+
+=== TEST 2: multiple bare-IPv6 nodes from discovery, least_conn
+--- apisix_yaml
+routes:
+  - uris:
+      - /hello
+    upstream_id: 1
+--- config
+    location /t {
+        content_by_lua_block {
+            local discovery = require("apisix.discovery.init").discovery
+            discovery.mock = {
+                nodes = function()
+                    return {
+                        {host = "::1", port = 1980, weight = 1},
+                        {host = "0:0:0:0:0:0:0:1", port = 1980, weight = 1},
+                    }
+                end
+            }
+            local http = require "resty.http"
+            local uri = "http://127.0.0.1:"; .. ngx.var.server_port .. "/hello"
+            local res = assert(http.new():request_uri(uri, {keepalive = 
false}))
+            ngx.say(res.status)
+            ngx.print(res.body)
+        }
+    }
+--- response_body
+200
+hello world
+--- no_error_log
+attempt to concatenate
+
+
+
+=== TEST 3: a discovery re-fetch of the same bare-IPv6 nodes is not seen as a 
change
+--- apisix_yaml
+routes:
+  - uris:
+      - /hello
+    upstream_id: 1
+--- config
+    location /t {
+        content_by_lua_block {
+            local discovery = require("apisix.discovery.init").discovery
+            discovery.mock = {
+                nodes = function()
+                    return {
+                        {host = "::1", port = 1980, weight = 1},
+                    }
+                end
+            }
+            local http = require "resty.http"
+            local uri = "http://127.0.0.1:"; .. ngx.var.server_port .. "/hello"
+            for _ = 1, 3 do
+                assert(http.new():request_uri(uri, {keepalive = false}))
+            end
+            ngx.say("ok")
+        }
+    }
+--- response_body
+ok
+--- grep_error_log eval
+qr/create_obj_fun\(\): upstream nodes:/
+--- grep_error_log_out
+create_obj_fun(): upstream nodes:
+
+
+
+=== TEST 4: bare-IPv6 node that already has port and priority is still 
bracketed
+--- apisix_yaml
+routes:
+  - uris:
+      - /hello
+    upstream_id: 1
+--- config
+    location /t {
+        content_by_lua_block {
+            -- port and priority both set, so the node skips the port/priority 
fill.
+            -- Only the IPv6 branch of the clone condition can trigger here
+            local discovery = require("apisix.discovery.init").discovery
+            discovery.mock = {
+                nodes = function()
+                    return {
+                        {host = "::1", port = 1980, weight = 1, priority = 0},
+                    }
+                end
+            }
+            local http = require "resty.http"
+            local uri = "http://127.0.0.1:"; .. ngx.var.server_port .. "/hello"
+            local res = assert(http.new():request_uri(uri, {keepalive = 
false}))
+            ngx.say(res.status)
+            ngx.print(res.body)
+        }
+    }
+--- response_body
+200
+hello world
+--- no_error_log
+attempt to concatenate
diff --git a/t/stream-node/least_conn.t b/t/stream-node/least_conn.t
new file mode 100644
index 000000000..15f8747be
--- /dev/null
+++ b/t/stream-node/least_conn.t
@@ -0,0 +1,74 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#     http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+use t::APISIX 'no_plan';
+
+# every repeat opens a new connection against the same worker, so the balancing
+# state built by the previous ones is still there. Without releasing the server
+# in the log phase the counts only grow and the picks start to drift away from
+# the node with the highest weight.
+repeat_each(4);
+log_level('info');
+no_root_location();
+no_shuffle();
+
+add_block_preprocessor(sub {
+    my ($block) = @_;
+
+    if ($block->apisix_yaml && !$block->yaml_config) {
+        my $yaml_config = <<_EOC_;
+apisix:
+    node_listen: 1984
+deployment:
+    role: data_plane
+    role_data_plane:
+        config_provider: yaml
+_EOC_
+
+        $block->set_value("yaml_config", $yaml_config);
+    }
+
+    $block->set_value("stream_enable", 1);
+
+    if (!$block->stream_request) {
+        $block->set_value("stream_request", "mmm");
+    }
+});
+
+run_tests();
+
+__DATA__
+
+=== TEST 1: release the finished connection in the stream log phase
+--- apisix_yaml
+stream_routes:
+  - id: 1
+    upstream:
+        type: least_conn
+        nodes:
+        - host: 127.0.0.1
+          port: 1995
+          weight: 3
+        - host: 127.0.0.2
+          port: 1995
+          weight: 1
+#END
+--- stream_response
+hello world
+--- grep_error_log eval
+qr/proxy request to \S+/
+--- grep_error_log_out
+proxy request to 127.0.0.1:1995

Reply via email to