From 01ac598e39422dad55d8517d1bb3da91a2db5e48 Mon Sep 17 00:00:00 2001 From: dami Date: Sat, 4 Jul 2026 22:43:20 +0800 Subject: [PATCH] Update webstats_common.lua --- plugins/webstats/lua/webstats_common.lua | 458 ++++++++++++++++--------------- 1 file changed, 239 insertions(+), 219 deletions(-) diff --git a/plugins/webstats/lua/webstats_common.lua b/plugins/webstats/lua/webstats_common.lua index 916d6787d..2a64ebd77 100644 --- a/plugins/webstats/lua/webstats_common.lua +++ b/plugins/webstats/lua/webstats_common.lua @@ -377,59 +377,67 @@ end function _M.cronPre(self) self:lock_working('cron_init_stat') - -- 获取当前小时和下一小时的存储键 - local time_key = self:get_store_key() - local time_key_next = self:get_store_key_with_time(ngx.time() + 3600) + local ok, err = pcall(function() + -- 获取当前小时和下一小时的存储键 + local time_key = self:get_store_key() + local time_key_next = self:get_store_key_with_time(ngx.time() + 3600) - -- 需要预创建的统计表 - local wc_stat = {'request_stat', 'client_stat', 'spider_stat'} + -- 需要预创建的统计表 + local wc_stat = {'request_stat', 'client_stat', 'spider_stat'} - -- 遍历所有站点 - for _, site_v in ipairs(sites) do - local input_sn = site_v["name"] + -- 遍历所有站点 + for _, site_v in ipairs(sites) do + local input_sn = site_v["name"] - -- 安全地初始化数据库连接 - local ok, db = pcall(function() return self:initDB(input_sn) end) - if not ok or not db then - self:D("cronPre initDB failed for " .. tostring(input_sn) .. ": " .. tostring(db)) - self:unlock_working('cron_init_stat') - return false - end - - -- 开启事务 - db:exec([[BEGIN TRANSACTION]]) - - local success = true - -- 预创建当前小时和下一小时的统计记录 - for _, ws_v in ipairs(wc_stat) do - if not self:_update_stat_pre(db, ws_v, time_key) then - success = false - break + -- 安全地初始化数据库连接 + local db_ok, db = pcall(function() return self:initDB(input_sn) end) + if not db_ok or not db then + self:D("cronPre initDB failed for " .. tostring(input_sn) .. ": " .. tostring(db)) + return false end - if not self:_update_stat_pre(db, ws_v, time_key_next) then - success = false - break + + -- 开启事务 + db:exec([[BEGIN TRANSACTION]]) + + local success = true + -- 预创建当前小时和下一小时的统计记录 + for _, ws_v in ipairs(wc_stat) do + if not self:_update_stat_pre(db, ws_v, time_key) then + success = false + break + end + if not self:_update_stat_pre(db, ws_v, time_key_next) then + success = false + break + end + end + + -- 提交或回滚事务 + if success then + pcall(function() db:execute([[COMMIT]]) end) + else + pcall(function() db:exec([[ROLLBACK]]) end) + end + + -- 安全关闭数据库 + pcall(function() db:close() end) + + if not success then + return false end end - -- 提交或回滚事务 - if success then - pcall(function() db:execute([[COMMIT]]) end) - else - pcall(function() db:exec([[ROLLBACK]]) end) - end - - -- 安全关闭数据库 - pcall(function() db:close() end) - - if not success then - self:unlock_working('cron_init_stat') - return false - end - end + return true + end) self:unlock_working('cron_init_stat') - return true + + if not ok then + self:D("cronPre error: " .. tostring(err)) + return false + end + + return err end --[[ @@ -484,212 +492,224 @@ function _M.cron(self) local time_key = self:get_store_key() - for _, site_v in ipairs(sites) do - local input_sn = site_v["name"] - if self:is_migrating(input_sn) then - self:unlock_working(cron_key) - return true - end - - if self:is_working('cron_init_stat') then - self:unlock_working(cron_key) - return true - end - - local ok, db = pcall(function() return self:initDB(input_sn) end) - if not ok or not db then - self:D("initDB failed for " .. input_sn .. ": " .. tostring(db)) - self:unlock_working(cron_key) - return true - end - - stat_fields[input_sn] = {} - - dbs[input_sn] = db - self:clean_stats(db, input_sn) - - local ok, stmt = pcall(function() - return db:prepare[[INSERT INTO web_logs( - time, ip, domain, server_name, method, status_code, uri, body_length, - referer, user_agent, protocol, request_time, is_spider, request_headers, ip_list, client_port) - VALUES(:time, :ip, :domain, :server_name, :method, :status_code, :uri, - :body_length, :referer, :user_agent, :protocol, :request_time, :is_spider, - :request_headers, :ip_list, :client_port)]] - end) - - if ok and stmt then - stmts[input_sn] = stmt - db:exec([[BEGIN TRANSACTION]]) - else - -- self:D("prepare statement failed for " .. input_sn .. ": " .. tostring(stmt)) - if db and db:isopen() then - db:close() + local function cleanup_and_unlock() + for _, site_v in ipairs(sites) do + local input_sn = site_v["name"] + local stmt = stmts[input_sn] + if stmt then + pcall(function() stmt:finalize() end) + end + local db = dbs[input_sn] + if db and db:isopen() then + pcall(function() db:close() end) end - dbs[input_sn] = nil - self:unlock_working(cron_key) - return true end + self:unlock_working(cron_key) end - local has_error = false - - for i = 1, llen do - local data = ngx.shared.mw_total:lpop(total_key) - if not data then - break + local ok, err = pcall(function() + for _, site_v in ipairs(sites) do + local input_sn = site_v["name"] + if self:is_migrating(input_sn) then + return + end + + if self:is_working('cron_init_stat') then + return + end + + local db_ok, db = pcall(function() return self:initDB(input_sn) end) + if not db_ok or not db then + self:D("initDB failed for " .. input_sn .. ": " .. tostring(db)) + return + end + + stat_fields[input_sn] = {} + + dbs[input_sn] = db + self:clean_stats(db, input_sn) + + local stmt_ok, stmt = pcall(function() + return db:prepare[[INSERT INTO web_logs( + time, ip, domain, server_name, method, status_code, uri, body_length, + referer, user_agent, protocol, request_time, is_spider, request_headers, ip_list, client_port) + VALUES(:time, :ip, :domain, :server_name, :method, :status_code, :uri, + :body_length, :referer, :user_agent, :protocol, :request_time, :is_spider, + :request_headers, :ip_list, :client_port)]] + end) + + if stmt_ok and stmt then + stmts[input_sn] = stmt + db:exec([[BEGIN TRANSACTION]]) + else + if db and db:isopen() then + db:close() + end + dbs[input_sn] = nil + return + end end - local ok, info = pcall(function() return json.decode(data) end) - if not ok then - self:D("json decode failed: " .. tostring(info)) - ngx.shared.mw_total:rpush(total_key, data) - has_error = true - break - end + for i = 1, llen do + local data = ngx.shared.mw_total:lpop(total_key) + if not data then + break + end - local input_sn = info['server_name'] - local db = dbs[input_sn] - if not db then - ngx.shared.mw_total:rpush(total_key, data) - has_error = true - break - end + local decode_ok, info = pcall(function() return json.decode(data) end) + if not decode_ok then + self:D("json decode failed: " .. tostring(info)) + ngx.shared.mw_total:rpush(total_key, data) + return + end - local stmt = stmts[input_sn] - if not stmt then - ngx.shared.mw_total:rpush(total_key, data) - has_error = true - break - end + local input_sn = info['server_name'] + local db = dbs[input_sn] + if not db then + ngx.shared.mw_total:rpush(total_key, data) + return + end - local insert_ok = self:store_logs_line(db, stmt, input_sn, info) - if not insert_ok then - ngx.shared.mw_total:rpush(total_key, data) - rollback_sites[input_sn] = true - has_error = true - break - end + local stmt = stmts[input_sn] + if not stmt then + ngx.shared.mw_total:rpush(total_key, data) + return + end - local log_kv = info["log_kv"] - local excluded = log_kv['excluded'] - local stat_tmp_fields = info['stat_fields'] + local insert_ok = self:store_logs_line(db, stmt, input_sn, info) + if not insert_ok then + ngx.shared.mw_total:rpush(total_key, data) + rollback_sites[input_sn] = true + return + end - local stat_fields_is = stat_fields[input_sn] - for stf_k, stf_v in pairs(stat_tmp_fields) do - if excluded then - if stf_k == "spider_stat_fields" or stf_k == "client_stat_fields" then - break + local log_kv = info["log_kv"] + local excluded = log_kv['excluded'] + local stat_tmp_fields = info['stat_fields'] + + local stat_fields_is = stat_fields[input_sn] + for stf_k, stf_v in pairs(stat_tmp_fields) do + if excluded then + if stf_k == "spider_stat_fields" or stf_k == "client_stat_fields" then + break + end + end + + local stf_is = stat_fields_is[stf_k] + if not stf_is then + stf_is = {} + stat_fields_is[stf_k] = stf_is + end + + for sv_k, sv_v in pairs(stf_v) do + stf_is[sv_k] = (stf_is[sv_k] or 0) + sv_v end end - local stf_is = stat_fields_is[stf_k] - if not stf_is then - stf_is = {} - stat_fields_is[stf_k] = stf_is - end - - for sv_k, sv_v in pairs(stf_v) do - stf_is[sv_k] = (stf_is[sv_k] or 0) + sv_v + if not excluded then + local ip = log_kv['ip'] + local body_length = log_kv["body_length"] + + local ip_stats_sn = ip_stats[input_sn] + if not ip_stats_sn then + ip_stats_sn = {} + ip_stats[input_sn] = ip_stats_sn + end + + local ip_stat = ip_stats_sn[ip] + if not ip_stat then + ip_stats_sn[ip] = {ip_num = 1, body_length = body_length} + else + ip_stat.ip_num = ip_stat.ip_num + 1 + ip_stat.body_length = ip_stat.body_length + body_length + end + + local url_stats_sn = url_stats[input_sn] + if not url_stats_sn then + url_stats_sn = {} + url_stats[input_sn] = url_stats_sn + end + + local request_uri = log_kv["request_uri"] + local request_uri_md5 = ngx.md5(request_uri) + local url_stat = url_stats_sn[request_uri_md5] + if not url_stat then + url_stats_sn[request_uri_md5] = {url_num = 1, uri = request_uri, body_length = body_length} + else + url_stat.url_num = url_stat.url_num + 1 + url_stat.body_length = url_stat.body_length + body_length + end end end - if not excluded then - local ip = log_kv['ip'] - local body_length = log_kv["body_length"] + local save_day = config['global']["save_day"] + local now_date = os.date("*t") + local save_date_timestamp = os.time{year = now_date.year, month = now_date.month, day = now_date.day - save_day, hour = 0} + local delete_sql = "DELETE FROM web_logs WHERE time<" .. tostring(save_date_timestamp) - local ip_stats_sn = ip_stats[input_sn] - if not ip_stats_sn then - ip_stats_sn = {} - ip_stats[input_sn] = ip_stats_sn + for _, site_v in ipairs(sites) do + local input_sn = site_v["name"] + + local stmt = stmts[input_sn] + if stmt then + pcall(function() stmt:finalize() end) end - local ip_stat = ip_stats_sn[ip] - if not ip_stat then - ip_stats_sn[ip] = {ip_num = 1, body_length = body_length} - else - ip_stat.ip_num = ip_stat.ip_num + 1 - ip_stat.body_length = ip_stat.body_length + body_length - end + local stat_fields_is = stat_fields[input_sn] + local db = dbs[input_sn] - local url_stats_sn = url_stats[input_sn] - if not url_stats_sn then - url_stats_sn = {} - url_stats[input_sn] = url_stats_sn - end - - local request_uri = log_kv["request_uri"] - local request_uri_md5 = ngx.md5(request_uri) - local url_stat = url_stats_sn[request_uri_md5] - if not url_stat then - url_stats_sn[request_uri_md5] = {url_num = 1, uri = request_uri, body_length = body_length} - else - url_stat.url_num = url_stat.url_num + 1 - url_stat.body_length = url_stat.body_length + body_length - end - end - end - - local save_day = config['global']["save_day"] - local now_date = os.date("*t") - local save_date_timestamp = os.time{year = now_date.year, month = now_date.month, day = now_date.day - save_day, hour = 0} - local delete_sql = "DELETE FROM web_logs WHERE time<" .. tostring(save_date_timestamp) - - for _, site_v in ipairs(sites) do - local input_sn = site_v["name"] - - local stmt = stmts[input_sn] - if stmt then - pcall(function() stmt:finalize() end) - end - - local stat_fields_is = stat_fields[input_sn] - local db = dbs[input_sn] - - if db then - if rollback_sites[input_sn] then - pcall(function() db:exec([[ROLLBACK]]) end) - else - for sti_k, sti_v in pairs(stat_fields_is) do - local vkk = "" - for sv_k, sv_v in pairs(sti_v) do - vkk = vkk .. sv_k .. "=" .. sv_k .. "+" .. sv_v .. "," - end - if vkk ~= "" then - vkk = string.sub(vkk, 1, string.len(vkk) - 1) - - if sti_k == 'request_stat_fields' then - self:update_stat(db, "request_stat", time_key, vkk) - elseif sti_k == 'client_stat_fields' then - self:update_stat(db, "client_stat", time_key, vkk) - elseif sti_k == 'spider_stat_fields' then - self:update_stat(db, "spider_stat", time_key, vkk) + if db then + if rollback_sites[input_sn] then + pcall(function() db:exec([[ROLLBACK]]) end) + else + for sti_k, sti_v in pairs(stat_fields_is) do + local vkk = "" + for sv_k, sv_v in pairs(sti_v) do + vkk = vkk .. sv_k .. "=" .. sv_k .. "+" .. sv_v .. "," + end + if vkk ~= "" then + vkk = string.sub(vkk, 1, string.len(vkk) - 1) + + if sti_k == 'request_stat_fields' then + self:update_stat(db, "request_stat", time_key, vkk) + elseif sti_k == 'client_stat_fields' then + self:update_stat(db, "client_stat", time_key, vkk) + elseif sti_k == 'spider_stat_fields' then + self:update_stat(db, "spider_stat", time_key, vkk) + end end end - end - local local_ip_stats = ip_stats[input_sn] - if local_ip_stats then - for ip_addr, ip_val in pairs(local_ip_stats) do - self:update_statistics_ip(db, ip_addr, ip_val.ip_num, ip_val.body_length) + local local_ip_stats = ip_stats[input_sn] + if local_ip_stats then + for ip_addr, ip_val in pairs(local_ip_stats) do + self:update_statistics_ip(db, ip_addr, ip_val.ip_num, ip_val.body_length) + end end - end - local local_url_stats = url_stats[input_sn] - if local_url_stats then - for url_md5, url_val in pairs(local_url_stats) do - self:update_statistics_uri(db, url_val.uri, url_md5, url_val.url_num, url_val.body_length) + local local_url_stats = url_stats[input_sn] + if local_url_stats then + for url_md5, url_val in pairs(local_url_stats) do + self:update_statistics_uri(db, url_val.uri, url_md5, url_val.url_num, url_val.body_length) + end end + + db:exec(delete_sql) + + pcall(function() db:execute([[COMMIT]]) end) end + end - db:exec(delete_sql) - - pcall(function() db:execute([[COMMIT]]) end) + if db and db:isopen() then + pcall(function() db:close() end) end end + end) - if db and db:isopen() then - pcall(function() db:close() end) - end + if not ok then + self:D("cron process error: " .. tostring(err)) + cleanup_and_unlock() + return true end self:unlock_working(cron_key) @@ -943,7 +963,7 @@ end function _M.lock_working(self, sign) local working_key = sign.."_working" - cache:set(working_key, true, 60) + cache:set(working_key, true, 10) end function _M.unlock_working(self, sign)