diff --git a/plugins/webstats/bak/webstats_common.lua b/plugins/webstats/bak/webstats_common.lua new file mode 100644 index 000000000..d73c0142c --- /dev/null +++ b/plugins/webstats/bak/webstats_common.lua @@ -0,0 +1,1220 @@ + +local setmetatable = setmetatable +local _M = { _VERSION = '1.0' } +local mt = { __index = _M } + +local json = require "cjson" +local sqlite3 = require "lsqlite3" +local config = require "webstats_config" +local sites = require "webstats_sites" + +local debug_mode = true +local total_key = "log_kv_total" + +local unset_server_name = "unset" +local max_log_id = 99999999999999 +local cache = ngx.shared.mw_total + +local today = ngx.re.gsub(ngx.today(),'-','') +local request_header = ngx.req.get_headers() +local method = ngx.req.get_method() + +local day = os.date("%d") +local number_day = tonumber(day) +local day_column = "day"..number_day +local flow_column = "flow"..number_day +local spider_column = "spider_flow"..number_day + +-- _M.setInputSn | need +local auto_config = nil + +local log_dir = "{$SERVER_APP}/logs" + +function _M.new(self) + local self = { + total_key = total_key, + params = nil, + site_config = nil, + config = nil, + } + -- self.dbs = {} + return setmetatable(self, mt) +end + + +-- function _M.getInstance(self) +-- if rawget(self, "instance") == nil then +-- rawset(self, "instance", self.new()) +-- self.cron() +-- end +-- assert(self.instance ~= nil) +-- return self.instance +-- 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.initDB(self, input_sn) + local path = log_dir .. '/' .. input_sn .. "/logs.db" + db, err = sqlite3.open(path) + + if err then + return nil + end + + db:exec([[PRAGMA synchronous = 0]]) + db:exec([[PRAGMA cache_size = 8000]]) + db:exec([[PRAGMA page_size = 32768]]) + db:exec([[PRAGMA journal_mode = wal]]) + db:exec([[PRAGMA journal_size_limit = 21474836480]]) + return db +end + +function _M.getTotalKey(self) + return self.total_key +end + +function _M.to_json(self, msg) + return json.encode(msg) +end + +function _M.setConfData( self, config, site_config ) + self.config = config + self.site_config = site_config +end + +function _M.setParams( self, params ) + self.params = params +end + +function _M.setInputSn(self, input_sn) + local global_config = config["global"] + if config[input_sn] == nil then + auto_config = global_config + else + auto_config = config[input_sn] + for k, v in pairs(global_config) do + if auto_config[k] == nil then + auto_config[k] = v + end + end + end + return auto_config +end + +function _M.get_domain(self) + local domain = ngx.req.get_headers()['host'] + -- domain = ngx.re.gsub(domain, "_", ".") + if domain == nil then + domain = "unknown" + end + return domain +end + +function _M.split(self, str, reps) + local arr = {} + -- 修复反向代理代过来的数据 + if "table" == type(str) then + return str + end + string.gsub(str,'[^'..reps..']+',function(w) table.insert(arr,w) end) + return arr +end + +function _M.arrlen(self, arr) + if not arr then return 0 end + local count = 0 + for _,v in ipairs(arr) do + count = count + 1 + end + return count +end + +function _M.is_ipaddr(self, client_ip) + local cipn = self:split(client_ip,'.') + if self:arrlen(cipn) < 4 then return false end + for _,v in ipairs({1,2,3,4}) + do + local ipv = tonumber(cipn[v]) + if ipv == nil then return false end + if ipv > 255 or ipv < 0 then return false end + end + return true +end + + +function _M.get_sn(self, input_sn) + local dst_name = cache:get(input_sn) + if dst_name then return dst_name end + + -- self:D(json.encode(sites)) + for _,v in ipairs(sites) + do + if input_sn == v["name"] then + cache:set(input_sn, v['name'], 86400) + return v["name"] + end + -- self:D("get_sn:"..json.encode(v)) + for _,dst_domain in ipairs(v['domains']) + do + if input_sn == dst_domain then + cache:set(input_sn, v['name'], 86400) + return v['name'] + elseif string.find(dst_domain, "*") then + local new_domain = string.gsub(dst_domain, '*', '.*') + if string.find(input_sn, new_domain) then + dst_domain = v['name'] + cache:set(input_sn, dst_domain, 86400) + end + end + end + end + + cache:set(input_sn, unset_server_name, 86400) + return unset_server_name +end + + +function _M.get_store_key(self) + return os.date("%Y%m%d%H", ngx.time()) +end + +function _M.get_store_key_with_time(self, htime) + return os.date("%Y%m%d%H", htime) +end + +function _M.get_length(self) + local clen = ngx.var.body_bytes_sent + if clen == nil then clen = 0 end + return tonumber(clen) +end + +function _M.get_last_id(self, input_sn) + local last_insert_id_key = input_sn .. "_last_id" + local new_id, err = cache:incr(last_insert_id_key, 1, 0) + cache:incr(cache_count_id_key, 1, 0) + if new_id >= max_log_id then + cache:set(last_insert_id_key, 1) + new_id = cache:get(last_insert_id_key) + end + return new_id +end + +function _M.get_http_origin(self) + local data = "" + local headers = ngx.req.get_headers() + if not headers then return data end + local req_method = ngx.req.get_method() + if req_method ~='GET' then + + -- proxy_pass, fastcgi_pass, uwsgi_pass, and scgi_pass + data = ngx.var.request_body + if not data then + data = ngx.req.get_body_data() + end + + -- API disabled in the context of log_by_lua + -- if not data then + -- ngx.req.read_body() + -- data = ngx.req.get_post_args(1000000) + -- end + + if "string" == type(data) then + headers["payload"] = data + end + + if "table" == type(data) then + -- headers = table.concat(headers, data) + headers["payload"] = table.concat(data, "&") + end + end + return json.encode(headers) +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) + + + for site_k, site_v in ipairs(sites) do + local input_sn = site_v["name"] + + local db = self:initDB(input_sn) + + local wc_stat = { + 'request_stat', + 'client_stat', + 'spider_stat' + } + + local v1 = true + local v2 = true + for _,ws_v in pairs(wc_stat) do + v1 = self:_update_stat_pre(db, ws_v, time_key) + v2 = self:_update_stat_pre(db, ws_v, time_key_next) + end + + if db and db:isopen() then + db:execute([[COMMIT]]) + db:close() + end + + if not v1 or not v2 then + return false + end + end + + self:unlock_working('cron_init_stat') + + return true +end + +-- 后台任务 +function _M.cron(self) + + local timer_every_get_data = function (premature) + + local llen, _ = ngx.shared.mw_total:llen(total_key) + -- self:D("PID:"..tostring(ngx.worker.id())..",llen:"..tostring(llen)) + if llen == 0 then + return true + end + + -- self:D("dedebide:cron task is busy!") + local ready_ok = self:cronPre() + if not ready_ok then + -- self:D("cron task is busy!") + return true + end + + local cron_key = 'cron_every_1s' + if self:is_working(cron_key) then + return true + end + + ngx.update_time() + local begin = ngx.now() + + local dbs = {} + local stmts = {} + local stat_fields = {} + local ip_stats = {} + local url_stats = {} + + local time_key = self:get_store_key() + + for site_k, site_v in ipairs(sites) do + local input_sn = site_v["name"] + -- self:D("input_sn:"..input_sn) + -- 迁移合并时不执行 + if self:is_migrating(input_sn) then + return true + end + + -- 初始化统计表时不执行 + if self:is_working('cron_init_stat') then + return true + end + + local db = self:initDB(input_sn) + + stat_fields[input_sn] = {} + + if db then + dbs[input_sn] = db + self:clean_stats(db,input_sn) + + local tmp_stmt = {} + local stmt = 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)]] + tmp_stmt["web_logs"] = stmt + stmts[input_sn] = tmp_stmt + + db:exec([[BEGIN TRANSACTION]]) + end + end + + self:lock_working(cron_key) + + -- 每秒100条 + for i=1,llen do + local data, _ = ngx.shared.mw_total:lpop(total_key) + if not data then + self:unlock_working(cron_key) + break + end + + local info = json.decode(data) + + -- self:D("info:"..json.encode(info)) + local input_sn = info['server_name'] + -- self:D("insert data input_sn:"..input_sn) + local db = dbs[input_sn] + local stat_fields_is = stat_fields[input_sn] + if not db then + ngx.shared.mw_total:rpush(total_key, data) + self:unlock_working(cron_key) + break + end + + local input_stmts = stmts[input_sn]["web_logs"] + if not input_stmts then + ngx.shared.mw_total:rpush(total_key, data) + self:unlock_working(cron_key) + break + end + + local insert_ok = self:store_logs_line(db, input_stmts, input_sn, info) + if not insert_ok then + ngx.shared.mw_total:rpush(total_key, data) + self:unlock_working(cron_key) + break + end + + local excluded = info["log_kv"]['excluded'] + local stat_tmp_fields = info['stat_fields'] + + -- 合并统计数据 + 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 + + if not stat_fields_is[stf_k] then + stat_fields_is[stf_k] = {} + end + + for sv_k,sv_v in pairs(stf_v) do + if not stat_fields_is[stf_k][sv_k] then + stat_fields_is[stf_k][sv_k] = sv_v + else + stat_fields_is[stf_k][sv_k] = stat_fields_is[stf_k][sv_k] + 1 + end + end + end + + stat_fields[input_sn] = stat_fields_is + -- ip 统计合并 + -- url 统计 + if not excluded then + local ip = info["log_kv"]['ip'] + local body_length = info["log_kv"]["body_length"] + + if not ip_stats[input_sn] then + ip_stats[input_sn] = {} + end + + if not ip_stats[input_sn][ip] then + local tmp = { + ip_num=1, + body_length=body_length + } + ip_stats[input_sn][ip] = tmp + else + ip_stats[input_sn][ip]["ip_num"] = ip_stats[input_sn][ip]["ip_num"]+1 + ip_stats[input_sn][ip]["body_length"] = ip_stats[input_sn][ip]["body_length"]+body_length + end + + + -- uri统计 + if not url_stats[input_sn] then + url_stats[input_sn] = {} + end + local request_uri = info["log_kv"]["request_uri"] + local request_uri_md5 = ngx.md5(request_uri) + if not url_stats[input_sn][request_uri_md5] then + local tmp = { + url_num=1, + uri=request_uri, + body_length=body_length + } + url_stats[input_sn][request_uri_md5] = tmp + else + url_stats[input_sn][request_uri_md5]["url_num"] = url_stats[input_sn][request_uri_md5]["url_num"]+1 + url_stats[input_sn][request_uri_md5]["body_length"] = url_stats[input_sn][request_uri_md5]["body_length"]+body_length + end + end + -- self:D("url_stats:.."..json.encode(url_stats)) + end + + for site_k, site_v in ipairs(sites) do + local input_sn = site_v["name"] + + if stmts[input_sn] then + for stmts_k,stmts_v in pairs(stmts[input_sn]) do + -- self:D("stmts_k:"..tostring(stmts_k)) + local res, err = stmts_v:finalize() + if tostring(res) == "5" then + self:D(stmts_k..":Finalize res:"..tostring(res)..",Finalize err:"..tostring(err)) + end + end + end + + + local stat_fields_is = stat_fields[input_sn] + local db = dbs[input_sn] + local local_ip_stats = ip_stats[input_sn] + local local_url_stats = url_stats[input_sn] + + if db then + -- 统计【spider_stat,client_stat,request_stat】 + 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) + end + if sti_k == 'request_stat_fields' and vkk ~= '' then + self:update_stat( db, "request_stat", time_key, vkk) + end + + if sti_k == 'client_stat_fields' and vkk ~= '' then + self:update_stat( db, "client_stat", time_key, vkk) + end + + if sti_k == 'spider_stat_fields' and vkk ~= '' then + self:update_stat( db, "spider_stat", time_key, vkk) + end + end + + -- ip统计 + 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 + + -- url统计 + if local_url_stats then + for url_md5,url_val in pairs(local_url_stats) do + self:update_statistics_uri( db, tostring(url_val["uri"]),url_md5, url_val["url_num"], url_val["body_length"]) + end + end + + -- delete expire data + local now_date = os.date("*t") + local save_day = config['global']["save_day"] + local save_date_timestamp = os.time{year=now_date.year, + month=now_date.month, day=now_date.day-save_day, hour=0} + + db:exec("DELETE FROM web_logs WHERE time<"..tostring(save_date_timestamp)) + end + + if db and db:isopen() then + db:execute([[COMMIT]]) + db:close() + end + end + + self:unlock_working(cron_key) + ngx.update_time() + -- self:D("PID:"..tostring(ngx.worker.id()).."--【"..tostring(llen).."】, elapsed: " .. tostring(ngx.now() - begin)) + end + + + function timer_every_get_data_try() + local presult, err = pcall( function() timer_every_get_data() end) + if not presult then + self:D("debug cron error on :"..tostring(err)) + return true + end + end + + ngx.timer.every(0.5, timer_every_get_data_try) +end + + +function _M.store_logs_line(self, db, stmt, input_sn, info) + local logline = info['log_kv'] + + local time = logline["time"] + local id = logline["id"] + local protocol = logline["protocol"] + local client_port = logline["client_port"] + local status_code = logline["status_code"] + local uri = logline["uri"] + local request_uri = logline["request_uri"] + local method = logline["method"] + local body_length = logline["body_length"] + local referer = logline["referer"] + local ip = logline["ip"] + local ip_list = logline["ip_list"] + local request_time = logline["request_time"] + local is_spider = logline["is_spider"] + local domain = logline["domain"] + local server_name = logline["server_name"] + local user_agent = logline["user_agent"] + local request_headers = logline["request_headers"] + local excluded = logline["excluded"] + local time_key = logline["time_key"] + + if "table" == type(user_agent) then + user_agent = self:to_json(user_agent) + end + + if not excluded then + stmt:bind_names { + time=time, + ip=ip, + domain=domain, + server_name=server_name, + method=method, + status_code=status_code, + uri=request_uri, + body_length=body_length, + referer=referer, + user_agent=user_agent, + protocol=protocol, + request_time=request_time, + is_spider=is_spider, + request_headers=request_headers, + ip_list=ip_list, + client_port=client_port, + } + + local res, err = stmt:step() + if tostring(res) == "5" then + -- self:D("json:"..json.encode(logline)) + -- self:D("the step database connection is busy, so it will be stored later | step res:"..tostring(res) ..",step err:"..tostring(err)) + if stmt then + stmt:reset() + end + return false + end + stmt:reset() + end + return true +end + +function _M.statistics_ipc(self, input_sn, ip) + -- 判断IP是否重复的时间限定范围是请求的当前时间+24小时 + local ipc = 0 + local ip_token = input_sn..'_'..ip + if not cache:get(ip_token) then + ipc = 1 + cache:set(ip_token,1, self:get_end_time()) + end + return ipc +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 + + pvc = 1 + if ua 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 + 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() + 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 + + 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 + end + end + end + end + return pvc, uvc +end + +--------------------- db start --------------------------- + + +function _M.update_statistics_uri(self, db, uri, uri_md5, day_num, body_length) + -- count the number of URI requests and traffic + local open_statistics_uri = config['global']["statistics_uri"] + if not open_statistics_uri then return true end + + local stat_sql = nil + stat_sql = "INSERT INTO uri_stat(uri_md5,uri) SELECT \""..uri_md5.."\",\""..uri.."\" WHERE NOT EXISTS (SELECT uri_md5 FROM uri_stat WHERE uri_md5=\""..uri_md5.."\");" + local res, err = db:exec(stat_sql) + + stat_sql = "UPDATE uri_stat SET "..day_column.."="..day_column.."+"..day_num..","..flow_column.."="..flow_column.."+"..body_length.." WHERE uri_md5=\""..uri_md5.."\";" + local res, err = db:exec(stat_sql) + return true +end + +function _M.statistics_uri(self, db, uri, uri_md5, body_length) + -- count the number of URI requests and traffic + local open_statistics_uri = config['global']["statistics_uri"] + if not open_statistics_uri then return true end + + local stat_sql = nil + stat_sql = "INSERT INTO uri_stat(uri_md5,uri) SELECT \""..uri_md5.."\",\""..uri.."\" WHERE NOT EXISTS (SELECT uri_md5 FROM uri_stat WHERE uri_md5=\""..uri_md5.."\");" + local res, err = db:exec(stat_sql) + + stat_sql = "UPDATE uri_stat SET "..day_column.."="..day_column.."+1,"..flow_column.."="..flow_column.."+"..body_length.." WHERE uri_md5=\""..uri_md5.."\";" + local res, err = db:exec(stat_sql) + return true +end + + +function _M.update_statistics_ip(self, db, ip, day_num ,body_length) + local open_statistics_ip = config['global']["statistics_ip"] + if not open_statistics_ip then return true end + + local stat_sql = nil + stat_sql = "INSERT INTO ip_stat(ip) SELECT \""..ip.."\" WHERE NOT EXISTS (SELECT ip FROM ip_stat WHERE ip=\""..ip.."\");" + local res, err = db:exec(stat_sql) + + stat_sql = "UPDATE ip_stat SET "..day_column.."="..day_column.."+"..day_num..","..flow_column.."="..flow_column.."+"..body_length.." WHERE ip=\""..ip.."\"" + local res, err = db:exec(stat_sql) + return true +end + +function _M.statistics_ip(self, db, ip, body_length) + local open_statistics_ip = config['global']["statistics_ip"] + if not open_statistics_ip then return true end + + local stat_sql = nil + stat_sql = "INSERT INTO ip_stat(ip) SELECT \""..ip.."\" WHERE NOT EXISTS (SELECT ip FROM ip_stat WHERE ip=\""..ip.."\");" + local res, err = db:exec(stat_sql) + + stat_sql = "UPDATE ip_stat SET "..day_column.."="..day_column.."+1,"..flow_column.."="..flow_column.."+"..body_length.." WHERE ip=\""..ip.."\"" + local res, err = db:exec(stat_sql) + return true +end + + +function _M._update_stat_pre(self, db, stat_table, key) + local local_sql = string.format("INSERT INTO %s(time) SELECT :time WHERE NOT EXISTS(SELECT time FROM %s WHERE time=:time);", stat_table, stat_table) + local update_stat_stmt = db:prepare(local_sql) + + if update_stat_stmt then + update_stat_stmt:bind_names{time=key} + update_stat_stmt:step() + update_stat_stmt:finalize() + return true + end + return false +end + +function _M.update_stat_quick(self, db, stat_table, key,columns) + if not columns then return end + local update_sql = "UPDATE ".. stat_table .. " SET " .. columns .. " WHERE time=" .. key + return db:exec(update_sql) +end + + +function _M.update_stat(self, db, stat_table, key, columns) + -- 根据指定表名,更新统计数据 + if not columns then return end + local local_sql = string.format("INSERT INTO %s(time) SELECT :time WHERE NOT EXISTS(SELECT time FROM %s WHERE time=:time);", stat_table, stat_table) + local stmt = db:prepare(local_sql) + stmt:bind_names{time=key} + local res, err = stmt:step() + stmt:finalize() + local update_sql = "UPDATE ".. stat_table .. " SET " .. columns .. " WHERE time=" .. key + return db:exec(update_sql) +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') + 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:flush() + fp:close() + return true +end + +function _M.is_migrating(self, input_sn) + local file = io.open("{$SERVER_APP}/migrating", "rb") + if file then return true end + local file = io.open("{$SERVER_APP}/logs/"..input_sn.."/migrating", "rb") + if file then return true end + return false +end + +function _M.is_working(self,sign) + local work_status = cache:get(sign.."_working") + if work_status ~= nil and work_status == true then + return true + end + return false +end + +function _M.lock_working(self, sign) + local working_key = sign.."_working" + cache:set(working_key, true, 60) +end + +function _M.unlock_working(self, sign) + local working_key = sign.."_working" + cache:set(working_key, false) +end + +function _M.write_file(self, filename, body, mode) + local fp = io.open(filename, mode) + if fp == nil then + return nil + end + fp:write(body) + fp:flush() + fp:close() + return true + end + +function _M.read_file_body(self, filename) + local fp = io.open(filename,'rb') + if not fp then + return nil + end + fbody = fp:read("*a") + fp:close() + if fbody == '' then + return nil + end + return fbody +end + +function _M.load_update_day(self, input_sn) + local _file = "{$SERVER_APP}/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" + return self:write_file(_file, today, "w") +end + +function _M.clean_stats(self, db, input_sn) + -- 清空 uri,ip 汇总的信息[昨日] + local update_day = self:load_update_day(input_sn) + if not update_day or update_day ~= today then + + local update_sql = "UPDATE uri_stat SET "..day_column.."=0,"..flow_column.."=0" + db:exec(update_sql) + + update_sql = "UPDATE ip_stat SET "..day_column.."=0,"..flow_column.."=0" + db:exec(update_sql) + self:write_update_day(input_sn) + end +end + +function _M.lpop(self) + local cache = ngx.shared.mw_total + return cache:lpop(total_key) +end + +function _M.rpop(self) + local cache = ngx.shared.mw_total + return cache:rpop(total_key) +end + + +function _M.get_update_field(self, field, value) + return field.."="..field.."+"..tostring(value) +end + +function _M.get_request_time(self) + local request_time = math.floor((ngx.now() - ngx.req.start_time()) * 1000) + if request_time == 0 then request_time = 1 end + return request_time +end + + +function _M.get_end_time(self) + local s_time = ngx.time() + local n_date = os.date("*t",s_time + 86400) + n_date.hour = 0 + n_date.min = 0 + n_date.sec = 0 + local d_time = ngx.time(n_date) + return d_time - s_time +end + + +function _M.match_spider(self, ua) + -- 匹配蜘蛛请求 + local is_spider = false + local spider_name = "" + local spider_match = "" + + local spider_table = { + ["baidu"] = 1, -- check + ["bing"] = 2, -- check + ["qh360"] = 3, -- check + ["google"] = 4, + ["bytes"] = 5, -- check + ["sogou"] = 6, -- check + ["youdao"] = 7, + ["soso"] = 8, + ["dnspod"] = 9, + ["yandex"] = 10, + ["yisou"] = 11, + ["other"] = 12, + ["mpcrawler"] = 13, + ["yahoo"] = 14, -- check + ["duckduckgo"] = 15 + } + + local find_spider, _ = ngx.re.match(ua, "(Baiduspider|Bytespider|360Spider|Sogou web spider|Sosospider|Googlebot|bingbot|AdsBot-Google|Google-Adwords|YoudaoBot|Yandex|DNSPod-Monitor|YisouSpider|mpcrawler)", "ijo") + if find_spider then + is_spider = true + spider_match = string.lower(find_spider[0]) + if string.find(spider_match, "baidu", 1, true) then + spider_name = "baidu" + elseif string.find(spider_match, "bytes", 1, true) then + spider_name = "bytes" + elseif string.find(spider_match, "360", 1, true) then + spider_name = "qh360" + elseif string.find(spider_match, "sogou", 1, true) then + spider_name = "sogou" + elseif string.find(spider_match, "soso", 1, true) then + spider_name = "soso" + elseif string.find(spider_match, "google", 1, true) then + spider_name = "google" + elseif string.find(spider_match, "bingbot", 1, true) then + spider_name = "bing" + elseif string.find(spider_match, "youdao", 1, true) then + spider_name = "youdao" + elseif string.find(spider_match, "dnspod", 1, true) then + spider_name = "dnspod" + elseif string.find(spider_match, "yandex", 1, true) then + spider_name = "yandex" + elseif string.find(spider_match, "yisou", 1, true) then + spider_name = "yisou" + elseif string.find(spider_match, "mpcrawler", 1, true) then + spider_name = "mpcrawler" + end + end + + if is_spider then + return is_spider, spider_name, spider_table[spider_name] + end + + -- Curl|Yahoo|HeadlessChrome|包含bot|Wget|Spider|Crawler|Scrapy|zgrab|python|java|Adsbot|DuckDuckGo + find_spider, _ = ngx.re.match(ua, "(Yahoo|Slurp|DuckDuckGo|Semrush|Spider|Bot|crawler)", "ijo") + if find_spider then + spider_match = string.lower(find_spider[0]) + if string.find(spider_match, "yahoo", 1, true) then + spider_name = "other" + elseif string.find(spider_match, "slurp", 1, true) then + spider_name = "other" + elseif string.find(spider_match, "duckduckgo", 1, true) then + spider_name = "other" + elseif string.find(spider_match, "semrush", 1, true) then + spider_name = "other" + elseif string.find(spider_match, "spider", 1, true) then + spider_name = "other" + elseif string.find(spider_match, "bot", 1, true) then + spider_name = "other" + elseif string.find(spider_match, "crawler", 1, true) then + spider_name = "other" + end + return true, spider_name, spider_table[spider_name] + end + return false, "", 0 +end + + +function _M.match_client(self, ua) + local client_stat_fields = "" + + if not ua then + return client_stat_fields + end + + local clients_map = { + ["android"] = "android", + ["iphone"] = "iphone", + ["ipod"] = "iphone", + ["ipad"] = "iphone", + ["firefox"] = "firefox", + ["msie"] = "msie", + ["trident"] = "msie", + ["360se"] = "qh360", + ["360ee"] = "qh360", + ["360browser"] = "qh360", + ["qihoo"] = "qh360", + ["the world"] = "theworld", + ["theworld"] = "theworld", + ["tencenttraveler"] = "tt", + ["maxthon"] = "maxthon", + ["opera"] = "opera", + ["qqbrowser"] = "qq", + ["ucweb"] = "uc", + ["ubrowser"] = "uc", + ["safari"] = "safari", + ["chrome"] = "chrome", + ["metasr"] = "metasr", + ["2345explorer"] = "pc2345", + ["edge"] = "edeg", + ["edg"] = "edeg", + ["windows"] = "windows", + ["linux"] = "linux", + ["macintosh"] = "mac", + ["mobile"] = "mobile" + } + local mobile_regx = "(Mobile|Android|iPhone|iPod|iPad)" + local mobile_res = ngx.re.match(ua, mobile_regx, "ijo") + --mobile + if mobile_res then + client_stat_fields = client_stat_fields..","..self:get_update_field("mobile", 1) + mobile_res = string.lower(mobile_res[0]) + if mobile_res ~= "mobile" then + client_stat_fields = client_stat_fields..","..self:get_update_field(clients_map[mobile_res], 1) + end + else + --pc + -- 匹配结果的顺序,与ua中关键词的顺序有关 + -- lua的正则不支持|语法 + -- 短字符串string.find效率要比ngx正则高 + local pc_regx1 = "(360SE|360EE|360browser|Qihoo|TheWorld|TencentTraveler|Maxthon|Opera|QQBrowser|UCWEB|UBrowser|MetaSr|2345Explorer|Edg[e]*)" + local pc_res = ngx.re.match(ua, pc_regx1, "ijo") + local cls_pc = nil + if not pc_res then + if ngx.re.find(ua, "[Ff]irefox") then + cls_pc = "firefox" + elseif string.find(ua, "MSIE") or string.find(ua, "Trident") then + cls_pc = "msie" + elseif string.find(ua, "[Cc]hrome") then + cls_pc = "chrome" + elseif string.find(ua, "[Ss]afari") then + cls_pc = "safari" + end + else + cls_pc = string.lower(pc_res[0]) + end + -- self::D("UA:"..ua) + -- D("PC cls:"..tostring(cls_pc)) + if cls_pc then + client_stat_fields = client_stat_fields..","..self:get_update_field(clients_map[cls_pc], 1) + else + -- machine and other + local machine_res, err = ngx.re.match(ua, "(ApacheBench|[Cc]url|HeadlessChrome|[a-zA-Z]+[Bb]ot|[Ww]get|[Ss]pider|[Cc]rawler|[Ss]crapy|zgrab|[Pp]ython|java)", "ijo") + if machine_res then + client_stat_fields = client_stat_fields..","..self:get_update_field("machine", 1) + else + -- 移动端+PC端+机器以外 归类到 其他 + client_stat_fields = client_stat_fields..","..self:get_update_field("other", 1) + end + end + + local os_regx = "(Windows|Linux|Macintosh)" + local os_res = ngx.re.match(ua, os_regx, "ijo") + if os_res then + os_res = string.lower(os_res[0]) + client_stat_fields = client_stat_fields..","..self:get_update_field(clients_map[os_res], 1) + end + end + + local other_regx = "MicroMessenger" + local other_res = ngx.re.find(ua, other_regx) + if other_res then + client_stat_fields = client_stat_fields..","..self:get_update_field("weixin", 1) + end + if client_stat_fields then + client_stat_fields = string.sub(client_stat_fields, 2) + end + return client_stat_fields +end + + +function _M.match_client_arr(self, ua) + local client_stat_fields = {} + + if not ua then + return client_stat_fields + end + + local clients_map = { + ["android"] = "android", + ["iphone"] = "iphone", + ["ipod"] = "iphone", + ["ipad"] = "iphone", + ["firefox"] = "firefox", + ["msie"] = "msie", + ["trident"] = "msie", + ["360se"] = "qh360", + ["360ee"] = "qh360", + ["360browser"] = "qh360", + ["qihoo"] = "qh360", + ["the world"] = "theworld", + ["theworld"] = "theworld", + ["tencenttraveler"] = "tt", + ["maxthon"] = "maxthon", + ["opera"] = "opera", + ["qqbrowser"] = "qq", + ["ucweb"] = "uc", + ["ubrowser"] = "uc", + ["safari"] = "safari", + ["chrome"] = "chrome", + ["metasr"] = "metasr", + ["2345explorer"] = "pc2345", + ["edge"] = "edeg", + ["edg"] = "edeg", + ["windows"] = "windows", + ["linux"] = "linux", + ["macintosh"] = "mac", + ["mobile"] = "mobile" + } + local mobile_regx = "(Mobile|Android|iPhone|iPod|iPad)" + local mobile_res = ngx.re.match(ua, mobile_regx, "ijo") + --mobile + if mobile_res then + client_stat_fields['mobile'] = 1 + mobile_res = string.lower(mobile_res[0]) + if mobile_res ~= "mobile" then + client_stat_fields[clients_map[mobile_res]] = 1 + end + else + --pc + -- 匹配结果的顺序,与ua中关键词的顺序有关 + -- lua的正则不支持|语法 + -- 短字符串string.find效率要比ngx正则高 + local pc_regx1 = "(360SE|360EE|360browser|Qihoo|TheWorld|TencentTraveler|Maxthon|Opera|QQBrowser|UCWEB|UBrowser|MetaSr|2345Explorer|Edg[e]*)" + local pc_res = ngx.re.match(ua, pc_regx1, "ijo") + local cls_pc = nil + + -- self:D("UA-JSON:"..self:to_json(ua)) + if "table" == type(ua) then + ua = self:to_json(ua) + end + + if not pc_res then + if ngx.re.find(ua, "[Ff]irefox") then + cls_pc = "firefox" + elseif string.find(ua, "MSIE") or string.find(ua, "Trident") then + cls_pc = "msie" + elseif string.find(ua, "[Cc]hrome") then + cls_pc = "chrome" + elseif string.find(ua, "[Ss]afari") then + cls_pc = "safari" + end + else + cls_pc = string.lower(pc_res[0]) + end + -- self:D("UA:"..ua) + -- self:D("PC cls:"..tostring(cls_pc)) + if cls_pc then + client_stat_fields[clients_map[cls_pc]] = 1 + else + -- machine and other + local machine_res, err = ngx.re.match(ua, "(ApacheBench|[Cc]url|HeadlessChrome|[a-zA-Z]+[Bb]ot|[Ww]get|[Ss]pider|[Cc]rawler|[Ss]crapy|zgrab|[Pp]ython|java)", "ijo") + if machine_res then + client_stat_fields["machine"] = 1 + else + -- 移动端+PC端+机器以外 归类到 其他 + client_stat_fields["other"] = 1 + end + end + + local os_regx = "(Windows|Linux|Macintosh)" + local os_res = ngx.re.match(ua, os_regx, "ijo") + if os_res then + os_res = string.lower(os_res[0]) + client_stat_fields[clients_map[os_res]] = 1 + end + end + + local other_regx = "MicroMessenger" + local other_res = ngx.re.find(ua, other_regx) + if other_res then + client_stat_fields["weixin"] = 1 + end + return client_stat_fields +end + +function _M.get_client_ip(self) + local client_ip = "unknown" + local cdn = auto_config['cdn'] + local request_header = ngx.req.get_headers() + if cdn == true then + for _,v in ipairs(auto_config['cdn_headers']) do + if request_header[v] ~= nil and request_header[v] ~= "" then + local ip_list = request_header[v] + 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 + client_ip = "unknown" + end + end + + return client_ip +end + +return _M \ No newline at end of file diff --git a/plugins/webstats/bak/webstats_log.lua b/plugins/webstats/bak/webstats_log.lua new file mode 100644 index 000000000..693424770 --- /dev/null +++ b/plugins/webstats/bak/webstats_log.lua @@ -0,0 +1,667 @@ +log_by_lua_block { + + 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 ver = '0.2.4' + local debug_mode = true + + local __C = require "webstats_common" + local C = __C:getInstance() + + -- cache start --- + local cache = ngx.shared.mw_total + local function cache_set(server_name, id ,key, val) + local line_kv = "log_kv_"..server_name..'_'..id.."_"..key + -- cache:set(line_kv, val) + cache:set(line_kv, val) + end + + local function cache_clear(server_name, id, key) + local line_kv = "log_kv_"..server_name..'_'..id.."_"..key + cache:delete(line_kv) + end + + local function cache_get(server_name, id, key) + local line_kv = "log_kv_"..server_name..'_'..id.."_"..key + local value = cache:get(line_kv) + return value + end + + + -- cache end --- + + -- domain config is import + local db = nil + local json = require "cjson" + local sqlite3 = require "lsqlite3" + local config = require "webstats_config" + local sites = require "webstats_sites" + + -- string.gsub(C:get_sn(ngx.var.server_name),'_','.') + 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 excluded = false + + local day = os.date("%d") + local number_day = tonumber(day) + local day_column = "day"..number_day + local flow_column = "flow"..number_day + local spider_column = "spider_flow"..number_day + --- default common var end --- + + local function init_var() + return true + end + + --------------------- exclude_func start -------------------------- + local function load_global_exclude_ip() + local load_key = "global_exclude_ip_load" + -- update global exclude ip + local global_exclude_ip = auto_config["exclude_ip"] + if global_exclude_ip then + for i, _ip in pairs(global_exclude_ip) + do + -- global + -- D("set global exclude ip: ".._ip) + if not cache:get("global_exclude_ip_".._ip) then + cache:set("global_exclude_ip_".._ip, true) + end + end + end + -- set tag + cache:set(load_key, true) + 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] + + local site_exclude_ip = nil + if site_config then + site_exclude_ip = site_config["exclude_ip"] + end + + -- update server_name exclude ip + if site_exclude_ip then + for i, _ip in pairs(site_exclude_ip) + do + cache:set(input_server_name .. "_exclude_ip_".._ip, true) + end + end + + -- set tag + cache:set(load_key, true) + return true + end + + 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 + + local function exclude_extension() + if not ngx.var.uri then return false end + if not auto_config['exclude_extension'] then return false end + for _,v in ipairs(auto_config['exclude_extension']) + do + if ngx.re.find(ngx.var.uri,"[.]"..v.."$",'ijo') then + return true + end + end + return false + end + + local function exclude_url() + if not ngx.var.uri then return false end + if not ngx.var.request_uri then return false end + if not auto_config['exclude_url'] then return false end + local the_uri = string.sub(ngx.var.request_uri, 2) + local url_conf = auto_config["exclude_url"] + for i,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 exclude_ip(input_server_name, ip) + -- 排除IP匹配,分网站单独的配置和全局配置两种方式 + local site_config = config[input_server_name] + local server_exclude_ips = nil + local check_server_exclude_ip = false + if site_config then + server_exclude_ips = site_config["exclude_ip"] + if not server_exclude_ips then + return false + end + for k, _ip in pairs(server_exclude_ips) + do + check_server_exclude_ip = true + break + end + end + -- D("server[" ..input_server_name.."]check exclude ip : "..tostring(check_server_exclude_ip)) + if check_server_exclude_ip then + if cache:get(input_server_name .. "_exclude_ip_"..ip) then + -- D("-Exclude server ip:"..ip) + return true + end + else + if cache:get("global_exclude_ip_"..ip) then + -- D("*Excluded global ip:"..ip) + return true + end + end + return false + end + --------------------- exclude_func end --------------------------- + + local function cache_logs_old(server_name) + + -- make new id + local new_id = C:get_last_id(server_name) + + local excluded = false + local ip = C:get_client_ip() + excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(server_name, ip) + + local ip_list = request_header["x-forwarded-for"] + if ip and not ip_list then + ip_list = ip + end + + local remote_addr = ngx.var.remote_addr + if not string.find(ip_list, remote_addr) then + if remote_addr then + ip_list = ip_list .. "," .. remote_addr + end + end + + -- local request_time = ngx.var.request_time + local request_time = C:get_request_time() + local client_port = ngx.var.remote_port + local real_server_name = ngx.var.server_name + local uri = ngx.var.uri + local status_code = ngx.status + local protocol = ngx.var.server_protocol + local request_uri = ngx.var.request_uri + local time_key = C:get_store_key() + local method = method + local body_length = C:get_length() + local domain = C:get_domain() + local referer = ngx.var.http_referer + + local kv = { + id=new_id, + time_key=time_key, + time=ngx.time(), + ip=ip, + domain=domain, + server_name=server_name, + real_server_name=real_server_name, + method=method, + status_code=status_code, + uri=uri, + request_uri=request_uri, + body_length=body_length, + referer=referer, + user_agent=request_header['user-agent'], + protocol=protocol, + is_spider=0, + request_time=request_time, + excluded=excluded, + request_headers='', + ip_list=ip_list, + client_port=client_port + } + + local request_stat_fields = "req=req+1,length=length+"..body_length + local spider_stat_fields = "x" + local client_stat_fields = "x" + + if not excluded then + + if status_code == 500 or (method=="POST" and config["record_post_args"] == true) or (status_code==403 and config["record_get_403_args"] == true) then + local data = "" + local ok, err = pcall(function() data = C:get_http_origin() end) + if ok and not err then + kv["request_headers"] = data + end + end + + if ngx.re.find("500,501,502,503,504,505,506,507,509,510,400,401,402,403,404,405,406,407,408,409,410,411,412,413,414,415,416,417,418,421,422,423,424,425,426,449,451,499", tostring(status_code), "jo") then + local field = "status_"..status_code + request_stat_fields = request_stat_fields .. ","..field.."="..field.."+1" + end + + local lower_method = string.lower(method) + if ngx.re.find("get,post,put,patch,delete", lower_method, "ijo") then + local field = "http_"..lower_method + request_stat_fields = request_stat_fields .. ","..field.."="..field.."+1" + end + + + local ipc = 0 + local pvc = 0 + local uvc = 0 + + local is_spider, request_spider, spider_index = C:match_spider(kv['user_agent']) + if not is_spider then + + client_stat_fields = C:match_client(kv['user_agent']) + if not client_stat_fields or #client_stat_fields == 0 then + client_stat_fields = request_stat_fields..",other=other+1" + end + + pvc, uvc = C:statistics_request(ip, is_spider,body_length) + ipc = C:statistics_ipc(server_name,ip) + else + kv["is_spider"] = spider_index + local field = "spider" + spider_stat_fields = request_spider.."="..request_spider.."+"..1 + request_stat_fields = request_stat_fields .. ","..field.."="..field.."+"..1 + end + + if ipc > 0 then + request_stat_fields = request_stat_fields..",ip=ip+1" + end + if uvc > 0 then + request_stat_fields = request_stat_fields..",uv=uv+1" + end + if pvc > 0 then + request_stat_fields = request_stat_fields..",pv=pv+1" + end + end + + local stat_fields = request_stat_fields..";"..client_stat_fields..";"..spider_stat_fields + + cache_set(server_name, new_id, "stat_fields", stat_fields) + cache_set(server_name, new_id, "log_kv", json.encode(kv)) + end + + local function cache_logs(input_sn) + + -- make new id + local new_id = C:get_last_id(input_sn) + + local excluded = false + local ip = C:get_client_ip() + excluded = filter_status() or exclude_extension() or exclude_url() or exclude_ip(input_sn, ip) + + local ip_list = request_header["x-forwarded-for"] + if ip and not ip_list then + ip_list = ip + end + + local remote_addr = ngx.var.remote_addr + + -- 修复反向代理代过来的数据 + if "table" == type(ip_list) then + ip_list = json.encode(ip_list) + end + + if not string.find(ip_list, remote_addr) then + if remote_addr then + ip_list = ip_list .. "," .. remote_addr + end + end + + -- local request_time = ngx.var.request_time + local request_time = C:get_request_time() + local client_port = ngx.var.remote_port + local real_server_name = input_sn + 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 time_key = C:get_store_key() + local method = method + local body_length = C:get_length() + local domain = C:get_domain() + local referer = ngx.var.http_referer + + local kv = { + id=new_id, + time_key=time_key, + time=ngx.time(), + ip=ip, + domain=domain, + server_name=input_sn, + real_server_name=real_server_name, + method=method, + status_code=status_code, + uri=uri, + request_uri=request_uri, + body_length=body_length, + referer=referer, + user_agent=request_header['user-agent'], + protocol=protocol, + is_spider=0, + request_time=request_time, + excluded=excluded, + request_headers='', + ip_list=ip_list, + client_port=client_port + } + + -- C:D(json.encode(kv)) + 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 + local data = "" + local ok, err = pcall(function() data = C:get_http_origin() end) + if ok and not err then + kv["request_headers"] = data + else + C:D("debug request_headers error:"..tostring(err)) + end + end + + if ngx.re.find("500,501,502,503,504,505,506,507,509,510,400,401,402,403,404,405,406,407,408,409,410,411,412,413,414,415,416,417,418,421,422,423,424,425,426,449,451,499", tostring(status_code), "jo") then + local field = "status_"..status_code + request_stat_fields[field] = 1 + end + + -- D("method:"..method) + local lower_method = string.lower(method) + if ngx.re.find("get,post,put,patch,delete", lower_method, "ijo") then + local field = "http_"..lower_method + request_stat_fields[field] = 1 + end + + + local ipc = 0 + local pvc = 0 + local uvc = 0 + + local is_spider, request_spider, spider_index = C:match_spider(kv['user_agent']) + if not is_spider then + + client_stat_fields = C:match_client_arr(kv['user_agent']) + pvc, uvc = C:statistics_request(ip, is_spider,body_length) + ipc = C:statistics_ipc(input_sn,ip) + else + kv["is_spider"] = spider_index + local field = "spider" + spider_stat_fields[request_spider] = 1 + request_stat_fields[field] = 1 + 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 + end + + local stat_fields = { + request_stat_fields=request_stat_fields, + client_stat_fields=client_stat_fields, + spider_stat_fields=spider_stat_fields, + } + + local data = { + server_name=input_sn, + stat_fields=stat_fields, + log_kv=kv, + } + + local push_data = json.encode(data) + -- C:D(json.encode(push_data)) + local key = C:getTotalKey() + ngx.shared.mw_total:rpush(key, push_data) + end + + local function store_logs_line(db, stmt, input_server_name, lineno) + local logvalue = cache_get(input_server_name, lineno, "log_kv") + if not logvalue then return false end + local logline = json.decode(logvalue) + + local time = logline["time"] + local id = logline["id"] + local protocol = logline["protocol"] + local client_port = logline["client_port"] + local status_code = logline["status_code"] + local uri = logline["uri"] + local request_uri = logline["request_uri"] + local method = logline["method"] + local body_length = logline["body_length"] + local referer = logline["referer"] + local ip = logline["ip"] + local ip_list = logline["ip_list"] + local request_time = logline["request_time"] + local is_spider = logline["is_spider"] + local domain = logline["domain"] + local server_name = logline["server_name"] + local user_agent = logline["user_agent"] + local request_headers = logline["request_headers"] + local excluded = logline["excluded"] + + local request_stat_fields = nil + local client_stat_fields = nil + local spider_stat_fields = nil + local stat_fields = cache_get(input_server_name, id, "stat_fields") + if stat_fields == nil then + -- D("Log stat fields is nil.") + -- D("Logdata:"..logvalue) + else + stat_fields = C:split(stat_fields, ";") + request_stat_fields = stat_fields[1] + client_stat_fields = stat_fields[2] + spider_stat_fields = stat_fields[3] + + if "x" == client_stat_fields then + client_stat_fields = nil + end + + if "x" == spider_stat_fields then + spider_stat_fields = nil + end + end + + local time_key = logline["time_key"] + if not excluded then + stmt:bind_names{ + time=time, + ip=ip, + domain=domain, + server_name=server_name, + method=method, + status_code=status_code, + uri=request_uri, + body_length=body_length, + referer=referer, + user_agent=user_agent, + protocol=protocol, + request_time=request_time, + is_spider=is_spider, + request_headers=request_headers, + ip_list=ip_list, + client_port=client_port, + } + + local res, err = stmt:step() + if tostring(res) == "5" then + -- D("step res:"..tostring(res)) + -- D("step err:"..tostring(err)) + -- D("the step database connection is busy, so it will be stored later.") + return false + end + stmt:reset() + -- D("store_logs_line ok") + + C:update_stat( db, "client_stat", time_key, client_stat_fields) + C:update_stat( db, "spider_stat", time_key, spider_stat_fields) + -- C:D("stat ok") + + -- only count non spider requests + local ok, err = C:statistics_uri(db, request_uri, ngx.md5(request_uri), body_length) + local ok, err = C:statistics_ip(db, ip, body_length) + end + + C:update_stat( db, "request_stat", time_key, request_stat_fields) + return true + end + + local function store_logs(input_sn) + if C:is_migrating(input_sn) == true then + -- D("migrating...") + return + end + + local last_insert_id_key = input_sn.."_last_id" + local store_start_id_key = input_sn.."_store_start" + local last_id = cache:get(last_insert_id_key) + local store_start = cache:get(store_start_id_key) + if store_start == nil then + store_start = 1 + end + local store_end = last_id + if store_end == nil then + store_end = 1 + end + + local worker_id = ngx.worker.id() + if C:is_working(input_sn) then + -- D("other workers are being stored, please store later.") + -- cache:delete(flush_data_key) + return true + end + C:lock_working(input_sn) + + local log_dir = "{$SERVER_APP}/logs" + local db_path = log_dir .. '/' .. input_sn .. "/logs.db" + local db, err = sqlite3.open(db_path) + + if tostring(err) ~= 'nil' then + C:D("sqlite3 open error:"..tostring(err)) + return true + end + + db:exec([[PRAGMA synchronous = 0]]) + db:exec([[PRAGMA page_size = 4096]]) + db:exec([[PRAGMA journal_mode = wal]]) + db:exec([[PRAGMA journal_size_limit = 1073741824]]) + + local stmt2 = nil + if db ~= nil then + stmt2 = 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 db == nil or stmt2 == nil then + -- D("web data db error") + -- cache:set(storing_key, false) + if db and db:isopen() then + db:close() + end + return true + end + + status, errorString = db:exec([[BEGIN TRANSACTION]]) + + C:clean_stats(db, input_sn) + + if store_end >= store_start then + for i=store_start, store_end, 1 do + -- D("store_start:"..store_start..":store_end:".. store_end) + if store_logs_line(db, stmt2, input_sn, i) then + cache_clear(input_sn, i, "log_kv") + cache_clear(input_sn, i, "stat_fields") + end + end + end + + local res, err = stmt2:finalize() + if tostring(res) == "5" then + C:D("Finalize res:"..tostring(res)..",Finalize err:"..tostring(err)) + end + + local now_date = os.date("*t") + local save_day = config['global']["save_day"] + local save_date_timestamp = os.time{year=now_date.year, + month=now_date.month, day=now_date.day-save_day, hour=0} + -- delete expire data + db:exec("DELETE FROM web_logs WHERE time<"..tostring(save_date_timestamp)) + + local res, err = db:execute([[COMMIT]]) + if db and db:isopen() then + db:close() + end + + cache:set(store_start_id_key, store_end+1) + C:unlock_working(input_sn) + end + + local function run_app() + -- C:D("------------ debug start ------------") + init_var() + + load_global_exclude_ip() + load_exclude_ip(server_name) + + cache_logs(server_name) + -- C:cron() + + -- cache_logs_old(server_name) + -- store_logs(server_name) + -- C:D("------------ debug 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() +} +