mirror of
https://github.com/midoks/mdserver-web.git
synced 2026-10-10 00:39:25 +08:00
@@ -281,6 +281,7 @@ def status():
|
||||
return 'start'
|
||||
|
||||
|
||||
# ps -ef|grep nginx|awk '{print $2}'|xargs kill
|
||||
def restyOp(method):
|
||||
file = initDreplace()
|
||||
|
||||
@@ -292,6 +293,7 @@ def restyOp(method):
|
||||
|
||||
current_os = mw.getOs()
|
||||
if current_os == "darwin":
|
||||
# print(file + ' ' + method)
|
||||
data = mw.execShell(file + ' ' + method)
|
||||
if data[1] == '':
|
||||
return 'ok'
|
||||
|
||||
@@ -49,10 +49,12 @@ case "$1" in
|
||||
mPID=`cat $PIDFILE`
|
||||
isStart=`ps ax | awk '{ print $1 }' | grep -e "^${mPID}$"`
|
||||
if [ "$isStart" = '' ];then
|
||||
ps -ef|grep nginx|grep -v grep |awk '{print $2}'| xargs kill
|
||||
echo "$NAME is not running."
|
||||
exit 1
|
||||
fi
|
||||
else
|
||||
ps -ef|grep nginx|grep -v grep |awk '{print $2}'| xargs kill
|
||||
echo "$NAME is not running."
|
||||
exit 1
|
||||
fi
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
lua_shared_dict mw_total 100m;
|
||||
lua_need_request_body off;
|
||||
include {$SERVER_APP}/lua/webstats_log.lua;
|
||||
@@ -1,3 +1,3 @@
|
||||
lua_shared_dict mw_total 100m;
|
||||
lua_need_request_body off;
|
||||
include {$SERVER_APP}/lua/webstats_log.lua;
|
||||
log_by_lua_file {$SERVER_APP}/lua/webstats_log.lua;
|
||||
@@ -226,6 +226,14 @@ def makeSiteConfig():
|
||||
return data
|
||||
|
||||
|
||||
def contentReplace(content):
|
||||
root_dir = mw.getServerDir()
|
||||
service_path = root_dir + "/webstats"
|
||||
content = content.replace('{$SERVER_APP}', service_path)
|
||||
content = content.replace('{$ROOT_PATH}', root_dir)
|
||||
return content
|
||||
|
||||
|
||||
def initDreplace():
|
||||
|
||||
service_path = getServerDir()
|
||||
@@ -236,8 +244,7 @@ def initDreplace():
|
||||
path_tpl = getPluginDir() + '/conf/webstats.conf'
|
||||
if not os.path.exists(path):
|
||||
content = mw.readFile(path_tpl)
|
||||
content = content.replace('{$SERVER_APP}', service_path)
|
||||
content = content.replace('{$ROOT_PATH}', mw.getServerDir())
|
||||
content = contentReplace(content)
|
||||
mw.writeFile(path, content)
|
||||
|
||||
# 已经安装的
|
||||
@@ -285,6 +292,7 @@ def start():
|
||||
import tool_task
|
||||
tool_task.createBgTask()
|
||||
|
||||
makeOpDstRunLua()
|
||||
# issues:326
|
||||
luaRestart()
|
||||
return 'ok'
|
||||
@@ -299,6 +307,7 @@ def stop():
|
||||
tool_task.removeBgTask()
|
||||
|
||||
luaRestart()
|
||||
makeOpDstStopLua()
|
||||
return 'ok'
|
||||
|
||||
|
||||
@@ -320,11 +329,40 @@ def reload():
|
||||
loadLuaFileReload(fl)
|
||||
|
||||
loadDebugLogFile()
|
||||
makeOpDstRunLua()
|
||||
|
||||
luaRestart()
|
||||
return 'ok'
|
||||
|
||||
|
||||
def makeOpDstRunLua(conf_reload=False):
|
||||
root_worker_dir = mw.getServerDir() + '/web_conf/nginx/lua/init_worker_by_lua_file'
|
||||
path = getServerDir()
|
||||
path_tpl = getPluginDir()
|
||||
|
||||
worker_dst = root_worker_dir + "/webstats_init_worker.lua"
|
||||
if not os.path.exists(worker_dst) or conf_reload:
|
||||
worker_dst_tpl = path_tpl + "/lua/webstats_init_worker.lua"
|
||||
content = mw.readFile(worker_dst_tpl)
|
||||
content = contentReplace(content)
|
||||
mw.writeFile(worker_dst, content)
|
||||
|
||||
mw.opLuaMakeAll()
|
||||
return True
|
||||
|
||||
|
||||
def makeOpDstStopLua():
|
||||
root_worker_dir = mw.getServerDir() + '/web_conf/nginx/lua/init_worker_by_lua_file'
|
||||
path = getServerDir()
|
||||
path_tpl = getPluginDir()
|
||||
|
||||
worker_dst = root_worker_dir + "/webstats_init_worker.lua"
|
||||
if os.path.exists(worker_dst):
|
||||
os.remove(worker_dst)
|
||||
|
||||
mw.opLuaMakeAll()
|
||||
return True
|
||||
|
||||
def getGlobalConf():
|
||||
conf = getConf()
|
||||
content = mw.readFile(conf)
|
||||
|
||||
@@ -90,9 +90,10 @@ Install_App()
|
||||
BREW_DIR=`which brew`
|
||||
BREW_DIR=${BREW_DIR/\/bin\/brew/}
|
||||
echo "BREW_DIR:"${BREW_DIR}
|
||||
LIB_SQLITE_DIR=/opt/homebrew/opt/sqlite
|
||||
find_cfg=`cat Makefile | grep 'SQLITE_DIR'`
|
||||
if [ "$find_cfg" == "" ];then
|
||||
LIB_SQLITE_DIR=`brew info sqlite | grep ${BREW_DIR}/Cellar/sqlite | cut -d \ -f 1 | awk 'END {print}'`
|
||||
# LIB_SQLITE_DIR=`brew info sqlite | grep /opt/homebrew/opt/sqlite | cut -d \ -f 1 | awk 'END {print}'`
|
||||
echo "LIB_SQLITE_DIR:"${LIB_SQLITE_DIR}
|
||||
sed -i $BAK "s#\$(ROCKSPEC)#\$(ROCKSPEC) SQLITE_DIR=${LIB_SQLITE_DIR}#g" Makefile
|
||||
fi
|
||||
|
||||
@@ -27,7 +27,7 @@
|
||||
]]
|
||||
|
||||
local setmetatable = setmetatable
|
||||
local _M = { _VERSION = '1.0' }
|
||||
local _M = { _VERSION = '0.2.5' }
|
||||
local mt = { __index = _M }
|
||||
|
||||
-- 依赖模块
|
||||
@@ -56,9 +56,6 @@ local day_column = "day" .. number_day
|
||||
-- 数据库列名:flow + 日期(如 flow01)
|
||||
local flow_column = "flow" .. number_day
|
||||
|
||||
-- 日志文件存储目录
|
||||
local log_dir = "{$SERVER_APP}/logs"
|
||||
|
||||
-- 当前站点的配置(模块级变量,供多个函数共享)
|
||||
local auto_config = nil
|
||||
|
||||
@@ -71,6 +68,7 @@ local today = ngx.re.gsub(ngx.today(), '-', '')
|
||||
]]
|
||||
function _M.new(self)
|
||||
local self = {
|
||||
app_dir = "", ---
|
||||
total_key = total_key, -- 共享内存队列键名
|
||||
params = nil, -- 请求参数
|
||||
site_config = nil, -- 站点配置
|
||||
@@ -88,12 +86,23 @@ end
|
||||
function _M.getInstance(self)
|
||||
if self.instance == nil then
|
||||
self.instance = self:new()
|
||||
self:cron() -- 启动定时任务
|
||||
end
|
||||
assert(self.instance ~= nil)
|
||||
return self.instance
|
||||
end
|
||||
|
||||
function _M.start_cron(self)
|
||||
if ngx.worker.id() ~= 0 then
|
||||
return
|
||||
end
|
||||
self:cron()
|
||||
end
|
||||
|
||||
|
||||
function _M.setAppDir(self, dir)
|
||||
self.app_dir = dir
|
||||
end
|
||||
|
||||
|
||||
--[[
|
||||
初始化 SQLite 数据库连接
|
||||
@@ -108,6 +117,11 @@ end
|
||||
- journal_size_limit = 20GB: WAL 文件大小限制
|
||||
]]
|
||||
function _M.initDB(self, input_sn)
|
||||
log_dir = self.app_dir .. "/logs"
|
||||
if log_dir == "" then
|
||||
return nil
|
||||
end
|
||||
|
||||
local path = log_dir .. '/' .. input_sn .. "/logs.db"
|
||||
local db, err = sqlite3.open(path)
|
||||
|
||||
@@ -120,6 +134,7 @@ function _M.initDB(self, input_sn)
|
||||
db:exec([[PRAGMA page_size = 32768]])
|
||||
db:exec([[PRAGMA journal_mode = wal]])
|
||||
db:exec([[PRAGMA journal_size_limit = 21474836480]])
|
||||
db:exec([[PRAGMA busy_timeout = 5000]])
|
||||
return db
|
||||
end
|
||||
|
||||
@@ -185,7 +200,7 @@ end
|
||||
@return string 域名,失败返回 "unknown"
|
||||
]]
|
||||
function _M.get_domain(self)
|
||||
local domain = ngx.req.get_headers()['host']
|
||||
local domain = ngx.var.host
|
||||
if domain == nil then
|
||||
domain = "unknown"
|
||||
end
|
||||
@@ -341,21 +356,18 @@ end
|
||||
]]
|
||||
function _M.get_http_origin(self)
|
||||
local data = ""
|
||||
local headers = ngx.req.get_headers()
|
||||
if not headers then return data end
|
||||
local headers = {
|
||||
user_agent = ngx.var.http_user_agent,
|
||||
referer = ngx.var.http_referer,
|
||||
host = ngx.var.host,
|
||||
x_forwarded_for = ngx.var.http_x_forwarded_for
|
||||
}
|
||||
|
||||
local req_method = ngx.req.get_method()
|
||||
-- 仅记录非 GET 请求的请求体
|
||||
local req_method = ngx.var.request_method
|
||||
if req_method ~= 'GET' then
|
||||
data = ngx.var.request_body
|
||||
if not data then
|
||||
data = ngx.req.get_body_data()
|
||||
end
|
||||
|
||||
if "string" == type(data) then
|
||||
if data and "string" == type(data) then
|
||||
headers["payload"] = data
|
||||
elseif "table" == type(data) then
|
||||
headers["payload"] = table.concat(data, "&")
|
||||
end
|
||||
end
|
||||
return json.encode(headers)
|
||||
@@ -377,59 +389,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 " .. 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
|
||||
|
||||
--[[
|
||||
@@ -451,7 +471,6 @@ end
|
||||
7. 提交事务并关闭连接
|
||||
]]
|
||||
function _M.cron(self)
|
||||
|
||||
--[[
|
||||
定时任务核心处理函数
|
||||
@param premature boolean 是否为定时器提前触发
|
||||
@@ -484,212 +503,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
|
||||
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
|
||||
|
||||
for i = 1, llen do
|
||||
local data = ngx.shared.mw_total:lpop(total_key)
|
||||
if not data then
|
||||
break
|
||||
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
|
||||
if not excluded then
|
||||
local ip = log_kv['ip']
|
||||
local body_length = log_kv["body_length"]
|
||||
|
||||
for sv_k, sv_v in pairs(stf_v) do
|
||||
stf_is[sv_k] = (stf_is[sv_k] or 0) + sv_v
|
||||
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
|
||||
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)
|
||||
|
||||
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 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)
|
||||
@@ -697,13 +728,37 @@ function _M.cron(self)
|
||||
|
||||
|
||||
function timer_every_get_data_try()
|
||||
local presult, err = pcall( function() timer_every_get_data() end)
|
||||
local presult, err = pcall(function() timer_every_get_data() end)
|
||||
if not presult then
|
||||
self:D("debug cron error on :"..tostring(err))
|
||||
self:D("debug cron error on :" .. tostring(err))
|
||||
return true
|
||||
end
|
||||
end
|
||||
|
||||
function timer_every_get_data_try_debug()
|
||||
ngx.update_time()
|
||||
local start_time = ngx.now()
|
||||
local start_llen = ngx.shared.mw_total:llen(total_key)
|
||||
|
||||
local presult, err = pcall(function() timer_every_get_data() end)
|
||||
|
||||
ngx.update_time()
|
||||
local end_time = ngx.now()
|
||||
local end_llen = ngx.shared.mw_total:llen(total_key)
|
||||
local duration = (end_time - start_time) * 1000
|
||||
local processed = start_llen - end_llen
|
||||
|
||||
if not presult then
|
||||
self:D("debug cron error on :" .. tostring(err))
|
||||
return true
|
||||
end
|
||||
|
||||
if start_llen > 0 then
|
||||
self:D(string.format("cron process: took %.2fms, start=%d, end=%d, processed=%d",
|
||||
duration, start_llen, end_llen, processed))
|
||||
end
|
||||
end
|
||||
|
||||
ngx.timer.every(0.5, timer_every_get_data_try)
|
||||
end
|
||||
|
||||
@@ -761,23 +816,14 @@ function _M.statistics_ipc(self, input_sn, ip)
|
||||
end
|
||||
|
||||
function _M.statistics_request(self, ip, is_spider, body_length)
|
||||
-- 计算pv uv
|
||||
local pvc = 0
|
||||
local uvc = 0
|
||||
local request_header = ngx.req.get_headers()
|
||||
if not is_spider and ngx.status == 200 and body_length > 0 then
|
||||
local ua = ''
|
||||
if request_header['user-agent'] then
|
||||
if "table" == type(request_header['user-agent']) then
|
||||
ua = self:to_json(request_header['user-agent'])
|
||||
else
|
||||
ua = request_header['user-agent']
|
||||
end
|
||||
ua = string.lower(ua)
|
||||
end
|
||||
local ua = ngx.var.http_user_agent or ''
|
||||
ua = string.lower(ua)
|
||||
|
||||
pvc = 1
|
||||
if ua then
|
||||
if ua and ua ~= '' then
|
||||
local today = ngx.today()
|
||||
local uv_token = ngx.md5(ip .. ua .. today)
|
||||
if not cache:get(uv_token) then
|
||||
@@ -789,31 +835,23 @@ function _M.statistics_request(self, ip, is_spider, body_length)
|
||||
return pvc, uvc
|
||||
end
|
||||
|
||||
-- 仅计算GET/HTML
|
||||
function _M.statistics_request_old(self, ip, is_spider, body_length)
|
||||
-- 计算pv uv
|
||||
local pvc = 0
|
||||
local uvc = 0
|
||||
local method = ngx.req.get_method()
|
||||
local method = ngx.var.request_method
|
||||
if not is_spider and method == 'GET' and ngx.status == 200 and body_length > 512 then
|
||||
local ua = ''
|
||||
if request_header['user-agent'] then
|
||||
ua = string.lower(request_header['user-agent'])
|
||||
end
|
||||
local ua = ngx.var.http_user_agent or ''
|
||||
ua = string.lower(ua)
|
||||
|
||||
out_header = ngx.resp.get_headers()
|
||||
if out_header['content-type'] then
|
||||
if string.find(out_header['content-type'], 'text/html', 1, true) then
|
||||
pvc = 1
|
||||
if request_header['user-agent'] then
|
||||
if string.find(ua,'mozilla') then
|
||||
local today = ngx.today()
|
||||
local uv_token = ngx.md5(ip .. request_header['user-agent'] .. today)
|
||||
if not cache:get(uv_token) then
|
||||
uvc = 1
|
||||
cache:set(uv_token,1, self:get_end_time())
|
||||
end
|
||||
end
|
||||
local content_type = ngx.var.content_type
|
||||
if content_type and string.find(content_type, 'text/html', 1, true) then
|
||||
pvc = 1
|
||||
if ua and string.find(ua,'mozilla') then
|
||||
local today = ngx.today()
|
||||
local uv_token = ngx.md5(ip .. ua .. today)
|
||||
if not cache:get(uv_token) then
|
||||
uvc = 1
|
||||
cache:set(uv_token,1, self:get_end_time())
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -907,28 +945,23 @@ function _M.update_stat(self, db, stat_table, key, columns)
|
||||
end
|
||||
--------------------- db end ---------------------------
|
||||
|
||||
-- debug func
|
||||
function _M.D(self,msg)
|
||||
if not debug_mode then return true end
|
||||
local fp = io.open('{$SERVER_APP}/debug.log', 'ab')
|
||||
local fp = io.open(self.app_dir..'/debug.log', 'ab')
|
||||
if fp == nil then
|
||||
return nil
|
||||
end
|
||||
local localtime = os.date("%Y-%m-%d %H:%M:%S")
|
||||
if server_name then
|
||||
fp:write(tostring(msg) .. "\n")
|
||||
else
|
||||
fp:write(localtime..":"..tostring(msg) .. "\n")
|
||||
end
|
||||
fp:write(localtime..":"..tostring(msg) .. "\n")
|
||||
fp:flush()
|
||||
fp:close()
|
||||
return true
|
||||
end
|
||||
|
||||
function _M.is_migrating(self, input_sn)
|
||||
local file = io.open("{$SERVER_APP}/migrating", "rb")
|
||||
local file = io.open(self.app_dir.."/migrating", "rb")
|
||||
if file then return true end
|
||||
local file = io.open("{$SERVER_APP}/logs/"..input_sn.."/migrating", "rb")
|
||||
local file = io.open(self.app_dir.."/logs/"..input_sn.."/migrating", "rb")
|
||||
if file then return true end
|
||||
return false
|
||||
end
|
||||
@@ -943,7 +976,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)
|
||||
@@ -976,12 +1009,12 @@ function _M.read_file_body(self, filename)
|
||||
end
|
||||
|
||||
function _M.load_update_day(self, input_sn)
|
||||
local _file = "{$SERVER_APP}/logs/"..input_sn.."/update_day.log"
|
||||
local _file = self.app_dir.."/logs/"..input_sn.."/update_day.log"
|
||||
return self:read_file_body(_file)
|
||||
end
|
||||
|
||||
function _M.write_update_day(self, input_sn)
|
||||
local _file = "{$SERVER_APP}/logs/"..input_sn.."/update_day.log"
|
||||
local _file = self.app_dir.."/logs/"..input_sn.."/update_day.log"
|
||||
return self:write_file(_file, today, "w")
|
||||
end
|
||||
|
||||
@@ -1015,7 +1048,7 @@ function _M.get_update_field(self, field, value)
|
||||
end
|
||||
|
||||
function _M.get_request_time(self)
|
||||
local request_time = math.floor((ngx.now() - ngx.req.start_time()) * 1000)
|
||||
local request_time = math.floor(tonumber(ngx.var.request_time or 0) * 1000)
|
||||
if request_time == 0 then request_time = 1 end
|
||||
return request_time
|
||||
end
|
||||
@@ -1243,23 +1276,21 @@ end
|
||||
function _M.get_client_ip(self)
|
||||
local client_ip = "unknown"
|
||||
if auto_config and auto_config['cdn'] == true then
|
||||
local request_header = ngx.req.get_headers()
|
||||
for _, v in ipairs(auto_config['cdn_headers'] or {}) do
|
||||
if request_header[v] ~= nil and request_header[v] ~= "" then
|
||||
local ip_list = request_header[v]
|
||||
local header_var = "http_" .. string.gsub(string.lower(v), "-", "_")
|
||||
local ip_list = ngx.var[header_var]
|
||||
if ip_list ~= nil and ip_list ~= "" then
|
||||
client_ip = self:split(ip_list, ',')[1]
|
||||
break
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
-- ipv6
|
||||
if type(client_ip) == 'table' then client_ip = "" end
|
||||
if client_ip ~= "unknown" and ngx.re.match(client_ip,"^([a-fA-F0-9]*):") then
|
||||
return client_ip
|
||||
end
|
||||
|
||||
-- ipv4
|
||||
if not ngx.re.match(client_ip,"\\d+\\.\\d+\\.\\d+\\.\\d+") == nil or not self:is_ipaddr(client_ip) then
|
||||
client_ip = ngx.var.remote_addr
|
||||
if client_ip == nil then
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
-- -- 添加 Lua 模块搜索路径
|
||||
local cpath = "{$SERVER_APP}/lua/"
|
||||
if not package.cpath:find(cpath) then
|
||||
package.cpath = cpath .. "?.so;" .. package.cpath
|
||||
end
|
||||
if not package.path:find(cpath) then
|
||||
package.path = cpath .. "?.lua;" .. package.path
|
||||
end
|
||||
|
||||
local app_dir = "{$SERVER_APP}"
|
||||
local __C = require "webstats_common"
|
||||
local C = __C:getInstance()
|
||||
C:setAppDir(app_dir)
|
||||
C:start_cron()
|
||||
@@ -1,394 +1,287 @@
|
||||
log_by_lua_block {
|
||||
--[[
|
||||
webstats_log.lua - WebStats 日志采集模块
|
||||
=========================================
|
||||
--[[
|
||||
webstats_log.lua - WebStats 日志采集模块
|
||||
=========================================
|
||||
|
||||
本模块在 Nginx 的 log_by_lua_block 阶段执行,负责:
|
||||
- 采集请求日志数据
|
||||
- 过滤不需要统计的请求
|
||||
- 识别爬虫和客户端类型
|
||||
- 计算 PV/UV/IP 等指标
|
||||
- 将日志数据发送到共享内存队列
|
||||
本模块在 Nginx 的 log_by_lua_file 阶段执行,负责:
|
||||
- 采集请求日志数据
|
||||
- 过滤不需要统计的请求
|
||||
- 识别爬虫和客户端类型
|
||||
- 计算 PV/UV/IP 等指标
|
||||
- 将日志数据发送到共享内存队列
|
||||
|
||||
执行时机:每次请求完成后(log_by_lua_block 阶段)
|
||||
执行时机:每次请求完成后(log_by_lua_file 阶段)
|
||||
|
||||
主要数据流向:
|
||||
请求完成 -> log_by_lua_block -> 数据采集 -> 过滤处理 -> 爬虫/客户端识别 -> 统计计算 -> 入队
|
||||
|
||||
依赖模块:
|
||||
- cjson: JSON 编解码
|
||||
- webstats_common: 核心工具模块
|
||||
- webstats_config: 全局配置
|
||||
- webstats_sites: 站点配置
|
||||
]]
|
||||
|
||||
-- 添加 Lua 模块搜索路径
|
||||
local cpath = "{$SERVER_APP}/lua/"
|
||||
if not package.cpath:find(cpath) then
|
||||
package.cpath = cpath .. "?.so;" .. package.cpath
|
||||
end
|
||||
if not package.path:find(cpath) then
|
||||
package.path = cpath .. "?.lua;" .. package.path
|
||||
end
|
||||
|
||||
-- 调试模式开关
|
||||
local debug_mode = true
|
||||
|
||||
-- 引入核心模块(单例模式)
|
||||
local __C = require "webstats_common"
|
||||
local C = __C:getInstance()
|
||||
|
||||
-- 依赖模块
|
||||
local json = require "cjson"
|
||||
local config = require "webstats_config"
|
||||
local sites = require "webstats_sites"
|
||||
|
||||
-- 获取站点名称(通过域名匹配)
|
||||
local server_name = C:get_sn(ngx.var.server_name)
|
||||
|
||||
-- 设置配置数据
|
||||
C:setConfData(config, sites)
|
||||
-- 获取合并后的站点配置
|
||||
local auto_config = C:setInputSn(server_name)
|
||||
|
||||
-- 获取请求头部和方法
|
||||
local request_header = ngx.req.get_headers()
|
||||
local method = ngx.req.get_method()
|
||||
|
||||
-- 共享内存字典实例(必须在使用前定义)
|
||||
local cache = ngx.shared.mw_total
|
||||
-- 日志队列键名
|
||||
local total_key = "log_kv_total"
|
||||
|
||||
--[[
|
||||
状态码过滤表
|
||||
需要记录详细信息的状态码(4xx 和 5xx 错误)
|
||||
]]
|
||||
local status_codes_to_log = {
|
||||
["400"] = true, ["401"] = true, ["402"] = true, ["403"] = true, ["404"] = true,
|
||||
["405"] = true, ["406"] = true, ["407"] = true, ["408"] = true, ["409"] = true,
|
||||
["410"] = true, ["411"] = true, ["412"] = true, ["413"] = true, ["414"] = true,
|
||||
["415"] = true, ["416"] = true, ["417"] = true, ["418"] = true, ["421"] = true,
|
||||
["422"] = true, ["423"] = true, ["424"] = true, ["425"] = true, ["426"] = true,
|
||||
["449"] = true, ["451"] = true, ["499"] = true, ["500"] = true, ["501"] = true,
|
||||
["502"] = true, ["503"] = true, ["504"] = true, ["505"] = true, ["506"] = true,
|
||||
["507"] = true, ["509"] = true, ["510"] = true
|
||||
}
|
||||
|
||||
--[[
|
||||
HTTP 方法过滤表
|
||||
需要统计的 HTTP 方法
|
||||
]]
|
||||
local http_methods = {["get"] = true, ["post"] = true, ["put"] = true, ["patch"] = true, ["delete"] = true}
|
||||
|
||||
--[[
|
||||
排除函数模块
|
||||
============
|
||||
|
||||
以下函数用于判断请求是否应该被排除(不统计),包括:
|
||||
- IP 排除:全局和站点级别的排除 IP
|
||||
- 状态码排除:指定状态码的请求不统计
|
||||
- 扩展名排除:指定扩展名的文件不统计
|
||||
- URL 排除:指定 URL 路径不统计
|
||||
注意:log_by_lua* 上下文中部分 API 不可用,需使用 ngx.var.* 变量替代
|
||||
]]
|
||||
|
||||
--[[
|
||||
加载全局排除 IP 列表到缓存
|
||||
将配置中的全局排除 IP 存入共享内存,避免每次请求都读取配置
|
||||
]]
|
||||
local function load_global_exclude_ip()
|
||||
local load_key = "global_exclude_ip_load"
|
||||
local global_exclude_ip = auto_config["exclude_ip"]
|
||||
if global_exclude_ip then
|
||||
for _, _ip in pairs(global_exclude_ip) do
|
||||
if not cache:get("global_exclude_ip_" .. _ip) then
|
||||
cache:set("global_exclude_ip_" .. _ip, true)
|
||||
end
|
||||
end
|
||||
end
|
||||
cache:set(load_key, true)
|
||||
end
|
||||
-- -- 添加 Lua 模块搜索路径
|
||||
local cpath = "{$SERVER_APP}/lua/"
|
||||
if not package.cpath:find(cpath) then
|
||||
package.cpath = cpath .. "?.so;" .. package.cpath
|
||||
end
|
||||
if not package.path:find(cpath) then
|
||||
package.path = cpath .. "?.lua;" .. package.path
|
||||
end
|
||||
|
||||
--[[
|
||||
加载站点级排除 IP 列表到缓存
|
||||
@param input_server_name string 站点名称
|
||||
@return boolean 是否成功
|
||||
]]
|
||||
local function load_exclude_ip(input_server_name)
|
||||
local load_key = input_server_name .. "_exclude_ip_load"
|
||||
local site_config = config[input_server_name]
|
||||
-- 调试模式开关
|
||||
local debug_mode = true
|
||||
local app_dir = "{$SERVER_APP}"
|
||||
|
||||
local site_exclude_ip = nil
|
||||
if site_config then
|
||||
site_exclude_ip = site_config["exclude_ip"]
|
||||
end
|
||||
local function run_app()
|
||||
local __C = require "webstats_common"
|
||||
local C = __C:getInstance()
|
||||
C:setAppDir(app_dir)
|
||||
|
||||
if site_exclude_ip then
|
||||
for _, _ip in pairs(site_exclude_ip) do
|
||||
cache:set(input_server_name .. "_exclude_ip_" .. _ip, true)
|
||||
end
|
||||
end
|
||||
local json = require "cjson"
|
||||
local config = require "webstats_config"
|
||||
local sites = require "webstats_sites"
|
||||
|
||||
cache:set(load_key, true)
|
||||
return true
|
||||
end
|
||||
local server_name = C:get_sn(ngx.var.server_name)
|
||||
|
||||
--[[
|
||||
根据状态码判断是否排除
|
||||
检查当前请求的状态码是否在排除列表中
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function filter_status()
|
||||
if not auto_config['exclude_status'] then return false end
|
||||
local the_status = tostring(ngx.status)
|
||||
for _, v in ipairs(auto_config['exclude_status']) do
|
||||
if the_status == v then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
C:setConfData(config, sites)
|
||||
local auto_config = C:setInputSn(server_name)
|
||||
|
||||
--[[
|
||||
根据文件扩展名判断是否排除
|
||||
检查请求的 URI 是否以排除的扩展名结尾
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function exclude_extension()
|
||||
local uri = ngx.var.uri
|
||||
if not uri then return false end
|
||||
-- 在 log_by_lua* 上下文中使用 ngx.var 替代不可用的 API
|
||||
local method = ngx.var.request_method
|
||||
local user_agent = ngx.var.http_user_agent
|
||||
local x_forwarded_for = ngx.var.http_x_forwarded_for
|
||||
local referer = ngx.var.http_referer
|
||||
local host = ngx.var.host
|
||||
|
||||
local exclude_ext = auto_config['exclude_extension']
|
||||
if not exclude_ext then return false end
|
||||
local cache = ngx.shared.mw_total
|
||||
local total_key = "log_kv_total"
|
||||
|
||||
local lower_uri = string.lower(uri)
|
||||
for _, ext in ipairs(exclude_ext) do
|
||||
if string.find(lower_uri, "." .. ext .. "$", 1, true) then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
local status_codes_to_log = {
|
||||
["400"] = true, ["401"] = true, ["402"] = true, ["403"] = true, ["404"] = true,
|
||||
["405"] = true, ["406"] = true, ["407"] = true, ["408"] = true, ["409"] = true,
|
||||
["410"] = true, ["411"] = true, ["412"] = true, ["413"] = true, ["414"] = true,
|
||||
["415"] = true, ["416"] = true, ["417"] = true, ["418"] = true, ["421"] = true,
|
||||
["422"] = true, ["423"] = true, ["424"] = true, ["425"] = true, ["426"] = true,
|
||||
["449"] = true, ["451"] = true, ["499"] = true, ["500"] = true, ["501"] = true,
|
||||
["502"] = true, ["503"] = true, ["504"] = true, ["505"] = true, ["506"] = true,
|
||||
["507"] = true, ["509"] = true, ["510"] = true
|
||||
}
|
||||
|
||||
--[[
|
||||
根据 URL 判断是否排除
|
||||
支持精确匹配和正则匹配两种模式
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function exclude_url()
|
||||
local request_uri = ngx.var.request_uri
|
||||
if not request_uri then return false end
|
||||
local http_methods = {["get"] = true, ["post"] = true, ["put"] = true, ["patch"] = true, ["delete"] = true}
|
||||
|
||||
local url_conf = auto_config['exclude_url']
|
||||
if not url_conf then return false end
|
||||
local function load_global_exclude_ip()
|
||||
local load_key = "global_exclude_ip_load"
|
||||
local global_exclude_ip = auto_config["exclude_ip"]
|
||||
if global_exclude_ip then
|
||||
for _, _ip in pairs(global_exclude_ip) do
|
||||
if not cache:get("global_exclude_ip_" .. _ip) then
|
||||
cache:set("global_exclude_ip_" .. _ip, true)
|
||||
end
|
||||
end
|
||||
end
|
||||
cache:set(load_key, true)
|
||||
end
|
||||
|
||||
-- 去掉开头的 '/'
|
||||
local the_uri = string.sub(request_uri, 2)
|
||||
for _, conf in ipairs(url_conf) do
|
||||
local mode = conf["mode"]
|
||||
local url = conf["url"]
|
||||
if mode == "regular" then
|
||||
-- 正则匹配模式
|
||||
if ngx.re.find(the_uri, url, "ijo") then
|
||||
return true
|
||||
end
|
||||
else
|
||||
-- 精确匹配模式
|
||||
if the_uri == url then
|
||||
return true
|
||||
end
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
local function load_exclude_ip(input_server_name)
|
||||
local load_key = input_server_name .. "_exclude_ip_load"
|
||||
local site_config = config[input_server_name]
|
||||
|
||||
--[[
|
||||
根据 IP 判断是否排除
|
||||
先检查站点级排除 IP,再检查全局排除 IP
|
||||
@param input_server_name string 站点名称
|
||||
@param ip string 客户端 IP
|
||||
@return boolean 是否排除
|
||||
]]
|
||||
local function exclude_ip(input_server_name, ip)
|
||||
local site_config = config[input_server_name]
|
||||
if site_config then
|
||||
local server_exclude_ips = site_config["exclude_ip"]
|
||||
if server_exclude_ips then
|
||||
for _, _ip in pairs(server_exclude_ips) do
|
||||
if cache:get(input_server_name .. "_exclude_ip_" .. ip) then
|
||||
return true
|
||||
end
|
||||
break
|
||||
end
|
||||
end
|
||||
end
|
||||
local site_exclude_ip = nil
|
||||
if site_config then
|
||||
site_exclude_ip = site_config["exclude_ip"]
|
||||
end
|
||||
|
||||
return cache:get("global_exclude_ip_" .. ip) ~= nil
|
||||
end
|
||||
if site_exclude_ip then
|
||||
for _, _ip in pairs(site_exclude_ip) do
|
||||
cache:set(input_server_name .. "_exclude_ip_" .. _ip, true)
|
||||
end
|
||||
end
|
||||
|
||||
--[[
|
||||
日志采集函数
|
||||
采集请求数据,识别爬虫/客户端,计算统计指标,然后入队
|
||||
cache:set(load_key, true)
|
||||
return true
|
||||
end
|
||||
|
||||
执行流程:
|
||||
1. 获取日志 ID 和客户端 IP
|
||||
2. 判断是否需要排除(状态码/扩展名/URL/IP)
|
||||
3. 构建 IP 列表(含代理链)
|
||||
4. 采集请求数据(URI、状态码、响应大小等)
|
||||
5. 构建日志记录
|
||||
6. 识别爬虫/客户端类型
|
||||
7. 计算 PV/UV/IP 指标
|
||||
8. 将数据入队到共享内存
|
||||
local function filter_status()
|
||||
if not auto_config['exclude_status'] then return false end
|
||||
local the_status = tostring(ngx.status)
|
||||
for _, v in ipairs(auto_config['exclude_status']) do
|
||||
if the_status == v then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
@param input_sn string 站点名称
|
||||
]]
|
||||
local function cache_logs(input_sn)
|
||||
-- 获取日志唯一 ID
|
||||
local new_id = C:get_last_id(input_sn)
|
||||
-- 获取客户端真实 IP(支持 CDN 转发)
|
||||
local ip = C:get_client_ip()
|
||||
-- 判断是否排除(状态码/扩展名/URL/IP)
|
||||
local excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(input_sn, ip)
|
||||
local function exclude_extension()
|
||||
local uri = ngx.var.uri
|
||||
if not uri then return false end
|
||||
|
||||
-- 构建 IP 列表(包含代理链)
|
||||
local ip_list = request_header["x-forwarded-for"]
|
||||
if ip and not ip_list then
|
||||
ip_list = ip
|
||||
end
|
||||
local exclude_ext = auto_config['exclude_extension']
|
||||
if not exclude_ext then return false end
|
||||
|
||||
local remote_addr = ngx.var.remote_addr
|
||||
if "table" == type(ip_list) then
|
||||
ip_list = json.encode(ip_list)
|
||||
end
|
||||
local lower_uri = string.lower(uri)
|
||||
for _, ext in ipairs(exclude_ext) do
|
||||
if string.find(lower_uri, "." .. ext .. "$", 1, true) then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
if remote_addr and not string.find(ip_list, remote_addr, 1, true) then
|
||||
ip_list = ip_list .. "," .. remote_addr
|
||||
end
|
||||
local function exclude_url()
|
||||
local request_uri = ngx.var.request_uri
|
||||
if not request_uri then return false end
|
||||
|
||||
-- 采集请求数据
|
||||
local request_time = C:get_request_time() -- 请求耗时(毫秒)
|
||||
local client_port = ngx.var.remote_port -- 客户端端口
|
||||
local uri = tostring(ngx.var.uri) -- 请求路径
|
||||
local status_code = ngx.status -- 状态码
|
||||
local protocol = ngx.var.server_protocol -- 协议(HTTP/1.1 等)
|
||||
local request_uri = ngx.var.request_uri -- 完整请求路径(含参数)
|
||||
local body_length = C:get_length() -- 响应体大小
|
||||
local domain = C:get_domain() -- 域名
|
||||
local referer = ngx.var.http_referer -- 来源页
|
||||
local user_agent = request_header['user-agent'] -- 用户代理
|
||||
local now_time = ngx.time() -- 当前时间戳
|
||||
local url_conf = auto_config['exclude_url']
|
||||
if not url_conf then return false end
|
||||
|
||||
-- 构建日志记录
|
||||
local kv = {
|
||||
id = new_id,
|
||||
time_key = os.date("%Y%m%d%H", now_time), -- 小时级存储键
|
||||
time = now_time,
|
||||
ip = ip,
|
||||
domain = domain,
|
||||
server_name = input_sn,
|
||||
real_server_name = input_sn,
|
||||
method = method,
|
||||
status_code = status_code,
|
||||
uri = uri,
|
||||
request_uri = request_uri,
|
||||
body_length = body_length,
|
||||
referer = referer,
|
||||
user_agent = user_agent,
|
||||
protocol = protocol,
|
||||
is_spider = 0, -- 是否爬虫
|
||||
request_time = request_time,
|
||||
excluded = excluded, -- 是否排除
|
||||
request_headers = '', -- 请求头(异常时记录)
|
||||
ip_list = ip_list, -- IP 列表(含代理链)
|
||||
client_port = client_port
|
||||
}
|
||||
local the_uri = string.sub(request_uri, 2)
|
||||
for _, conf in ipairs(url_conf) do
|
||||
local mode = conf["mode"]
|
||||
local url = conf["url"]
|
||||
if mode == "regular" then
|
||||
if ngx.re.find(the_uri, url, "ijo") then
|
||||
return true
|
||||
end
|
||||
else
|
||||
if the_uri == url then
|
||||
return true
|
||||
end
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
|
||||
-- 初始化统计字段
|
||||
local request_stat_fields = {req = 1, length = body_length}
|
||||
local spider_stat_fields = {}
|
||||
local client_stat_fields = {}
|
||||
local function exclude_ip(input_server_name, ip)
|
||||
local site_config = config[input_server_name]
|
||||
if site_config then
|
||||
local server_exclude_ips = site_config["exclude_ip"]
|
||||
if server_exclude_ips then
|
||||
for _, _ip in pairs(server_exclude_ips) do
|
||||
if cache:get(input_server_name .. "_exclude_ip_" .. ip) then
|
||||
return true
|
||||
end
|
||||
break
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
-- 非排除请求的额外处理
|
||||
if not excluded then
|
||||
-- 记录异常请求的原始数据(500/POST/403)
|
||||
if status_code == 500 or (method == "POST" and config['global']["record_post_args"] == true) or (status_code == 403 and config['global']["record_get_403_args"] == true) then
|
||||
local ok, data = pcall(function() return C:get_http_origin() end)
|
||||
if ok and data then
|
||||
kv["request_headers"] = data
|
||||
end
|
||||
end
|
||||
return cache:get("global_exclude_ip_" .. ip) ~= nil
|
||||
end
|
||||
|
||||
-- 统计状态码
|
||||
if status_codes_to_log[tostring(status_code)] then
|
||||
request_stat_fields["status_" .. status_code] = 1
|
||||
end
|
||||
local function cache_logs(input_sn)
|
||||
local new_id = C:get_last_id(input_sn)
|
||||
local ip = C:get_client_ip()
|
||||
local excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(input_sn, ip)
|
||||
|
||||
-- 统计 HTTP 方法
|
||||
local lower_method = string.lower(method)
|
||||
if http_methods[lower_method] then
|
||||
request_stat_fields["http_" .. lower_method] = 1
|
||||
end
|
||||
local ip_list = x_forwarded_for
|
||||
if ip and not ip_list then
|
||||
ip_list = ip
|
||||
end
|
||||
|
||||
-- 识别爬虫/客户端
|
||||
local is_spider, request_spider, spider_index = C:match_spider(user_agent)
|
||||
if not is_spider then
|
||||
-- 非爬虫:识别客户端类型,计算 PV/UV/IP
|
||||
client_stat_fields = C:match_client_arr(user_agent)
|
||||
local pvc, uvc = C:statistics_request(ip, is_spider, body_length)
|
||||
local ipc = C:statistics_ipc(input_sn, ip)
|
||||
local remote_addr = ngx.var.remote_addr
|
||||
if "table" == type(ip_list) then
|
||||
ip_list = json.encode(ip_list)
|
||||
end
|
||||
|
||||
if ipc > 0 then request_stat_fields["ip"] = 1 end
|
||||
if uvc > 0 then request_stat_fields["uv"] = 1 end
|
||||
if pvc > 0 then request_stat_fields["pv"] = 1 end
|
||||
else
|
||||
-- 爬虫:记录爬虫类型
|
||||
kv["is_spider"] = spider_index
|
||||
spider_stat_fields[request_spider] = 1
|
||||
request_stat_fields["spider"] = 1
|
||||
end
|
||||
end
|
||||
if remote_addr and not string.find(ip_list, remote_addr, 1, true) then
|
||||
ip_list = ip_list .. "," .. remote_addr
|
||||
end
|
||||
|
||||
-- 构建最终数据结构
|
||||
local data = {
|
||||
server_name = input_sn,
|
||||
stat_fields = {
|
||||
request_stat_fields = request_stat_fields,
|
||||
client_stat_fields = client_stat_fields,
|
||||
spider_stat_fields = spider_stat_fields,
|
||||
},
|
||||
log_kv = kv,
|
||||
}
|
||||
local request_time = C:get_request_time()
|
||||
local client_port = ngx.var.remote_port
|
||||
local uri = tostring(ngx.var.uri)
|
||||
local status_code = ngx.status
|
||||
local protocol = ngx.var.server_protocol
|
||||
local request_uri = ngx.var.request_uri
|
||||
local body_length = C:get_length()
|
||||
local domain = host or "unknown"
|
||||
|
||||
-- 入队到共享内存队列
|
||||
cache:rpush(total_key, json.encode(data))
|
||||
end
|
||||
local kv = {
|
||||
id = new_id,
|
||||
time_key = os.date("%Y%m%d%H", ngx.time()),
|
||||
time = ngx.time(),
|
||||
ip = ip,
|
||||
domain = domain,
|
||||
server_name = input_sn,
|
||||
real_server_name = input_sn,
|
||||
method = method,
|
||||
status_code = status_code,
|
||||
uri = uri,
|
||||
request_uri = request_uri,
|
||||
body_length = body_length,
|
||||
referer = referer,
|
||||
user_agent = user_agent,
|
||||
protocol = protocol,
|
||||
is_spider = 0,
|
||||
request_time = request_time,
|
||||
excluded = excluded,
|
||||
request_headers = '',
|
||||
ip_list = ip_list,
|
||||
client_port = client_port
|
||||
}
|
||||
|
||||
--[[
|
||||
应用入口函数
|
||||
加载排除 IP 列表,然后采集日志
|
||||
]]
|
||||
local function run_app()
|
||||
load_global_exclude_ip()
|
||||
load_exclude_ip(server_name)
|
||||
cache_logs(server_name)
|
||||
end
|
||||
local request_stat_fields = {req = 1, length = body_length}
|
||||
local spider_stat_fields = {}
|
||||
local client_stat_fields = {}
|
||||
|
||||
if not excluded then
|
||||
if status_code == 500 or (method == "POST" and config['global']["record_post_args"] == true) or (status_code == 403 and config['global']["record_get_403_args"] == true) then
|
||||
kv["request_headers"] = json.encode({
|
||||
user_agent = user_agent,
|
||||
referer = referer,
|
||||
host = host
|
||||
})
|
||||
end
|
||||
|
||||
--[[
|
||||
应用入口包装函数(带错误捕获)
|
||||
在调试模式下,使用 pcall 捕获异常并记录日志
|
||||
]]
|
||||
local function run_app_ok()
|
||||
if not debug_mode then return run_app() end
|
||||
if status_codes_to_log[tostring(status_code)] then
|
||||
request_stat_fields["status_" .. status_code] = 1
|
||||
end
|
||||
|
||||
local presult, err = pcall(function() run_app() end)
|
||||
if not presult then
|
||||
C:D("debug error on :" .. tostring(err))
|
||||
return true
|
||||
end
|
||||
end
|
||||
local lower_method = method and string.lower(method) or ""
|
||||
if lower_method ~= "" and http_methods[lower_method] then
|
||||
request_stat_fields["http_" .. lower_method] = 1
|
||||
end
|
||||
|
||||
-- 执行应用
|
||||
return run_app_ok()
|
||||
}
|
||||
local is_spider, request_spider, spider_index = C:match_spider(user_agent)
|
||||
if not is_spider then
|
||||
client_stat_fields = C:match_client_arr(user_agent)
|
||||
local pvc, uvc = C:statistics_request(ip, is_spider, body_length)
|
||||
local ipc = C:statistics_ipc(input_sn, ip)
|
||||
|
||||
if ipc > 0 then request_stat_fields["ip"] = 1 end
|
||||
if uvc > 0 then request_stat_fields["uv"] = 1 end
|
||||
if pvc > 0 then request_stat_fields["pv"] = 1 end
|
||||
else
|
||||
kv["is_spider"] = spider_index
|
||||
spider_stat_fields[request_spider] = 1
|
||||
request_stat_fields["spider"] = 1
|
||||
end
|
||||
end
|
||||
|
||||
local data = {
|
||||
server_name = input_sn,
|
||||
stat_fields = {
|
||||
request_stat_fields = request_stat_fields,
|
||||
client_stat_fields = client_stat_fields,
|
||||
spider_stat_fields = spider_stat_fields,
|
||||
},
|
||||
log_kv = kv,
|
||||
}
|
||||
|
||||
cache:rpush(total_key, json.encode(data))
|
||||
end
|
||||
|
||||
-- C:D("webstats_log run_app start, server_name=" .. tostring(server_name))
|
||||
-- C:D(tostring(total_key) ..":" .. tostring(ngx.shared.mw_total:llen(total_key)))
|
||||
load_global_exclude_ip()
|
||||
load_exclude_ip(server_name)
|
||||
cache_logs(server_name)
|
||||
-- C:D("webstats_log run_app end")
|
||||
end
|
||||
|
||||
local function run_app_ok()
|
||||
-- if not debug_mode then return run_app() end
|
||||
local presult, err = pcall(function() run_app() end)
|
||||
if not presult then
|
||||
C:D("debug error on :" .. tostring(err))
|
||||
return true
|
||||
end
|
||||
end
|
||||
|
||||
return run_app_ok()
|
||||
Reference in New Issue
Block a user