Files
MediaCrawler/tools/app_runner.py
程序员阿江(Relakkes) 8a0fd49b96 refactor: 抽离应用 runner 并优化退出清理
- 新增 tools/app_runner.py 统一信号/取消/清理超时逻辑
- main.py 精简为业务入口与资源清理实现
- CDPBrowserManager 不再覆盖已有 SIGINT/SIGTERM 处理器
2025-12-15 18:06:57 +08:00

110 lines
3.6 KiB
Python

# -*- coding: utf-8 -*-
# Copyright (c) 2025 relakkes@gmail.com
#
# This file is part of MediaCrawler project.
# Repository: https://github.com/NanmiCoder/MediaCrawler
# GitHub: https://github.com/NanmiCoder
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
from __future__ import annotations
import asyncio
import os
import signal
from collections.abc import Awaitable, Callable
from typing import Optional
AsyncFn = Callable[[], Awaitable[None]]
def run(
app_main: AsyncFn,
app_cleanup: AsyncFn,
*,
cleanup_timeout_seconds: float = 15.0,
on_first_interrupt: Optional[Callable[[], None]] = None,
force_exit_code: int = 130,
) -> None:
async def _cleanup_with_timeout() -> None:
try:
await asyncio.wait_for(asyncio.shield(app_cleanup()), timeout=cleanup_timeout_seconds)
except asyncio.TimeoutError:
print(f"[Main] 清理超时({cleanup_timeout_seconds}s),跳过剩余清理。")
async def _cancel_remaining_tasks(timeout_seconds: float = 2.0) -> None:
current = asyncio.current_task()
tasks = [t for t in asyncio.all_tasks() if t is not current and not t.done()]
if not tasks:
return
for t in tasks:
t.cancel()
try:
await asyncio.wait_for(
asyncio.gather(*tasks, return_exceptions=True),
timeout=timeout_seconds,
)
except asyncio.TimeoutError:
pass
async def _runner() -> None:
loop = asyncio.get_running_loop()
runner_task = asyncio.current_task()
if runner_task is None:
raise RuntimeError("Runner task not found")
shutdown_requested = False
def _on_signal(signum: int) -> None:
nonlocal shutdown_requested
if shutdown_requested:
print("[Main] 再次收到中断信号,强制退出。")
os._exit(force_exit_code)
shutdown_requested = True
print(f"\n[Main] 收到中断信号 {signum},正在退出(清理最多{cleanup_timeout_seconds}s)...")
if on_first_interrupt is not None:
try:
on_first_interrupt()
except Exception:
pass
runner_task.cancel()
try:
loop.add_signal_handler(signal.SIGINT, _on_signal, signal.SIGINT)
loop.add_signal_handler(signal.SIGTERM, _on_signal, signal.SIGTERM)
except NotImplementedError:
signal.signal(signal.SIGINT, lambda signum, _frame: _on_signal(signum))
signal.signal(signal.SIGTERM, lambda signum, _frame: _on_signal(signum))
cancelled = False
try:
await app_main()
except asyncio.CancelledError:
cancelled = True
finally:
try:
await _cleanup_with_timeout()
except Exception as e:
print(f"[Main] 清理时出错: {e}")
await _cancel_remaining_tasks()
if cancelled:
return
asyncio.run(_runner())