fix: 避免服务重启后重复执行 Agent 定时任务 (#6302)

This commit is contained in:
InfinityPacer
2026-08-13 19:25:11 +08:00
committed by GitHub
parent dc9f8ef168
commit e014604b1c
6 changed files with 446 additions and 29 deletions

View File

@@ -109,6 +109,7 @@ class UpdateAgentTaskTool(MoviePilotTool):
if task.last_status == "running":
return {"error": f"Agent 定时任务 {payload.task_id} 正在执行,请稍后再修改"}
has_schedule_update = payload.trigger_type is not None
trigger_type = payload.trigger_type or task.trigger_type
trigger_value = payload.trigger
if trigger_type == "date" and payload.delay_minutes is not None:
@@ -120,20 +121,33 @@ class UpdateAgentTaskTool(MoviePilotTool):
trigger_value = (
task.cron_expression if trigger_type == "cron" else task.run_at
)
enabled = task.enabled if payload.enabled is None else payload.enabled
normalized_type, normalized_trigger = TimerUtils.normalize_schedule_trigger(
trigger_type=trigger_type,
trigger_value=trigger_value,
timezone_name=settings.TZ,
require_future=bool(enabled and trigger_type == "date"),
normalized_type = trigger_type
normalized_trigger = trigger_value
validates_existing_date_schedule = bool(
payload.enabled
and task.last_status != "interrupted"
and trigger_type == "date"
)
if has_schedule_update or validates_existing_date_schedule:
normalized_type, normalized_trigger = TimerUtils.normalize_schedule_trigger(
trigger_type=trigger_type,
trigger_value=trigger_value,
timezone_name=settings.TZ,
require_future=bool(
trigger_type == "date"
and (
task.last_status == "interrupted"
or (task.enabled if payload.enabled is None else payload.enabled)
)
),
)
update_payload = {}
if payload.name is not None:
update_payload["name"] = payload.name.strip()
if payload.content is not None:
update_payload["content"] = payload.content.strip()
if payload.trigger_type is not None:
if has_schedule_update:
update_payload.update(
{
"trigger_type": normalized_type,
@@ -147,7 +161,7 @@ class UpdateAgentTaskTool(MoviePilotTool):
)
if payload.enabled is not None:
update_payload["enabled"] = payload.enabled
if payload.enabled:
if payload.enabled and task.last_status != "interrupted":
update_payload["last_status"] = "waiting"
oper.update(

View File

@@ -104,6 +104,16 @@ class AgentTaskOper(DbOper):
run_at=self._now(),
)
def mark_interrupted(self, task_id: int, result: str) -> bool:
"""
将遗留的运行中任务标记为中断且结果未知。
"""
return AgentTask.mark_interrupted(
self._db,
task_id=task_id,
result=(result or "")[:20000],
)
def finish(
self,
task_id: int,

View File

@@ -36,6 +36,7 @@ class AgentTask(Base):
last_status = Column(String, nullable=False, default="waiting")
last_run_at = Column(String)
last_result = Column(Text)
# 已收口执行次数;进程中断的未完成尝试不计入
run_count = Column(Integer, nullable=False, default=0)
created_at = Column(String, nullable=False)
updated_at = Column(String, nullable=False)
@@ -146,6 +147,30 @@ class AgentTask(Base):
)
)
@classmethod
@db_update
def mark_interrupted(cls, db: Session, task_id: int, result: str) -> bool:
"""
将服务重启时遗留的运行中任务标记为结果未知。
该状态保留原执行时间和计数,避免把可能已经产生副作用的执行误记为
从未开始或完整失败。
"""
return bool(
db.query(cls)
.filter(
cls.id == task_id,
cls.last_status == "running",
)
.update(
{
"last_status": "interrupted",
"last_result": result,
"updated_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
}
)
)
@classmethod
@db_update
def finish_task(

View File

@@ -307,10 +307,14 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass):
self._lock = threading.RLock()
# 各服务的运行状态
self._jobs = {}
# 进程启动时只对账一次,配置热重载不得改写仍在执行的任务状态
self._agent_task_interruptions_reconciled = False
# 用户认证失败次数
self._auth_count = 0
# 用户认证失败消息发送
self._auth_message = False
# 对账上个进程未收口的 Agent 任务
self._reconcile_agent_task_interruptions()
# 初始化
self.init()
@@ -1066,20 +1070,33 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass):
def init_agent_task_jobs(self) -> None:
"""
数据库恢复所有启用的 Agent 自主定时任务。
数据库当前状态注册所有启用的 Agent 自主定时任务。
"""
oper = AgentTaskOper()
for task in oper.list(enabled=True):
if task.last_status == "running":
oper.update(
task_id=task.id,
payload={
"last_status": "waiting",
"last_result": "服务重启后恢复调度",
},
)
for task in AgentTaskOper().list(enabled=True):
self.update_agent_task_job(task.id)
def _reconcile_agent_task_interruptions(self) -> None:
"""
将上个进程未收口的 Agent 任务标记为结果未知。
配置变更会在同一进程内重建调度器,因此该对账在实例生命周期内只能
成功执行一次,避免把当前进程仍在运行的任务误判为中断。
"""
with self._lock:
if self._agent_task_interruptions_reconciled:
return
oper = AgentTaskOper()
for task in oper.list():
if task.last_status == "running":
oper.mark_interrupted(
task_id=task.id,
result=(
"服务重启时任务执行被中断,结果未知,可能已有部分操作;"
"请先检查实际状态,再决定是否重新执行"
),
)
self._agent_task_interruptions_reconciled = True
def update_agent_task_job(self, task_id: int) -> Optional[str]:
"""
按数据库中的最新配置新增或替换 Agent 自主定时任务。
@@ -1100,15 +1117,18 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass):
trigger_value = (
task.cron_expression if task.trigger_type == "cron" else task.run_at
)
try:
trigger = TimerUtils.build_schedule_trigger(
trigger_type=task.trigger_type,
trigger_value=trigger_value,
timezone_name=settings.TZ,
)
except (TypeError, ValueError) as err:
logger.error(f"Agent 定时任务 {task_id} 的触发配置无效:{str(err)}")
return None
manual_only = task.trigger_type == "date" and task.last_status == "interrupted"
trigger = None
if not manual_only:
try:
trigger = TimerUtils.build_schedule_trigger(
trigger_type=task.trigger_type,
trigger_value=trigger_value,
timezone_name=settings.TZ,
)
except (TypeError, ValueError) as err:
logger.error(f"Agent 定时任务 {task_id} 的触发配置无效:{str(err)}")
return None
job_id = self._get_agent_task_job_id(task_id)
with self._lock:
@@ -1119,6 +1139,10 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass):
"running": False,
"kwargs": {"task_id": task_id},
}
# 已开始的一次任务在重启后结果未知,只保留显式执行入口,不能按
# 过期触发时间自动重放可能已经发生的外部副作用。
if manual_only:
return None
self._scheduler.add_job(
self.start,
trigger=trigger,
@@ -1165,6 +1189,8 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass):
task = AgentTaskOper().get(task_id)
if not task or not task.enabled:
return None
if task.trigger_type == "date" and task.last_status == "interrupted":
return None
trigger_value = (
task.cron_expression if task.trigger_type == "cron" else task.run_at
)

View File

@@ -316,6 +316,8 @@ Agent 音乐流程与影视共用同一采集管线,但实体边界不同:
`trigger_type=date` 表示单次执行“30 分钟后检查”这类相对时间传 `delay_minutes=30`,由后端计算精确时间;固定时间则传 ISO 8601 `trigger`,支持精确到秒。`trigger_type=cron` 使用标准五段 cron分、时、日、月、周适合周期检查。未显式携带时区的时间按 MoviePilot 的 `TZ` 配置解释。任务由内存调度器精确触发,配置持久化到数据库,服务重启后会自动恢复;触发后 Agent 在原会话中执行 `content`,执行过程及最终结果均不绑定创建任务时的消息渠道,而是通过 MoviePilot 已配置的通知渠道广播。如果 Agent 在执行过程中已通过消息工具发送完整结果,任务结束时不会再次发送相同的最终回复。
服务重启时仍处于运行中的任务会显示为 `interrupted`,表示上次结果未知且可能已有部分操作。中断的一次任务不会自动补跑,暂停后恢复也仍保留中断状态;需要先核对实际结果,再用 `run_agent_task` 明确立即重跑,或通过 `update_agent_task` 提供新的 `trigger_type` 与未来触发时间重新安排。
Agent 自主任务工具使用数据库中的整数 `task_id``query_schedulers``run_scheduler` 仅面向系统、插件和工作流注册的运行时定时服务,使用字符串 `job_id`,不会返回或执行 `agent-task-*`。两类 ID 不可混用;需要立即执行自主任务时,应先通过 `query_agent_tasks` 确认归属和状态,再调用 `run_agent_task`。立即执行只提交任务,不在当前工具调用内等待结果,从而避免同一 Agent 会话互相等待;执行结果仍按上述通知规则广播。
上述过滤只约束 Agent 工具,避免模型混用两类任务。前端系统设置和仪表盘使用的 `/api/v1/dashboard/schedule` 仍返回完整运行时列表,其中包含 `provider=[Agent]` 的自主任务;前端通过 `/api/v1/system/runscheduler` 立即执行这类列表项的行为也保持不变。

View File

@@ -31,7 +31,10 @@ from app.agent.tools.impl.query_schedulers import QuerySchedulersTool
from app.agent.tools.impl.run_agent_task import RunAgentTaskTool
from app.agent.tools.impl.run_scheduler import RunSchedulerTool
from app.agent.tools.impl.send_message import SendMessageTool
from app.agent.tools.impl.update_agent_task import UpdateAgentTaskTool
from app.agent.tools.impl.update_agent_task import (
UpdateAgentTaskInput,
UpdateAgentTaskTool,
)
from app.agent.tools.tags import ToolTag
from app.core.config import settings
from app.db.agenttask_oper import AgentTaskOper
@@ -88,6 +91,49 @@ def _future_time(minutes: int = 10) -> str:
)
def _past_time(minutes: int = 10) -> str:
"""生成系统时区内的过去时间字符串。"""
timezone = pytz.timezone(settings.TZ)
return (datetime.now(timezone) - timedelta(minutes=minutes)).isoformat(
timespec="seconds"
)
def _invalid_time() -> str:
"""返回无法构造自动触发器的遗留时间值。"""
return "invalid-run-at"
def _add_agent_task(trigger_type: str, trigger_value: str, prefix: str):
"""创建带隔离用户上下文的 Agent 定时任务。"""
user_id = f"{prefix}-{uuid4().hex}"
return AgentTaskOper().add(
name=f"{prefix} 检查",
content="检查资源并报告",
trigger_type=trigger_type,
cron_expression=trigger_value if trigger_type == "cron" else None,
run_at=trigger_value if trigger_type == "date" else None,
user_id=user_id,
username="admin",
session_id=f"session-{user_id}",
channel="Telegram",
source="telegram-test",
original_chat_id="chat-1",
)
def _build_agent_task_scheduler(reconcile: bool = False) -> Scheduler:
"""构造不启动后台线程的 Agent 任务调度器。"""
scheduler = object.__new__(Scheduler)
scheduler._lock = threading.RLock()
scheduler._jobs = {}
scheduler._scheduler = BackgroundScheduler(timezone=settings.TZ)
scheduler._agent_task_interruptions_reconciled = False
if reconcile:
scheduler._reconcile_agent_task_interruptions()
return scheduler
def _build_tool(tool_class, user_id: str):
"""构造带当前用户消息上下文的 Agent 工具。"""
tool = tool_class(session_id=f"session-{user_id}", user_id=user_id)
@@ -291,6 +337,300 @@ def test_scheduler_registers_and_removes_agent_task_job() -> None:
scheduler._scheduler.shutdown(wait=False)
@pytest.mark.parametrize(
"run_time_factory",
[_past_time, _future_time, _invalid_time],
)
def test_scheduler_restart_keeps_interrupted_date_task_manual_only(
run_time_factory,
) -> None:
"""重启不得自动补跑结果未知的一次任务,但必须保留显式重跑入口。"""
task = _add_agent_task("date", run_time_factory(), "restart-date")
assert AgentTaskOper().mark_running(task.id)
scheduler = _build_agent_task_scheduler(reconcile=True)
scheduler.init_agent_task_jobs()
scheduler.init_agent_task_jobs()
recovered = AgentTaskOper().get(task.id)
job_id = scheduler._get_agent_task_job_id(task.id)
assert recovered.last_status == "interrupted"
assert "可能已有部分操作" in recovered.last_result
assert recovered.last_run_at is not None
assert recovered.run_count == 0
assert job_id in scheduler._jobs
assert scheduler._scheduler.get_job(job_id) is None
assert scheduler.get_agent_task_next_run(task.id) is None
scheduler.start = Mock()
assert scheduler.start_agent_task(task.id) is True
scheduler.start.assert_called_once_with(job_id)
@pytest.mark.anyio
async def test_interrupted_date_task_manual_run_disables_and_removes_job(
monkeypatch,
) -> None:
"""用户显式重跑中断的一次任务后,应按原有单次任务语义正常收口。"""
task = _add_agent_task("date", _past_time(), "restart-date-manual")
assert AgentTaskOper().mark_running(task.id)
scheduler = _build_agent_task_scheduler(reconcile=True)
scheduler.init_agent_task_jobs()
process_message = AsyncMock(return_value="执行完成")
monkeypatch.setattr("app.agent.agent_manager.process_message", process_message)
assert await scheduler.execute_agent_task(task.id) == (True, "执行完成")
process_message.assert_awaited_once()
finished = AgentTaskOper().get(task.id)
job_id = scheduler._get_agent_task_job_id(task.id)
assert finished.last_status == "success"
assert finished.run_count == 1
assert finished.enabled is False
assert job_id not in scheduler._jobs
@pytest.mark.parametrize("run_time_factory", [_future_time, _invalid_time])
@pytest.mark.anyio
async def test_interrupted_date_task_enable_toggle_stays_manual_only(
monkeypatch,
run_time_factory,
) -> None:
"""暂停或恢复中断的一次任务不得重新注册原自动触发。"""
task = _add_agent_task("date", run_time_factory(), "interrupted-toggle")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
assert oper.mark_interrupted(task.id, "执行结果未知")
scheduler = _build_agent_task_scheduler()
scheduler.init_agent_task_jobs()
monkeypatch.setattr("app.scheduler.Scheduler", lambda: scheduler)
tool = _build_tool(UpdateAgentTaskTool, task.user_id)
paused = json.loads(await tool.run(task_id=task.id, enabled=False))
assert paused["enabled"] is False
assert paused["last_status"] == "interrupted"
assert paused["next_run_at"] is None
resumed = json.loads(await tool.run(task_id=task.id, enabled=True))
job_id = scheduler._get_agent_task_job_id(task.id)
assert resumed["enabled"] is True
assert resumed["last_status"] == "interrupted"
assert resumed["next_run_at"] is None
assert job_id in scheduler._jobs
assert scheduler._scheduler.get_job(job_id) is None
@pytest.mark.anyio
async def test_interrupted_date_task_new_trigger_rearms_schedule(monkeypatch) -> None:
"""用户明确重设触发时间时,可以清除中断状态并重新安排单次任务。"""
task = _add_agent_task("date", _future_time(), "interrupted-reschedule")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
assert oper.mark_interrupted(task.id, "执行结果未知")
scheduler = _build_agent_task_scheduler()
scheduler.init_agent_task_jobs()
monkeypatch.setattr("app.scheduler.Scheduler", lambda: scheduler)
updated = json.loads(
await _build_tool(UpdateAgentTaskTool, task.user_id).run(
task_id=task.id,
trigger_type="date",
delay_minutes=20,
enabled=True,
)
)
job_id = scheduler._get_agent_task_job_id(task.id)
assert updated["last_status"] == "waiting"
assert updated["last_result"] is None
assert updated["next_run_at"] is not None
assert scheduler._scheduler.get_job(job_id) is not None
@pytest.mark.anyio
async def test_interrupted_date_task_rejects_past_trigger_while_pausing(
monkeypatch,
) -> None:
"""中断的一次任务即使同时暂停,也只能用新的未来时间离开中断状态。"""
task = _add_agent_task("date", _future_time(), "interrupted-past-reschedule")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
assert oper.mark_interrupted(task.id, "执行结果未知")
scheduler = _build_agent_task_scheduler()
scheduler.init_agent_task_jobs()
monkeypatch.setattr("app.scheduler.Scheduler", lambda: scheduler)
tool = _build_tool(UpdateAgentTaskTool, task.user_id)
with pytest.raises(ValueError, match="必须晚于当前时间"):
await tool.run(
task_id=task.id,
trigger_type="date",
trigger=_past_time(),
enabled=False,
)
unchanged = oper.get(task.id)
assert unchanged.enabled is True
assert unchanged.last_status == "interrupted"
assert unchanged.last_result == "执行结果未知"
def test_interrupted_date_task_requires_explicit_reschedule_value() -> None:
"""仅重复 date 类型不构成重新排期,必须同时提供新的触发值。"""
with pytest.raises(ValueError, match="必须提供 trigger 或 delay_minutes"):
UpdateAgentTaskInput(task_id=1, trigger_type="date")
@pytest.mark.anyio
async def test_expired_date_task_rejects_enable_without_reschedule(
monkeypatch,
) -> None:
"""普通单次任务恢复启用时仍需校验现有触发时间,避免过期补跑。"""
task = _add_agent_task("date", _past_time(), "expired-enable")
oper = AgentTaskOper()
assert oper.update(task_id=task.id, payload={"enabled": False})
scheduler = _build_agent_task_scheduler()
scheduler.init_agent_task_jobs()
monkeypatch.setattr("app.scheduler.Scheduler", lambda: scheduler)
with pytest.raises(ValueError, match="必须晚于当前时间"):
await _build_tool(UpdateAgentTaskTool, task.user_id).run(
task_id=task.id,
enabled=True,
)
unchanged = oper.get(task.id)
assert unchanged.enabled is False
assert unchanged.last_status == "waiting"
assert scheduler._get_agent_task_job_id(task.id) not in scheduler._jobs
def test_agent_task_interruption_transition_preserves_execution_evidence() -> None:
"""中断迁移只能消费 running 状态,并保留既有执行次数与开始时间。"""
task = _add_agent_task("cron", "0 * * * *", "interrupt-state")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
running = oper.get(task.id)
assert running.last_run_at is not None
assert oper.mark_interrupted(task.id, "执行结果未知")
assert not oper.mark_interrupted(task.id, "不得覆盖")
interrupted = oper.get(task.id)
assert interrupted.last_status == "interrupted"
assert interrupted.last_result == "执行结果未知"
assert interrupted.last_run_at == running.last_run_at
assert interrupted.run_count == running.run_count == 0
def test_scheduler_restart_keeps_unstarted_date_task_misfire_behavior() -> None:
"""停机期间仅错过触发时间的一次任务仍应按原有语义补跑。"""
task = _add_agent_task("date", _past_time(), "restart-waiting")
scheduler = _build_agent_task_scheduler()
scheduler.init_agent_task_jobs()
recovered = AgentTaskOper().get(task.id)
job_id = scheduler._get_agent_task_job_id(task.id)
runtime_job = scheduler._scheduler.get_job(job_id)
assert recovered.last_status == "waiting"
assert runtime_job is not None
assert runtime_job.misfire_grace_time is None
@pytest.mark.parametrize("success", [True, False])
def test_scheduler_restart_keeps_finished_date_task_misfire_behavior(
success: bool,
) -> None:
"""已有单次任务终态仍沿用原调度语义,不被中断策略扩大影响。"""
task = _add_agent_task("date", _past_time(), f"restart-finished-{success}")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
assert oper.finish(task.id, success=success, result="已收口")
scheduler = _build_agent_task_scheduler()
scheduler.init_agent_task_jobs()
recovered = oper.get(task.id)
runtime_job = scheduler._scheduler.get_job(
scheduler._get_agent_task_job_id(task.id)
)
assert recovered.last_status == ("success" if success else "failed")
assert runtime_job is not None
assert runtime_job.misfire_grace_time is None
@pytest.mark.anyio
async def test_scheduler_config_reload_does_not_interrupt_running_agent_task() -> None:
"""同进程重建调度时不得把仍在执行的任务误判为进程中断。"""
scheduler = _build_agent_task_scheduler(reconcile=True)
task = _add_agent_task("cron", "0 * * * *", "reload-running")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
scheduler._reconcile_agent_task_interruptions()
scheduler.init_agent_task_jobs()
running = oper.get(task.id)
assert running.last_status == "running"
assert not oper.mark_running(task.id)
manager = AgentManager()
manager.process_message = AsyncMock()
success, result = await manager.execute_scheduled_task(task.id)
assert success is False
assert result == "Agent 定时任务当前不可执行"
manager.process_message.assert_not_awaited()
def test_scheduler_restart_keeps_interrupted_cron_future_schedule() -> None:
"""周期任务中断后只保留下次正常调度,不抹掉本轮中断事实。"""
task = _add_agent_task("cron", "0 * * * *", "restart-cron")
assert AgentTaskOper().mark_running(task.id)
scheduler = _build_agent_task_scheduler(reconcile=True)
scheduler.init_agent_task_jobs()
scheduler.init_agent_task_jobs()
recovered = AgentTaskOper().get(task.id)
job_id = scheduler._get_agent_task_job_id(task.id)
runtime_job = scheduler._scheduler.get_job(job_id)
assert recovered.last_status == "interrupted"
assert "可能已有部分操作" in recovered.last_result
assert recovered.last_run_at is not None
assert recovered.run_count == 0
assert runtime_job is not None
assert scheduler.get_agent_task_next_run(task.id) is not None
scheduler.start = Mock()
assert scheduler.start_agent_task(task.id) is True
scheduler.start.assert_called_once_with(job_id)
def test_scheduler_restart_reconciles_disabled_running_agent_task() -> None:
"""停用状态不应掩盖上个进程遗留的运行中事实。"""
task = _add_agent_task("cron", "0 * * * *", "restart-disabled")
oper = AgentTaskOper()
assert oper.mark_running(task.id)
assert oper.update(task_id=task.id, payload={"enabled": False})
scheduler = _build_agent_task_scheduler(reconcile=True)
scheduler.init_agent_task_jobs()
recovered = oper.get(task.id)
assert recovered.last_status == "interrupted"
assert recovered.enabled is False
assert scheduler._get_agent_task_job_id(task.id) not in scheduler._jobs
def test_scheduler_starts_registered_agent_task_without_waiting() -> None:
"""Agent 任务立即执行入口应只调用运行时 start并拒绝缺失或运行中的任务。"""
scheduler = object.__new__(Scheduler)