Metadata-Version: 2.4
Name: llm-async-scheduler
Version: 0.1.4
Summary: 基于 APScheduler 的异步任务调度通用库，支持 cron、interval 与 manual 触发模式。
License-Expression: MIT
License-File: LICENSE
Keywords: apscheduler,async,cron,scheduler,task
Classifier: Development Status :: 4 - Beta
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Requires-Python: >=3.12
Requires-Dist: apscheduler<4,>=3.10.4
Description-Content-Type: text/markdown

# llm-async-scheduler

基于 [APScheduler](https://apscheduler.readthedocs.io/) 的异步任务调度通用库，提供统一的任务抽象与 cron / interval / manual 三种触发模式。

> PyPI 包名：`llm-async-scheduler`，import 名：`async_scheduler`。

## 特性

- 统一的 `ScheduleTask` 基类，内置超时、重试与生命周期钩子
- 基于 `AsyncIOScheduler` 的 cron 与 interval 调度
- manual 模式支持按需手动触发
- cron / interval 支持 `immediate=True` 启动后立即执行一次
- 同一任务不支持并发执行（`max_instances=1`）
- **异步上下文管理器**：自动管理调度器生命周期
- Python 3.12+，纯 asyncio

## 安装

```bash
pip install llm-async-scheduler
```

或使用 [uv](https://docs.astral.sh/uv/)：

```bash
uv add llm-async-scheduler
```

## 快速开始

### 定义任务

继承 `ScheduleTask` 并实现 `run()`：

```python
import asyncio
from async_scheduler import AsyncTaskScheduler, ScheduleTask


class HelloTask(ScheduleTask):
    async def run(self) -> None:
        print(f"hello from {self.task_id}")


async def main() -> None:
    task = HelloTask("hello", timeout=30, retry_times=1, tags=["demo"])
    scheduler = AsyncTaskScheduler(task, "manual")
    await scheduler.run_manual_once()


asyncio.run(main())
```

### Cron 调度

```python
import asyncio
from async_scheduler import AsyncTaskScheduler, ScheduleTask


class SyncDataTask(ScheduleTask):
    async def run(self) -> None:
        # 业务逻辑
        ...


async def main() -> None:
    task = SyncDataTask("sync-data")
    scheduler = AsyncTaskScheduler(
        task,
        "cron",
        cron_expr="0 */6 * * *",  # 每 6 小时
    )
    scheduler.start()
    await scheduler.wait_until_done()  # 常驻进程，可通过 stop() 退出


asyncio.run(main())
```

### Interval 调度

```python
scheduler = AsyncTaskScheduler(
    task,
    "interval",
    interval_seconds=60,
    immediate=True,  # 启动后立即执行一次，之后按 interval 周期调度
)
scheduler.start()
await scheduler.wait_until_done()
```

## API 概览

### ScheduleTask

| 参数 | 类型 | 默认值 | 说明 |
|------|------|--------|------|
| `task_id` | `str` | — | 任务唯一标识 |
| `timeout` | `int` | `3600` | 单次执行超时（秒） |
| `retry_times` | `int` | `0` | 失败后重试次数 |
| `retry_delay` | `int` | `5` | 重试间隔（秒） |
| `tags` | `Sequence[str]` | `[]` | 可选标签，便于日志与监控 |

可覆写钩子：

- `before_run()`：执行前（可用于校验、初始化、加锁）
- `run()`：核心业务（**必须实现**）
- `after_run(success)`：执行后（`success` 表示是否成功）
- `on_error(exc)`：每次失败时（可自定义告警逻辑）

> ⚠️ **重要**：调度器调用 `safe_execute()`，不要直接调用 `run()`。

### AsyncTaskScheduler

| 参数 | 类型 | 说明 |
|------|------|------|
| `task` | `ScheduleTask` | 任务实例 |
| `trigger_mode` | `"cron" \| "interval" \| "manual"` | 触发模式 |
| `cron_expr` | `str \| None` | cron 表达式（5 段，如 `0 2 * * *`） |
| `interval_seconds` | `int \| None` | interval 间隔秒数 |
| `immediate` | `bool` | `False` | cron / interval 模式下启动后是否立即执行一次 |

cron / interval job 固定 `max_instances=1`，同一任务不会并发执行。若单次执行耗时超过调度间隔，后续触发会被跳过（misfire）。

常用方法：

| 方法 | 说明 |
|------|------|
| `start()` | 启动调度（manual 模式可跳过） |
| `stop()` | 停止调度并唤醒 `wait_until_done()` |
| `run_manual_once()` | manual 模式手动执行一次 |
| `wait_until_done()` | 阻塞等待结束或 shutdown |
| `request_shutdown()` | 仅唤醒 `wait_until_done()`，不停止 job |

#### 异步上下文管理器（推荐）

`AsyncTaskScheduler` 支持异步上下文管理器，自动处理启动和关闭：

```python
# ✅ 推荐：使用上下文管理器
async with AsyncTaskScheduler(task, "interval", interval_seconds=60) as scheduler:
    await scheduler.wait_until_done()
# 自动调用 start() 和 stop()

# ❌ 手动管理（容易遗漏 stop）
scheduler = AsyncTaskScheduler(task, "interval", interval_seconds=60)
scheduler.start()
try:
    await scheduler.wait_until_done()
finally:
    scheduler.stop()
```

## 错误处理

### 执行流程与错误传播

```
safe_execute()
├── before_run()
│   └── ❌ 异常 → 跳过 run()，直接进入 after_run(False)
├── run() (重试循环)
│   ├── ✅ 成功 → break
│   ├── ❌ 超时 (TimeoutError) → on_error() → 重试或退出
│   ├── ❌ 业务异常 → on_error() → 重试或退出
│   └── ❌ CancelledError → 直接抛出（不重试）
└── after_run(success)
    └── ❌ 异常 → 记录日志但不传播（避免吞掉原始错误）
```

### 关键设计决策

1. **CancelledError 不重试**：符合 Python asyncio 规范，取消信号必须立即传播
2. **after_run 异常隔离**：即使 after_run 失败也不会影响主流程，仅记录日志
3. **超时使用 `asyncio.wait_for`**：确保长时间运行的任务能被及时终止
4. **所有异常都经过 `on_error()`**：统一告警入口

### 常见问题排查

| 场景 | 可能原因 | 解决方案 |
|------|---------|---------|
| 任务从未执行 | cron 表达式错误 | 使用 [crontab.guru](https://crontab.guru/) 验证表达式 |
| 重试未生效 | `retry_times=0`（默认） | 显式设置 `retry_times>=1` |
| 任务卡住 | 超时时间过长 | 降低 `timeout` 值（默认 3600 秒） |
| 并发执行 | 多个调度器实例 | 确保同一 task_id 只创建一个调度器 |

## 最佳实践

### 1. 优雅关闭

对于长期运行的调度任务，建议结合信号处理实现优雅关闭：

```python
import signal
import asyncio
from async_scheduler import AsyncTaskScheduler, ScheduleTask


class LongRunningTask(ScheduleTask):
    async def run(self) -> None:
        # 检查是否需要提前终止
        await asyncio.sleep(10)


async def main() -> None:
    task = LongRunningTask("worker")

    loop = asyncio.get_running_loop()
    shutdown_event = asyncio.Event()

    for sig in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(sig, lambda: shutdown_event.set())

    async with AsyncTaskScheduler(
        task, "interval", interval_seconds=30, immediate=True
    ):
        await shutdown_event.wait()
        print("收到关闭信号，正在退出...")


asyncio.run(main())
```

### 2. 分布式锁（防重复）

如果多个进程可能运行相同任务，建议在 `before_run()` 中加锁：

```python
import redis.asyncio as redis
from async_scheduler import ScheduleTask


class DistributedTask(ScheduleTask):
    def __init__(self, task_id: str, redis_url: str) -> None:
        super().__init__(task_id)
        self.redis = redis.from_url(redis_url)
        self._lock = None

    async def before_run(self) -> None:
        # 尝试获取分布式锁
        self._lock = self.redis.lock(f"scheduler:{self.task_id}", timeout=300)
        if not await self._lock.acquire(blocking=False):
            raise RuntimeError(f"任务 {self.task_id} 已在其他节点运行")

    async def run(self) -> None:
        # 业务逻辑
        ...

    async def after_run(self, success: bool) -> None:
        if self._lock and self._lock.locked():
            await self._lock.release()
```

### 3. 监控与健康检查

利用钩子实现监控上报：

```python
import time
from async_scheduler import ScheduleTask


class MonitoredTask(ScheduleTask):
    def __init__(self, task_id: str) -> None:
        super().__init__(task_id)
        self.last_success_time: float | None = None
        self.consecutive_failures = 0

    async def run(self) -> None:
        # 业务逻辑
        start = time.monotonic()
        await self.do_work()
        duration = time.monotonic() - start
        # 上报指标
        await self.report_metrics(duration)

    async def do_work(self) -> None:
        ...

    async def report_metrics(self, duration: float) -> None:
        """上报到 Prometheus / Grafana 等"""
        ...

    async def after_run(self, success: bool) -> None:
        if success:
            self.last_success_time = time.monotonic()
            self.consecutive_failures = 0
        else:
            self.consecutive_failures += 1
            if self.consecutive_failures >= 3:
                await self.send_alert()

    async def send_alert(self) -> None:
        """连续失败 3 次发送告警"""
        ...
```

### 4. 测试策略

编写可靠的调度任务测试：

```python
import pytest
from async_scheduler import ScheduleTask


class TestableTask(ScheduleTask):
    def __init__(self) -> None:
        super().__init__("test-task", retry_times=2, retry_delay=0)
        self.call_count = 0

    async def run(self) -> None:
        self.call_count += 1
        if self.call_count < 3:
            raise RuntimeError("模拟失败")


@pytest.mark.asyncio
async def test_task_retries_on_failure() -> None:
    task = TestableTask()
    await task.safe_execute()

    # 第 1、2 次失败，第 3 次成功
    assert task.call_count == 3
```

## 性能考虑

### 调度精度

- **APScheduler 精度**：通常在毫秒级，受事件循环负载影响
- **Misfire 处理**：若上一次执行未完成，新触发会被跳过（`max_instances=1`）
- **建议**：对时间敏感的场景，使用独立的定时服务 + webhook 触发

### 内存占用

- 每个 `AsyncTaskScheduler` 创建一个 APScheduler 实例
- 大量任务时考虑共享调度器（需自行封装）

### 资源限制

- 默认单次执行超时 3600 秒（1 小时）
- 根据业务调整 `timeout` 参数，避免资源泄漏

## 许可证

[MIT](./LICENSE)
