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 88f37dc8c fix(etcd): do not advance the watch revision on a timeout, 
and make the recovery reload cheap (#13721)
88f37dc8c is described below

commit 88f37dc8ca7cc97f545606f723f6d6f7efc3f9b6
Author: AlinsRan <[email protected]>
AuthorDate: Thu Aug 6 08:02:08 2026 +0800

    fix(etcd): do not advance the watch revision on a timeout, and make the 
recovery reload cheap (#13721)
---
 apisix/core/config_etcd.lua |  67 +++++++-----
 t/core/config_etcd.t        | 250 ++++++++++++++++++++++++++++++++++++++++++--
 2 files changed, 280 insertions(+), 37 deletions(-)

diff --git a/apisix/core/config_etcd.lua b/apisix/core/config_etcd.lua
index d7d6f1971..e59dc6b32 100644
--- a/apisix/core/config_etcd.lua
+++ b/apisix/core/config_etcd.lua
@@ -26,6 +26,7 @@ local json         = require("apisix.core.json")
 local etcd_apisix  = require("apisix.core.etcd")
 local core_str     = require("apisix.core.string")
 local new_tab      = require("table.new")
+local nkeys        = require("table.nkeys")
 local inspect      = require("inspect")
 local process      = require("ngx.process")
 local check_schema = require("apisix.core.schema").check
@@ -196,18 +197,9 @@ local function do_run_watch(premature)
     opts.need_cancel = true
     opts.start_revision = watch_ctx.rev
 
-    -- get latest revision
-    local res, err = watch_ctx.cli:readdir(watch_ctx.prefix .. "/phantomkey")
-    if err then
-        log.error("failed to get latest revision, err: ", err)
-    end
-    local latest_rev
-    if res and res.body and res.body.header and res.body.header.revision then
-        latest_rev = tonumber(res.body.header.revision)
-    else
-        log.error("failed to get latest revision, res: ", 
json.delay_encode(res))
-    end
-
+    -- A watch timeout must not advance start_revision: it cannot tell an idle
+    -- prefix from a stream that died silently, and skipping ahead loses the
+    -- events etcd already wrote into that stream. See #13067.
     log.info("restart watchdir: start_revision=", opts.start_revision)
 
     local res_func, err, http_cli = watch_ctx.cli:watchdir(watch_ctx.prefix, 
opts)
@@ -227,12 +219,6 @@ local function do_run_watch(premature)
             then
                 log.error("wait watch event: ", err)
             end
-            if err == "timeout" then
-                if latest_rev and watch_ctx.rev < latest_rev + 1 then
-                    watch_ctx.rev = latest_rev + 1
-                    log.info("etcd watch timeout, upgrade revision to ", 
watch_ctx.rev)
-                end
-            end
             cancel_watch(http_cli)
             break
         end
@@ -560,6 +546,7 @@ end
 local function load_full_data(self, dir_res, headers, prev_values, 
prev_values_hash)
     local err
     local changed = false
+    local prev_keys_still_present = 0
 
     if self.single_item then
         self.values = new_tab(1, 0)
@@ -621,6 +608,22 @@ local function load_full_data(self, dir_res, headers, 
prev_values, prev_values_h
 
         for _, item in ipairs(values) do
             local key = short_key(self, item.key)
+            local prev_item = get_prev_item(prev_values, prev_values_hash, key)
+            if prev_item then
+                prev_keys_still_present = prev_keys_still_present + 1
+            end
+
+            -- Deliberately leaves `changed` alone, so a reload that changed
+            -- nothing does not bump conf_version and rebuild every router.
+            -- Same semantics as sync_data, which re-runs the checker and 
filter
+            -- only for the keys that changed.
+            if prev_item and prev_item.modifiedIndex == item.modifiedIndex then
+                insert_tab(self.values, prev_item)
+                self.values_hash[key] = #self.values
+                self:upgrade_version(item.modifiedIndex)
+                goto continue
+            end
+
             local data_valid = true
             err = nil
             if type(item.value) ~= "table" then
@@ -660,20 +663,26 @@ local function load_full_data(self, dir_res, headers, 
prev_values, prev_values_h
                     self.filter(item)
                 end
 
-            else
-                local prev_item = get_prev_item(prev_values, prev_values_hash, 
key)
-                if prev_item then
-                    -- keep serving with the last valid configuration instead 
of
-                    -- silently dropping the whole item on a full reload, see 
the
-                    -- incremental path in sync_data for the same semantics
-                    log.warn("failed to check item data of [", self.key, "/", 
key,
-                             "], keep the previous configuration, err: ", err)
-                    insert_tab(self.values, prev_item)
-                    self.values_hash[key] = #self.values
-                end
+            elseif prev_item then
+                -- keep serving with the last valid configuration instead of
+                -- silently dropping the whole item on a full reload, see the
+                -- incremental path in sync_data for the same semantics
+                log.warn("failed to check item data of [", self.key, "/", key,
+                         "], keep the previous configuration, err: ", err)
+                insert_tab(self.values, prev_item)
+                self.values_hash[key] = #self.values
             end
 
             self:upgrade_version(item.modifiedIndex)
+
+            ::continue::
+        end
+
+        -- A deletion leaves every surviving key untouched, so it has to be
+        -- detected separately or a reload that only deletes would keep
+        -- serving the removed items.
+        if prev_values_hash and prev_keys_still_present < 
nkeys(prev_values_hash) then
+            changed = true
         end
     end
 
diff --git a/t/core/config_etcd.t b/t/core/config_etcd.t
index a40f425b6..02fbee1a7 100644
--- a/t/core/config_etcd.t
+++ b/t/core/config_etcd.t
@@ -521,7 +521,7 @@ main etcd watcher initialised, revision=
 
 
 
-=== TEST 14: watch revision should be upgraded when timeout occurs
+=== TEST 14: watch revision must not be upgraded when the watch times out
 --- yaml_config
 deployment:
   role: traditional
@@ -538,6 +538,7 @@ nginx_config:
 --- config
     location /t {
         content_by_lua_block {
+            local core = require("apisix.core")
             local etcd = require("resty.etcd")
             local etcd_cli, err = etcd.new({
                 http_host = "http://127.0.0.1:2379";,
@@ -547,7 +548,8 @@ nginx_config:
                 return
             end
             ngx.sleep(2)
-            -- we will assert 4 lines of revision upgrade log because we have 
one worker and one privileged agent
+            -- write outside the watched prefix so that the global revision
+            -- moves on while the watch itself stays idle and keeps timing out
             for i = 1, 2 do
                local _, err = etcd_cli:set("/apache", "apisix")
                if err then
@@ -556,17 +558,50 @@ nginx_config:
                end
                ngx.sleep(1)
             end
-            ngx.say("passed")
+
+            -- The only other assertion here is that a log line is absent, 
which
+            -- would also hold if the watch were broken outright. So check 
delivery.
+            local _, err = etcd_cli:set("/apisix/routes/after-timeout", {
+                id = "after-timeout",
+                uri = "/after-timeout",
+                create_time = 1700000000,
+                update_time = 1700000000,
+                upstream = {type = "roundrobin", nodes = {["127.0.0.1:1980"] = 
1}}
+            })
+            if err then
+                ngx.say("failed to set route: ", err)
+                return
+            end
+            -- polled, not slept: a fixed wait would just guess at CI latency
+            local delivered = false
+            for _ = 1, 25 do
+                ngx.sleep(0.2)
+                local obj = core.config.fetch_created_obj("/routes")
+                for _, item in ipairs(obj and obj.values or {}) do
+                    if item and item.value and item.value.id == 
"after-timeout" then
+                        delivered = true
+                        break
+                    end
+                end
+                if delivered then
+                    break
+                end
+            end
+            ngx.say("update after timeout delivered: ", delivered)
+
+            local _, err = etcd_cli:delete("/apisix/routes/after-timeout")
+            if err then
+                ngx.log(ngx.WARN, "failed to clean up route: ", err)
+            end
         }
     }
+--- timeout: 20
 --- request
 GET /t
 --- response_body
-passed
---- grep_error_log eval
-qr/etcd watch timeout, upgrade revision to/
---- grep_error_log_out eval
-qr/(etcd watch timeout, upgrade revision to\n){2,}/
+update after timeout delivered: true
+--- no_error_log
+etcd watch timeout, upgrade revision to
 
 
 
@@ -862,3 +897,202 @@ GET /t
 invalid new item loaded: false
 --- no_error_log
 keep the previous configuration
+
+
+
+=== TEST 19: a full reload that changes nothing reuses the items and does not 
bump conf_version
+--- timeout: 25
+--- yaml_config
+deployment:
+  role: traditional
+  role_traditional:
+    config_provider: etcd
+  etcd:
+    host:
+      - "http://127.0.0.1:2379";
+    prefix: /apisix
+--- extra_yaml_config
+nginx_config:
+    worker_processes: 1
+--- config
+    location /t {
+        content_by_lua_block {
+            local core = require("apisix.core")
+            local etcd = require("resty.etcd")
+            local etcd_cli, err = etcd.new({
+                http_host = "http://127.0.0.1:2379";,
+            })
+            if not etcd_cli then
+                ngx.say("failed to create etcd client: ", err)
+                return
+            end
+
+            local _, err = etcd_cli:set("/apisix/global_rules/1", {
+                id = "1",
+                create_time = 1700000000,
+                update_time = 1700000000,
+                plugins = {["response-rewrite"] = {headers = {set = {["X-T"] = 
"a"}}}}
+            })
+            if err then
+                ngx.say("failed to set global_rules/1: ", err)
+                return
+            end
+            ngx.sleep(2)
+
+            local obj = core.config.fetch_created_obj("/global_rules")
+            local before_version = obj.conf_version
+
+            -- A reload always builds a fresh `values` array, so losing this 
probe
+            -- proves it ran; the incremental path only mutates elements.
+            obj.values.array_probe = "old"
+            -- The items inside must survive: reusing them is the whole point.
+            for _, item in ipairs(obj.values) do
+                if item and item.value and item.value.id == "1" then
+                    item.reload_probe = "kept"
+                end
+            end
+
+            -- Arm the recovery path taken after `compacted`. sync_data is 
parked
+            -- in waitdir, so the write below is what wakes it: /2 arrives
+            -- incrementally (+1), then the reload runs with /1 and /2 already 
in
+            -- memory at the revisions etcd reports, so it must not bump again.
+            obj.need_reload = true
+            local _, err = etcd_cli:set("/apisix/global_rules/2", {
+                id = "2",
+                create_time = 1700000000,
+                update_time = 1700000000,
+                plugins = {["response-rewrite"] = {headers = {set = {["X-T2"] 
= "b"}}}}
+            })
+            if err then
+                ngx.say("failed to set global_rules/2: ", err)
+                return
+            end
+            ngx.sleep(3)
+
+            local probe_kept = false
+            for _, item in ipairs(obj.values) do
+                if item and item.value and item.value.id == "1" then
+                    probe_kept = (item.reload_probe == "kept")
+                end
+            end
+
+            ngx.say("reload ran: ", obj.values.array_probe == nil)
+            ngx.say("item reused: ", probe_kept)
+            ngx.say("conf_version bumped once, not twice: ",
+                    obj.conf_version == before_version + 1)
+
+            for _, key in ipairs({"/apisix/global_rules/1", 
"/apisix/global_rules/2"}) do
+                local _, del_err = etcd_cli:delete(key)
+                if del_err then
+                    ngx.log(ngx.WARN, "failed to clean up ", key, ": ", 
del_err)
+                end
+            end
+            ngx.sleep(1)
+        }
+    }
+--- request
+GET /t
+--- response_body
+reload ran: true
+item reused: true
+conf_version bumped once, not twice: true
+
+
+
+=== TEST 20: a full reload that only deletes must still bump conf_version
+--- timeout: 25
+--- yaml_config
+deployment:
+  role: traditional
+  role_traditional:
+    config_provider: etcd
+  etcd:
+    host:
+      - "http://127.0.0.1:2379";
+    prefix: /apisix
+--- extra_yaml_config
+nginx_config:
+    worker_processes: 1
+--- config
+    location /t {
+        content_by_lua_block {
+            local core = require("apisix.core")
+            local etcd = require("resty.etcd")
+            local etcd_cli, err = etcd.new({
+                http_host = "http://127.0.0.1:2379";,
+            })
+            if not etcd_cli then
+                ngx.say("failed to create etcd client: ", err)
+                return
+            end
+
+            local _, err = etcd_cli:set("/apisix/global_rules/1", {
+                id = "1",
+                create_time = 1700000000,
+                update_time = 1700000000,
+                plugins = {["response-rewrite"] = {headers = {set = {["X-T"] = 
"a"}}}}
+            })
+            if err then
+                ngx.say("failed to set global_rules/1: ", err)
+                return
+            end
+            ngx.sleep(2)
+
+            local obj = core.config.fetch_created_obj("/global_rules")
+            obj.values.array_probe = "old"
+
+            -- Live in memory, gone from etcd. Every surviving key is untouched
+            -- and therefore reused, so without an explicit deletion check
+            -- conf_version would not move and the routers would keep serving 
it.
+            local ghost = {
+                key = "/apisix/global_rules/ghost",
+                modifiedIndex = 1,
+                value = {id = "ghost", plugins = {}},
+            }
+            core.table.insert(obj.values, ghost)
+            obj.values_hash["ghost"] = #obj.values
+
+            local before_version = obj.conf_version
+
+            -- same wake-up mechanism as TEST 19: /2 arrives incrementally
+            -- (+1), then the reload drops the ghost (+1)
+            obj.need_reload = true
+            local _, err = etcd_cli:set("/apisix/global_rules/2", {
+                id = "2",
+                create_time = 1700000000,
+                update_time = 1700000000,
+                plugins = {["response-rewrite"] = {headers = {set = {["X-T2"] 
= "b"}}}}
+            })
+            if err then
+                ngx.say("failed to set global_rules/2: ", err)
+                return
+            end
+            ngx.sleep(3)
+
+            local found_ghost = false
+            for _, item in ipairs(obj.values) do
+                if item and item.value and item.value.id == "ghost" then
+                    found_ghost = true
+                end
+            end
+
+            ngx.say("reload ran: ", obj.values.array_probe == nil)
+            ngx.say("ghost dropped: ", not found_ghost)
+            ngx.say("conf_version bumped for the deletion: ",
+                    obj.conf_version == before_version + 2)
+
+            for _, key in ipairs({"/apisix/global_rules/1", 
"/apisix/global_rules/2"}) do
+                local _, del_err = etcd_cli:delete(key)
+                if del_err then
+                    ngx.log(ngx.WARN, "failed to clean up ", key, ": ", 
del_err)
+                end
+            end
+            ngx.sleep(1)
+        }
+    }
+--- request
+GET /t
+--- response_body
+reload ran: true
+ghost dropped: true
+conf_version bumped for the deletion: true

Reply via email to