Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
179 changes: 173 additions & 6 deletions apisix/core/utils.lua
Original file line number Diff line number Diff line change
Expand Up @@ -494,19 +494,186 @@ function _M.check_tls_bool(fields, conf, plugin_name)
end


function _M.set_var_rate_limiting_info(ctx, key, limit, remaining, reset)
local json_escapes = {
['"'] = '\\"',
['\\'] = '\\\\',
['\b'] = '\\b',
['\f'] = '\\f',
['\n'] = '\\n',
['\r'] = '\\r',
['\t'] = '\\t',
}
local function json_escape_char(c)
return json_escapes[c] or str_format("\\u%04x", str_byte(c))
end


local function to_ms(seconds)
return math_floor(seconds * 1000 + 0.5)
end


local function int_or_null(value)
if value == nil then
return "null"
end
return str_format("%d", value)
end


local function ms_or_null(seconds)
if seconds == nil then
return "null"
end
return str_format("%d", to_ms(seconds))
end


local function bool_or_null(value)
if value == nil then
return "null"
end
return value and "true" or "false"
end


-- a non-negative number with `digits` decimals
local function fixed_point(value, digits)
local scale = 10 ^ digits
local scaled = math_floor(value * scale + 0.5)
return str_format("%d.%0" .. digits .. "d", math_floor(scaled / scale), scaled % scale)
end


-- sliding windows are aligned to the clock: the window containing `now` is
-- the id-th one since the epoch, and the previous window's count is weighted
-- by the share of it still inside the sliding range
local function sliding_window_of(now, window_size)
local id = math_floor(now / window_size)
return id, id * window_size, (window_size - now % window_size) / window_size
end


-- handles every case, including the null values, see format_common_detail
-- for the common ones
local function format_detail(detail)
local window_type = detail.window_type
local window_size = detail.window_size
local fields = str_format(',"window_type":"%s","window_size_ms":%d,"decision":"%s"',
window_type, to_ms(window_size), detail.decision)
if detail.decision == "error" then
return fields
end

fields = fields .. str_format(',"cost":%d,"evaluated_at_ms":%d',
detail.cost, to_ms(detail.now))

local delayed = detail.delayed_sync
if delayed then
fields = fields .. str_format(
',"delayed_sync":{"synced_at_ms":%s,"synced_count":%s,"local_delta":%s}',
ms_or_null(delayed.synced_at), int_or_null(delayed.synced_count),
int_or_null(delayed.local_delta))
end

if window_type ~= "sliding" then
local window_end = detail.window_end
return fields .. str_format(
',"current_window":{"start_ms":%s,"end_ms":%s,"count":%s,"created":%s}',
ms_or_null(window_end and window_end - window_size), ms_or_null(window_end),
int_or_null(detail.count), bool_or_null(detail.created))
end

local id, window_start, weight
if detail.window_now then
id, window_start, weight = sliding_window_of(detail.window_now, window_size)
end
local last_count = detail.last_count
return fields .. str_format(
',"current_window":{"id":%s,"start_ms":%s,"end_ms":%s,"count":%s}'
.. ',"previous_window":{"count":%s,"weight":%s,"weighted_count":%s}',
int_or_null(id), ms_or_null(window_start),
ms_or_null(window_start and window_start + window_size),
int_or_null(detail.count), int_or_null(last_count),
weight and fixed_point(weight, 6) or "null",
(weight and last_count) and fixed_point(last_count * weight, 3) or "null")
end


-- The variable is set on every rate limited request, so the common cases,
-- where the counter is not synced with a delay and every value is known, are
-- formatted in a single string.format call and with integers only, which
-- keeps their cost close to the four original fields alone.
local function format_common_detail(key, limit, remaining, reset, detail)
local decision = detail.decision
local count = detail.count
if decision == "error" or detail.delayed_sync or not count then
return nil
end

local window_size = detail.window_size
local window_size_ms = to_ms(window_size)
if detail.window_type ~= "sliding" then
local window_end = detail.window_end
if not window_end then
return nil
end
local end_ms = to_ms(window_end)
return str_format(
'{"rate_limiting_key":"%s","rate_limiting_limit":%d,'
.. '"rate_limiting_remaining":%d,"rate_limiting_reset":%d,'
.. '"window_type":"fixed","window_size_ms":%d,"decision":"%s","cost":%d,'
.. '"evaluated_at_ms":%d,"current_window":{"start_ms":%d,"end_ms":%d,'
.. '"count":%d,"created":%s}}',
key, limit, remaining, reset, window_size_ms, decision, detail.cost,
to_ms(detail.now), end_ms - window_size_ms, end_ms, count,
bool_or_null(detail.created))
end

local window_now = detail.window_now
local last_count = detail.last_count
if not window_now or not last_count then
return nil
end
local id, window_start, weight = sliding_window_of(window_now, window_size)
local start_ms = to_ms(window_start)
local weight_e6 = math_floor(weight * 1000000 + 0.5)
local weighted_e3 = math_floor(last_count * weight * 1000 + 0.5)
return str_format(
'{"rate_limiting_key":"%s","rate_limiting_limit":%d,'
.. '"rate_limiting_remaining":%d,"rate_limiting_reset":%d,'
.. '"window_type":"sliding","window_size_ms":%d,"decision":"%s","cost":%d,'
.. '"evaluated_at_ms":%d,"current_window":{"id":%d,"start_ms":%d,"end_ms":%d,'
.. '"count":%d},"previous_window":{"count":%d,"weight":%d.%06d,'
.. '"weighted_count":%d.%03d}}',
key, limit, remaining, reset, window_size_ms, decision, detail.cost,
to_ms(detail.now), id, start_ms, start_ms + window_size_ms, count, last_count,
math_floor(weight_e6 / 1000000), weight_e6 % 1000000,
math_floor(weighted_e3 / 1000), weighted_e3 % 1000)
end


-- `detail` (optional) describes the window the request was counted in, see
-- the $rate_limiting_info section of the limit-count plugin docs
function _M.set_var_rate_limiting_info(ctx, key, limit, remaining, reset, detail)
if not ctx then
return
end
key = key or ""
-- the key usually comes from a request variable, so escape it to keep
-- the value valid JSON
key = str_gsub(tostring(key or ""), '[%c"\\]', json_escape_char)
limit = limit or 0
remaining = tonumber(remaining) or 0
reset = reset or 0

ctx.var.rate_limiting_info = str_format(
'{"rate_limiting_key":"%s","rate_limiting_limit":%d,'
.. '"rate_limiting_remaining":%d,"rate_limiting_reset":%d}',
key, limit, remaining, reset)
local info = detail and format_common_detail(key, limit, remaining, reset, detail)
if not info then
info = str_format(
'{"rate_limiting_key":"%s","rate_limiting_limit":%d,'
.. '"rate_limiting_remaining":%d,"rate_limiting_reset":%d%s}',
key, limit, remaining, reset, detail and format_detail(detail) or "")
end
ctx.var.rate_limiting_info = info
end


Expand Down
45 changes: 31 additions & 14 deletions apisix/plugins/limit-count/delayed-syncer.lua
Original file line number Diff line number Diff line change
Expand Up @@ -98,12 +98,18 @@ function _M.key_remote_quota(self, key)
end


function _M.sync_to_shm(self, key, remaining, reset, local_delta)
function _M.sync_to_shm(self, key, remaining, reset, local_delta, info)
local quota = {
remaining = remaining,
reset = reset,
sync_at = ngx_now(),
}
-- window details of the synced counter, only reported in $rate_limiting_info
if info then
quota.count = info.count
quota.last_count = info.last_count
quota.window_now = info.now
end

local _, err, quota_json

Expand All @@ -124,6 +130,8 @@ function _M.sync_to_shm(self, key, remaining, reset, local_delta)
core.log.error("incr local delta shm to failed: ", err, ", key: ", key)
return err
end

return nil, quota
end


Expand All @@ -147,7 +155,8 @@ function _M.delayed_sync(self, key, cost, syncer_id)
end

-- wrap the delayed syncer call in a pcall to avoid the lock being held forever
local ok, remaining, reset, err = pcall(self._delayed_sync, self, key, cost, syncer_id)
local ok, remaining, reset, err, info = pcall(self._delayed_sync, self, key, cost,
syncer_id)
if not ok then
err = remaining
remaining = nil
Expand All @@ -159,12 +168,12 @@ function _M.delayed_sync(self, key, cost, syncer_id)
core.log.error("unlock key(" .. key .. ") failed: ", err_unlock)
end

return remaining, reset, err
return remaining, reset, err, info
end


function _M._delayed_sync(self, key, cost, syncer_id)
local _, reset, remote_quota_json
local _, reset, remote_quota_json, info
local local_delta, err = self.shd:get(self:key_local_delta(key))
if err then
return nil, nil, err
Expand Down Expand Up @@ -209,15 +218,15 @@ function _M._delayed_sync(self, key, cost, syncer_id)
-- is at/over the limit; the fixed-window backend has no commit() and its
-- incoming() already increments before reporting "rejected".
local flush = self.limiter.commit or self.limiter.incoming
_, remaining_or_err, reset = flush(self.limiter, key, local_delta)
_, remaining_or_err, reset, info = flush(self.limiter, key, local_delta)
if type(remaining_or_err) ~= "string" then
remote_remaining = remaining_or_err
remote_reset = reset
elseif remaining_or_err ~= "rejected" then
core.log.error("sync to redis failed: ", remaining_or_err, ", key: ", key)
if self.limiter.fallback_limiter then
core.log.warn("try use fallback limiter to do rate limiting")
_, remaining_or_err, reset =
_, remaining_or_err, reset, info =
self.limiter.fallback_limiter:incoming(key, local_delta)
if type(remaining_or_err) ~= "string" then
remote_remaining = remaining_or_err
Expand All @@ -240,7 +249,7 @@ function _M._delayed_sync(self, key, cost, syncer_id)

core.log.info("sync to shm, key: ", key, ", remote_remaining: ", remote_remaining,
", remote_reset: ", remote_reset)
err = self:sync_to_shm(key, remote_remaining, remote_reset, local_delta)
err, quota = self:sync_to_shm(key, remote_remaining, remote_reset, local_delta, info)
if err then
return nil, nil, err
end
Expand Down Expand Up @@ -330,7 +339,15 @@ function _M._delayed_sync(self, key, cost, syncer_id)
end
end

return remaining, reset
return remaining, reset, nil, {
delayed_sync = {
synced_at = quota.sync_at,
synced_count = quota.count,
local_delta = local_delta,
},
last_count = quota.last_count,
window_now = quota.window_now,
}
end


Expand All @@ -342,30 +359,30 @@ local function sync_key(self, key)

if delta then
local flush = self.limiter.commit or self.limiter.incoming
local _, remaining_or_err, reset = flush(self.limiter, key, delta)
local _, remaining_or_err, reset, info = flush(self.limiter, key, delta)
-- compat
if type(remaining_or_err) ~= "string" then
self:sync_to_shm(key, remaining_or_err, reset, delta)
self:sync_to_shm(key, remaining_or_err, reset, delta, info)
elseif remaining_or_err ~= "rejected" then
core.log.error("sync to redis failed: ", remaining_or_err, ", key: ", key)
if self.limiter.fallback_limiter then
core.log.warn("try use fallback limiter to do rate limiting")
if delta < 1 then
delta = 1
end
_, remaining_or_err, reset =
_, remaining_or_err, reset, info =
self.limiter.fallback_limiter:incoming(key, delta)
if type(remaining_or_err) ~= "string" then
self:sync_to_shm(key, remaining_or_err, reset, delta)
self:sync_to_shm(key, remaining_or_err, reset, delta, info)
elseif remaining_or_err ~= "rejected" then
core.log.error("sync to fallback_limiter failed: ",
remaining_or_err, ", key: ", key)
else
self:sync_to_shm(key, 0, reset, delta)
self:sync_to_shm(key, 0, reset, delta, info)
end
end
else
self:sync_to_shm(key, 0, reset, delta)
self:sync_to_shm(key, 0, reset, delta, info)
end
end
end
Expand Down
Loading
Loading