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 ''})