better-staridc-MNBT
1#!/usr/bin/python
2# coding: utf-8
3
4import json
5import os
6import signal
7import sqlite3
8import subprocess
9import sys
10import time
11from stats_api import StatsAPIMixin
12
13
14BASE_DIR = os.path.dirname(os.path.abspath(__file__))
15INDEX_DB_PATH = os.path.join(BASE_DIR, "file_index.db")
16
17
18class mnbt_connector_main(StatsAPIMixin):
19 def __init__(self):
20 self.base_dir = os.path.dirname(os.path.abspath(__file__))
21 self.config_path = os.path.join(self.base_dir, "config.json")
22 self.worker_path = os.path.join(self.base_dir, "worker.py")
23 self.pid_path = os.path.join(self.base_dir, "worker.pid")
24 self.log_path = os.path.join(self.base_dir, "worker.log")
25
26 def _try_write_log(self, msg):
27 try:
28 import public
29 public.WriteLog('MNBT连接器', msg)
30 except Exception:
31 pass
32
33 def get_site_list(self, args):
34 try:
35 import public
36 panel_sites = public.M('sites').field('name').select()
37 panel_names = {s["name"] for s in panel_sites}
38 except Exception:
39 panel_names = set()
40 if panel_names:
41 stats_data = super().get_site_list(args).get("data", [])
42 stats_map = {s["site_name"]: s for s in stats_data}
43 data = []
44 for name in sorted(panel_names):
45 merged = {"site_name": name, "pv": 0, "uv": 0, "total_bytes": 0}
46 if name in stats_map:
47 merged.update(stats_map[name])
48 data.append(merged)
49 extra_names = {s["site_name"] for s in stats_data} - panel_names
50 for name in sorted(extra_names):
51 data.append(stats_map[name])
52 else:
53 data = super().get_site_list(args).get("data", [])
54 return {"status": True, "data": data}
55
56 def reset_site_stats(self, args):
57 site = args.get("site", "").strip()
58 if not site:
59 return {"status": False, "msg": "缺少站点名称"}
60 if ".." in site or "/" in site or "\\" in site:
61 return {"status": False, "msg": "非法的站点名称"}
62 try:
63 from stats_collector import reset_site_stats as do_reset
64 deleted = do_reset(site)
65 self._try_write_log(f"重置站点统计: {site}")
66 return {
67 "status": True,
68 "msg": "站点统计已重置,下次采集时将重新回溯全部历史日志",
69 "deleted": deleted,
70 }
71 except Exception as exc:
72 return {"status": False, "msg": f"重置失败: {exc}"}
73
74 def _read_config(self):
75 if not os.path.exists(self.config_path):
76 return {}
77 with open(self.config_path, "r", encoding="utf-8") as handle:
78 return json.load(handle)
79
80 def _read_pid(self):
81 if not os.path.exists(self.pid_path):
82 return 0
83 try:
84 with open(self.pid_path, "r", encoding="utf-8") as handle:
85 return int(handle.read().strip() or "0")
86 except (OSError, ValueError):
87 return 0
88
89 def _is_process_running(self, pid):
90 if pid <= 0:
91 return False
92 try:
93 os.kill(pid, 0)
94 return True
95 except OSError:
96 return False
97
98 def _clear_stale_pid(self, pid):
99 if pid and not self._is_process_running(pid) and os.path.exists(self.pid_path):
100 try:
101 os.remove(self.pid_path)
102 except OSError:
103 pass
104
105 def _tail_log(self, max_bytes=4000):
106 if not os.path.exists(self.log_path):
107 return ""
108 try:
109 with open(self.log_path, "rb") as handle:
110 handle.seek(0, os.SEEK_END)
111 size = handle.tell()
112 handle.seek(max(0, size - max_bytes))
113 return handle.read().decode("utf-8", errors="ignore")
114 except OSError:
115 return ""
116
117 def _log_meta(self):
118 if not os.path.exists(self.log_path):
119 return {"log_size": 0, "log_mtime": ""}
120 try:
121 stat_info = os.stat(self.log_path)
122 return {
123 "log_size": stat_info.st_size,
124 "log_mtime": time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(stat_info.st_mtime)),
125 }
126 except OSError:
127 return {"log_size": 0, "log_mtime": ""}
128
129 @staticmethod
130 def _mask_secret_value(value):
131 if value is None:
132 return ""
133 text = str(value)
134 if len(text) <= 8:
135 return "****" if text else ""
136 return text[:4] + "****" + text[-4:]
137
138 def _mask_config(self, value):
139 secret_keys = {"platform_secret", "node_secret", "api_key", "secret", "token", "password", "key"}
140 if isinstance(value, dict):
141 masked = {}
142 for key, item in value.items():
143 key_text = str(key).lower()
144 if key_text in secret_keys or "secret" in key_text or "token" in key_text or "password" in key_text:
145 masked[key] = self._mask_secret_value(item)
146 else:
147 masked[key] = self._mask_config(item)
148 return masked
149 if isinstance(value, list):
150 return [self._mask_config(item) for item in value]
151 return value
152
153 def get_runtime_status(self, args):
154 pid = self._read_pid()
155 running = self._is_process_running(pid)
156 if not running:
157 self._clear_stale_pid(pid)
158 pid = 0
159 return {
160 "status": True,
161 "running": running,
162 "pid": pid,
163 "log_path": self.log_path,
164 "last_log": self._tail_log(),
165 }
166
167 def get_status(self, args):
168 config = self._read_config()
169 runtime = self.get_runtime_status(args)
170 return {
171 "status": True,
172 "msg": "MNBT 连接器运行中" if runtime["running"] else "MNBT 连接器已停止",
173 "node_id": config.get("node_id", ""),
174 "version": config.get("version", "0.1.0"),
175 "capabilities": config.get("capabilities", []),
176 "running": runtime["running"],
177 "pid": runtime["pid"],
178 "last_log": runtime["last_log"],
179 }
180
181 def get_logs(self, args):
182 max_bytes = int(args.get("max_bytes", 20000) or 20000)
183 max_bytes = min(max(max_bytes, 1000), 200000)
184 runtime = self.get_runtime_status(args)
185 meta = self._log_meta()
186 return {
187 "status": True,
188 "log_path": self.log_path,
189 "log_text": self._tail_log(max_bytes),
190 "running": runtime["running"],
191 "pid": runtime["pid"],
192 "max_bytes": max_bytes,
193 "log_size": meta["log_size"],
194 "log_mtime": meta["log_mtime"],
195 }
196
197 def get_log_list(self, args):
198 log_dir = os.path.dirname(self.log_path)
199 log_base = os.path.basename(self.log_path)
200 log_files = []
201 try:
202 for entry in sorted(os.listdir(log_dir), reverse=True):
203 if entry == log_base or entry.startswith(log_base + "."):
204 full_path = os.path.join(log_dir, entry)
205 try:
206 stat_info = os.stat(full_path)
207 log_files.append({
208 "name": entry,
209 "path": full_path,
210 "size": stat_info.st_size,
211 "size_mb": round(stat_info.st_size / 1024 / 1024, 2),
212 "mtime": time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(stat_info.st_mtime)),
213 })
214 except OSError:
215 pass
216 except OSError as exc:
217 return {"status": False, "msg": str(exc), "data": []}
218 return {
219 "status": True,
220 "data": log_files,
221 "total_count": len(log_files),
222 }
223
224 def get_log_content(self, args):
225 target = args.get("file", "")
226 offset = int(args.get("offset", 0) or 0)
227 limit = int(args.get("limit", 50000) or 50000)
228 limit = min(max(limit, 1000), 500000)
229 keyword = args.get("keyword", "").strip()
230 level = args.get("level", "").strip().upper()
231
232 if target and os.sep in target:
233 return {"status": False, "msg": "非法的日志文件名"}
234
235 log_dir = os.path.dirname(self.log_path)
236 log_path = os.path.join(log_dir, target) if target else self.log_path
237
238 if not os.path.exists(log_path):
239 return {"status": False, "msg": "日志文件不存在", "data": []}
240
241 try:
242 file_size = os.path.getsize(log_path)
243 with open(log_path, "r", encoding="utf-8", errors="ignore") as handle:
244 if offset > 0:
245 handle.seek(min(offset, file_size))
246 raw = handle.read(limit)
247 except OSError as exc:
248 return {"status": False, "msg": str(exc), "data": []}
249
250 lines = raw.splitlines()
251 filtered_lines = []
252 for line in lines:
253 if not line.strip():
254 continue
255 if keyword and keyword not in line:
256 continue
257 if level and f"[{level}]" not in line:
258 continue
259 filtered_lines.append(line)
260
261 current_offset = offset + len(raw.encode("utf-8", errors="ignore"))
262 has_more = current_offset < file_size
263
264 return {
265 "status": True,
266 "data": filtered_lines,
267 "total_lines": len(filtered_lines),
268 "file_size": file_size,
269 "current_offset": current_offset,
270 "has_more": has_more,
271 "keyword": keyword,
272 "level": level,
273 }
274
275 def clear_log(self, args):
276 target = args.get("file", "")
277 if target and os.sep in target:
278 return {"status": False, "msg": "非法的日志文件名"}
279
280 log_dir = os.path.dirname(self.log_path)
281 log_path = os.path.join(log_dir, target) if target else self.log_path
282
283 if not os.path.exists(log_path):
284 return {"status": True, "msg": "日志文件不存在,无需清空"}
285
286 try:
287 with open(log_path, "w", encoding="utf-8"):
288 pass
289 self._try_write_log("日志已清空")
290 return {"status": True, "msg": "日志已清空"}
291 except OSError as exc:
292 return {"status": False, "msg": f"清空失败: {exc}"}
293
294 def get_log_level(self, args):
295 config = self._read_config()
296 return {
297 "status": True,
298 "level": config.get("log_level", "INFO"),
299 "available_levels": ["DEBUG", "INFO", "WARNING", "ERROR"],
300 }
301
302 def set_log_level(self, args):
303 level = args.get("level", "").strip().upper()
304 if level not in ("DEBUG", "INFO", "WARNING", "ERROR"):
305 return {"status": False, "msg": "无效的日志级别,可选: DEBUG, INFO, WARNING, ERROR"}
306 config = self._read_config()
307 config["log_level"] = level
308 try:
309 with open(self.config_path, "w", encoding="utf-8") as handle:
310 json.dump(config, handle, ensure_ascii=False, indent=2)
311 self._try_write_log(f"日志级别已设置为 {level}")
312 return {"status": True, "msg": f"日志级别已设置为 {level},重启 worker 后生效"}
313 except OSError as exc:
314 return {"status": False, "msg": str(exc)}
315
316 def get_worker_log_stats(self, args):
317 target = args.get("file", "")
318 if target and os.sep in target:
319 return {"status": False, "msg": "非法的日志文件名"}
320
321 log_dir = os.path.dirname(self.log_path)
322 log_path = os.path.join(log_dir, target) if target else self.log_path
323
324 if not os.path.exists(log_path):
325 return {"status": True, "data": {"DEBUG": 0, "INFO": 0, "WARNING": 0, "ERROR": 0, "total": 0}}
326
327 counts = {"DEBUG": 0, "INFO": 0, "WARNING": 0, "ERROR": 0, "total": 0}
328 try:
329 with open(log_path, "r", encoding="utf-8", errors="ignore") as handle:
330 for line in handle:
331 counts["total"] += 1
332 for level in ("DEBUG", "INFO", "WARNING", "ERROR"):
333 if f"[{level}]" in line:
334 counts[level] += 1
335 break
336 return {"status": True, "data": counts}
337 except OSError as exc:
338 return {"status": False, "msg": str(exc), "data": counts}
339
340 def get_config_info(self, args):
341 config = self._read_config()
342 masked_config = self._mask_config(config)
343 return {
344 "status": True,
345 "config_path": self.config_path,
346 "config": masked_config,
347 "config_json": json.dumps(masked_config, ensure_ascii=False, indent=2),
348 "exists": os.path.exists(self.config_path),
349 }
350
351 def start(self, args):
352 if not os.path.exists(self.config_path):
353 return {"status": False, "msg": "config.json 配置文件不存在"}
354 runtime = self.get_runtime_status(args)
355 if runtime["running"]:
356 return {"status": True, "msg": "工作进程已在运行中", "pid": runtime["pid"]}
357 process = subprocess.Popen(
358 [sys.executable, self.worker_path],
359 cwd=self.base_dir,
360 stdout=open(self.log_path, "a", encoding="utf-8"),
361 stderr=subprocess.STDOUT,
362 stdin=subprocess.DEVNULL,
363 start_new_session=True,
364 )
365 with open(self.pid_path, "w", encoding="utf-8") as handle:
366 handle.write(str(process.pid))
367 time.sleep(0.3)
368 self._try_write_log(f'工作进程已启动 PID={process.pid}')
369 return {
370 "status": True,
371 "msg": "工作进程已启动",
372 "pid": process.pid,
373 }
374
375 def stop(self, args):
376 pid = self._read_pid()
377 if pid <= 0:
378 return {"status": True, "msg": "工作进程未运行"}
379 if self._is_process_running(pid):
380 try:
381 os.kill(pid, signal.SIGTERM)
382 except OSError as exc:
383 return {"status": False, "msg": str(exc)}
384 try:
385 if os.path.exists(self.pid_path):
386 os.remove(self.pid_path)
387 except OSError:
388 pass
389
390 return {"status": True, "msg": "工作进程已停止"}
391
392 def run_once(self, args):
393 if not os.path.exists(self.config_path):
394 return {"status": False, "msg": "config.json 配置文件不存在"}
395 try:
396 result = subprocess.run(
397 [sys.executable, self.worker_path, "--once"],
398 stdout=subprocess.PIPE,
399 stderr=subprocess.PIPE,
400 text=True,
401 timeout=120,
402 )
403 ok = result.returncode == 0
404 self._try_write_log(f'执行一次{"成功" if ok else "失败"}')
405 return {
406 "status": ok,
407 "msg": "执行成功" if ok else result.stderr,
408 "stdout": result.stdout,
409 }
410 except subprocess.TimeoutExpired as exc:
411 return {
412 "status": False,
413 "msg": "执行超时",
414 "stdout": exc.stdout or "",
415 "stderr": exc.stderr or "",
416 }
417
418 def get_index_stats(self, args):
419 if not os.path.exists(INDEX_DB_PATH):
420 return {
421 "status": True,
422 "data": {
423 "total_files": 0,
424 "total_size": 0,
425 "scan_modes": {},
426 "index_exists": False,
427 },
428 }
429
430 try:
431 conn = sqlite3.connect(INDEX_DB_PATH)
432 cursor = conn.cursor()
433
434 cursor.execute("SELECT COUNT(*) FROM file_index")
435 total_files = cursor.fetchone()[0]
436
437 cursor.execute("SELECT SUM(size) FROM file_index")
438 total_size = cursor.fetchone()[0] or 0
439
440 cursor.execute("SELECT scan_mode, COUNT(*) FROM file_index GROUP BY scan_mode")
441 scan_modes = {row[0]: row[1] for row in cursor.fetchall()}
442
443 conn.close()
444
445 return {
446 "status": True,
447 "data": {
448 "total_files": total_files,
449 "total_size": total_size,
450 "total_size_mb": round(total_size / 1024 / 1024, 2),
451 "scan_modes": scan_modes,
452 "index_exists": True,
453 },
454 }
455 except Exception as exc:
456 return {
457 "status": False,
458 "msg": str(exc),
459 "data": None,
460 }
461
462 def trigger_full_scan(self, args):
463 config = self._read_config()
464 if not config.get("mnbt_url"):
465 return {"status": False, "msg": "未配置 MNBT 地址"}
466
467 result = subprocess.run(
468 [sys.executable, self.worker_path, "--full-scan"],
469 stdout=subprocess.PIPE,
470 stderr=subprocess.PIPE,
471 text=True,
472 timeout=300,
473 )
474 return {
475 "status": result.returncode == 0,
476 "msg": "全量扫描已触发" if result.returncode == 0 else result.stderr,
477 }
478
479
480if __name__ == "__main__":
481 import sys
482 try:
483 payload = json.loads(sys.argv[1]) if len(sys.argv) > 1 else {}
484 except (IndexError, json.JSONDecodeError):
485 payload = {}
486
487 action = payload.get("s", "")
488 if not action:
489 print(json.dumps({"status": False, "msg": "缺少参数 s"}))
490 sys.exit(1)
491
492 args = {k: v for k, v in payload.items() if k not in ("s", "action")}
493
494 plugin = mnbt_connector_main()
495 ALLOWED_ACTIONS = {
496 "get_status", "start", "stop", "run_once", "get_logs", "get_config_info",
497 "get_site_list", "get_site_overview", "get_site_trend",
498 "get_site_ip_rank", "get_site_uri_rank", "get_site_error_logs",
499 "get_site_recent_logs", "get_site_spider_analysis",
500 "get_site_client_stats", "get_site_method_stats",
501 "get_index_stats",
502 "get_log_list", "get_log_content", "clear_log",
503 "get_log_level", "set_log_level", "get_worker_log_stats",
504 "reset_site_stats",
505 }
506 if action not in ALLOWED_ACTIONS:
507 print(json.dumps({"status": False, "msg": f"未知操作: {action}"}))
508 sys.exit(1)
509 method = getattr(plugin, action, None)
510
511 try:
512 result = method(args)
513 print(json.dumps(result, ensure_ascii=False, default=str))
514 except Exception as exc:
515 print(json.dumps({"status": False, "msg": str(exc)}, ensure_ascii=False))
516 sys.exit(1)