165 lines
4.6 KiB
Python
165 lines
4.6 KiB
Python
|
|
import logging
|
||
|
|
import threading
|
||
|
|
import time
|
||
|
|
import json
|
||
|
|
from collections import deque
|
||
|
|
from datetime import datetime
|
||
|
|
|
||
|
|
_log_deque = deque(maxlen=2000)
|
||
|
|
_subscribers = {}
|
||
|
|
_sub_lock = threading.Lock()
|
||
|
|
_sub_id_counter = 0
|
||
|
|
_log_id_lock = threading.Lock()
|
||
|
|
_log_id_counter = 0
|
||
|
|
|
||
|
|
SOURCE_ALIASES = {
|
||
|
|
'api': '外部API',
|
||
|
|
'frontend': '前端操作',
|
||
|
|
'db': '数据库',
|
||
|
|
'batch': '批量',
|
||
|
|
'system': '系统',
|
||
|
|
'外部API': '外部API',
|
||
|
|
'前端操作': '前端操作',
|
||
|
|
'数据库': '数据库',
|
||
|
|
'批量': '批量',
|
||
|
|
'系统': '系统',
|
||
|
|
}
|
||
|
|
|
||
|
|
|
||
|
|
def _next_log_id():
|
||
|
|
global _log_id_counter
|
||
|
|
with _log_id_lock:
|
||
|
|
_log_id_counter += 1
|
||
|
|
return _log_id_counter
|
||
|
|
|
||
|
|
|
||
|
|
def _resolve_source(token):
|
||
|
|
if not token:
|
||
|
|
return None
|
||
|
|
return SOURCE_ALIASES.get(token, token)
|
||
|
|
|
||
|
|
|
||
|
|
class SSELogHandler(logging.Handler):
|
||
|
|
def emit(self, record):
|
||
|
|
try:
|
||
|
|
msg = self.format(record)
|
||
|
|
level = record.levelname
|
||
|
|
if level == 'WARNING':
|
||
|
|
level = 'WARN'
|
||
|
|
source = getattr(record, 'log_source', '系统')
|
||
|
|
trace_id = getattr(record, 'trace_id', None) or ''
|
||
|
|
entry = {
|
||
|
|
'id': _next_log_id(),
|
||
|
|
'timestamp': datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3],
|
||
|
|
'time': datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3],
|
||
|
|
'level': level,
|
||
|
|
'source': source,
|
||
|
|
'message': msg,
|
||
|
|
'trace_id': trace_id,
|
||
|
|
}
|
||
|
|
_log_deque.append(entry)
|
||
|
|
_notify_subscribers(entry)
|
||
|
|
except Exception:
|
||
|
|
pass
|
||
|
|
|
||
|
|
|
||
|
|
def _notify_subscribers(entry):
|
||
|
|
with _sub_lock:
|
||
|
|
for q in _subscribers.values():
|
||
|
|
q.append(entry)
|
||
|
|
|
||
|
|
|
||
|
|
def subscribe():
|
||
|
|
global _sub_id_counter
|
||
|
|
with _sub_lock:
|
||
|
|
_sub_id_counter += 1
|
||
|
|
sub_id = _sub_id_counter
|
||
|
|
q = deque(maxlen=500)
|
||
|
|
_subscribers[sub_id] = q
|
||
|
|
return sub_id, q
|
||
|
|
|
||
|
|
|
||
|
|
def unsubscribe(sub_id):
|
||
|
|
with _sub_lock:
|
||
|
|
_subscribers.pop(sub_id, None)
|
||
|
|
|
||
|
|
|
||
|
|
def get_all_logs():
|
||
|
|
return list(_log_deque)
|
||
|
|
|
||
|
|
|
||
|
|
def query_logs(after_id=0, level=None, source=None, keyword=None, limit=200):
|
||
|
|
after_id = int(after_id or 0)
|
||
|
|
limit = max(1, min(int(limit or 200), 1000))
|
||
|
|
norm_source = _resolve_source(source) if source else None
|
||
|
|
keyword_lc = (keyword or '').lower()
|
||
|
|
|
||
|
|
matched = []
|
||
|
|
for entry in _log_deque:
|
||
|
|
if entry['id'] <= after_id:
|
||
|
|
continue
|
||
|
|
if level and entry.get('level') != level:
|
||
|
|
continue
|
||
|
|
if norm_source and entry.get('source') != norm_source:
|
||
|
|
continue
|
||
|
|
if keyword_lc and keyword_lc not in entry.get('message', '').lower():
|
||
|
|
continue
|
||
|
|
matched.append(entry)
|
||
|
|
matched = matched[-limit:] if len(matched) > limit else matched
|
||
|
|
last_id = matched[-1]['id'] if matched else after_id
|
||
|
|
return {'logs': matched, 'last_id': last_id}
|
||
|
|
|
||
|
|
|
||
|
|
def setup_logging():
|
||
|
|
from lib.config import get_log_config
|
||
|
|
log_cfg = get_log_config()
|
||
|
|
log_dir = log_cfg['file_path']
|
||
|
|
import os
|
||
|
|
if not os.path.isabs(log_dir):
|
||
|
|
project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||
|
|
log_dir = os.path.join(project_root, log_dir)
|
||
|
|
os.makedirs(log_dir, exist_ok=True)
|
||
|
|
|
||
|
|
log_level = getattr(logging, log_cfg['level'].upper(), logging.INFO)
|
||
|
|
|
||
|
|
root_logger = logging.getLogger()
|
||
|
|
root_logger.setLevel(log_level)
|
||
|
|
|
||
|
|
fmt = logging.Formatter('%(message)s')
|
||
|
|
|
||
|
|
sse_handler = SSELogHandler()
|
||
|
|
sse_handler.setLevel(log_level)
|
||
|
|
sse_handler.setFormatter(fmt)
|
||
|
|
root_logger.addHandler(sse_handler)
|
||
|
|
|
||
|
|
file_handler = logging.FileHandler(
|
||
|
|
os.path.join(log_dir, f"seal_{datetime.now().strftime('%Y%m%d')}.log"),
|
||
|
|
encoding='utf-8',
|
||
|
|
)
|
||
|
|
file_handler.setLevel(log_level)
|
||
|
|
file_handler.setFormatter(logging.Formatter(
|
||
|
|
'[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s',
|
||
|
|
datefmt='%Y-%m-%d %H:%M:%S'
|
||
|
|
))
|
||
|
|
root_logger.addHandler(file_handler)
|
||
|
|
|
||
|
|
console_handler = logging.StreamHandler()
|
||
|
|
console_handler.setLevel(log_level)
|
||
|
|
console_handler.setFormatter(logging.Formatter(
|
||
|
|
'[%(asctime)s] [%(levelname)s] %(message)s',
|
||
|
|
datefmt='%H:%M:%S'
|
||
|
|
))
|
||
|
|
root_logger.addHandler(console_handler)
|
||
|
|
|
||
|
|
|
||
|
|
def log_info(msg, source='前端操作', trace_id=None):
|
||
|
|
logging.info(msg, extra={'log_source': source, 'trace_id': trace_id or ''})
|
||
|
|
|
||
|
|
|
||
|
|
def log_warn(msg, source='前端操作', trace_id=None):
|
||
|
|
logging.warning(msg, extra={'log_source': source, 'trace_id': trace_id or ''})
|
||
|
|
|
||
|
|
|
||
|
|
def log_error(msg, source='前端操作', trace_id=None):
|
||
|
|
logging.error(msg, extra={'log_source': source, 'trace_id': trace_id or ''})
|