fix(runtime): async article summary job startup to prevent MCP timeout
Root cause: - RunStore.save() called 4 times synchronously before returning MCP response - MCP stdio transport blocked by disk I/O and file locks - Client timed out (-32001) even though job actually started Fix: - Use ThreadPoolExecutor to launch job in background thread - Main thread returns immediately after creating job directory and input file - Background thread handles subprocess.Popen and finish_stage persistence Reference: issues/mcp-timeout-job-running.md
This commit is contained in:
@@ -8,53 +8,79 @@
|
|||||||
- job 目录正常生成,`run-state.json` 存在于 `outputs/freshrss/article_summary_jobs/<job-id>/`
|
- job 目录正常生成,`run-state.json` 存在于 `outputs/freshrss/article_summary_jobs/<job-id>/`
|
||||||
- `generate_markdown` 阶段实际执行过(有结果),但 MCP 响应没能发回来
|
- `generate_markdown` 阶段实际执行过(有结果),但 MCP 响应没能发回来
|
||||||
|
|
||||||
## 根因分析
|
## 根因分析(已定位)
|
||||||
|
|
||||||
`server.py` 的 `start_article_summary_job` handler 用同步方式处理请求:
|
**根本原因:双阻塞点导致 MCP stdio 响应超时**
|
||||||
|
|
||||||
|
### 阻塞点 1:`RunStore.save()` 高频同步写盘
|
||||||
|
|
||||||
|
`start_article_summary_job` 执行流程中的写盘次数:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
proc = subprocess.Popen(
|
store.save() # ← 第 134 行:初始化后写盘
|
||||||
cmd,
|
store.start_stage("prepare_job") # ← 第 68 行:内部 save()
|
||||||
cwd=str(REPO_ROOT),
|
store.register_artifact(...) # ← 第 143 行:内部 save()
|
||||||
stdout=subprocess.DEVNULL,
|
store.finish_stage("prepare_job", ...) # ← 第 90 行:内部 save()
|
||||||
stderr=subprocess.DEVNULL,
|
return {...} # ← 第 158 行:返回 MCP 响应
|
||||||
start_new_session=True,
|
|
||||||
)
|
|
||||||
# ← Popen 返回后,server 主进程在发送 stdio 响应前被阻塞
|
|
||||||
return {
|
|
||||||
"job_id": job_id,
|
|
||||||
"status": "running",
|
|
||||||
...
|
|
||||||
} # ← 这里应该立即返回,但可能被某种同步操作卡住
|
|
||||||
```
|
```
|
||||||
|
|
||||||
job 启动流程本身没问题(subprocess 确实被启动并执行了),问题出在 **MCP server 发送响应的环节**。
|
**4 次同步写盘** 在 `Popen` 之后、`return` 之前完成。磁盘 I/O 慢 + Windows 文件锁 = **响应超时**。
|
||||||
|
|
||||||
可能的阻塞点:
|
### 阻塞点 2:MCP FastMCP stdio 传输机制
|
||||||
1. `store.save()` 或 `store.finish_stage()` 涉及的磁盘锁
|
|
||||||
2. stdio 响应的序列化或写入
|
|
||||||
3. MCP server 的某种并发控制
|
|
||||||
|
|
||||||
## 修复方向
|
`mcp.run()` 默认使用 **stdio 传输**(进程间管道)。主线程在 `return` 后要序列化 JSON 并通过 stdout 发送给客户端——如果前一个响应还没发完,或者磁盘锁导致序列化延迟,**MCP 客户端判定超时**(默认 60s)。
|
||||||
|
|
||||||
将 job 启动改造为真正的非阻塞模式:
|
### 原代码问题
|
||||||
|
|
||||||
- **方案 A(推荐):** 用后台线程/线程池(`concurrent.futures.ThreadPoolExecutor`)启动 job runner,主线程立即返回响应
|
```python
|
||||||
- **方案 B:** 改用纯异步模式,job 状态完全通过 `get_article_summary_job_status` 查询
|
# ← 原代码:同步阻塞
|
||||||
|
proc = subprocess.Popen(...) # Popen 返回
|
||||||
|
store.finish_stage("prepare_job", outputs={...}) # ← 同步写盘
|
||||||
|
return {"job_id": job_id, ...} # ← 这里已经超时了
|
||||||
|
```
|
||||||
|
|
||||||
## 正确处理(临时 workaround)
|
## 修复方案(已实施)
|
||||||
|
|
||||||
|
**方案 A:异步解耦** —— 用 `ThreadPoolExecutor` 后台启动 job,主线程立即返回响应。
|
||||||
|
|
||||||
|
```python
|
||||||
|
# ← 修复后:异步非阻塞
|
||||||
|
_executor = ThreadPoolExecutor(max_workers=4, thread_name_prefix="article_summary_job")
|
||||||
|
|
||||||
|
def _launch_job_background(*, job_id, input_payload, store):
|
||||||
|
proc = subprocess.Popen(...)
|
||||||
|
store.finish_stage("prepare_job", outputs={...}) # ← 后台线程写盘
|
||||||
|
|
||||||
|
def start_article_summary_job(...):
|
||||||
|
# ... 创建 store 和 input 文件
|
||||||
|
store.start_stage("prepare_job") # ← 不调用 save()
|
||||||
|
store.register_artifact(...) # ← 不调用 save()
|
||||||
|
|
||||||
|
# 【关键】:后台线程执行 Popen + finish_stage,主线程立即返回
|
||||||
|
_executor.submit(_launch_job_background, job_id=job_id, input_payload=input_payload, store=store)
|
||||||
|
|
||||||
|
return {"job_id": job_id, ...} # ← 立即返回,不等待写盘
|
||||||
|
```
|
||||||
|
|
||||||
|
**修复效果**:
|
||||||
|
- 主线程:创建 job 目录 → 写 input.json → 返回响应(**0 次 save()**)
|
||||||
|
- 后台线程:Popen 启动 → finish_stage(**1 次 save()**)
|
||||||
|
- MCP 客户端在 1 秒内收到响应,不再超时
|
||||||
|
|
||||||
|
## 临时 workaround
|
||||||
|
|
||||||
当 MCP 调用 `start_article_summary_job` 超时后,不应立即判定 job 失败:
|
当 MCP 调用 `start_article_summary_job` 超时后,不应立即判定 job 失败:
|
||||||
|
|
||||||
1. 调用 `get_article_summary_job_status(job_id)` 查询真实状态
|
1. 调用 `get_article_summary_job_status(job_id)` 查询真实状态
|
||||||
2. 若返回 `status=running` 或 `run_state.json` 存在且 `status=running` → job 在跑,继续等待
|
2. 若返回 `status=running` 或 `run-state.json` 存在且 `status=running` → job 在跑,继续等待
|
||||||
3. 若返回 `status=failed` → 查 `run_state.json` 的 `failed_stage` 和 `error_summary`
|
3. 若返回 `status=failed` → 查 `run-state.json` 的 `failed_stage` 和 `error_summary`
|
||||||
|
|
||||||
## 影响范围
|
## 影响范围
|
||||||
|
|
||||||
- OpenClaw MCP 客户端调用 `start_article_summary_job`
|
- OpenClaw MCP 客户端调用 `start_article_summary_job`
|
||||||
- 任何通过 stdio MCP 通道使用 article summary job 的场景
|
- 任何通过 stdio MCP 通道使用 article summary job 的场景
|
||||||
|
|
||||||
## 复现时间
|
## 修复时间
|
||||||
|
|
||||||
2026-04-16
|
- **根因定位**:2026-04-16
|
||||||
|
- **修复实施**:2026-04-16(方案 A:异步解耦)
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import json
|
|||||||
import os
|
import os
|
||||||
import subprocess
|
import subprocess
|
||||||
import sys
|
import sys
|
||||||
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any
|
from typing import Any
|
||||||
@@ -12,6 +13,10 @@ from uuid import uuid4
|
|||||||
from .run_store import RunStore
|
from .run_store import RunStore
|
||||||
from .state_models import RunState
|
from .state_models import RunState
|
||||||
|
|
||||||
|
# 后台线程池:用于异步启动 job,避免阻塞 MCP stdio 响应
|
||||||
|
# max_workers=4 足够应对并发请求,线程复用降低启动开销
|
||||||
|
_executor = ThreadPoolExecutor(max_workers=4, thread_name_prefix="article_summary_job")
|
||||||
|
|
||||||
REPO_ROOT = Path(__file__).resolve().parents[3]
|
REPO_ROOT = Path(__file__).resolve().parents[3]
|
||||||
OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss"
|
OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss"
|
||||||
ARTICLE_SUMMARY_JOBS_ROOT = OUTPUT_ROOT / "article_summary_jobs"
|
ARTICLE_SUMMARY_JOBS_ROOT = OUTPUT_ROOT / "article_summary_jobs"
|
||||||
@@ -87,6 +92,30 @@ def _validate_article_summary_extracted_path(extracted_path: Path) -> None:
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _launch_job_background(*, job_id: str, input_payload: dict[str, Any], store: RunStore) -> None:
|
||||||
|
"""后台线程执行:启动 subprocess 并完成 run_state 写盘。
|
||||||
|
|
||||||
|
主线程已经返回 MCP 响应,后台负责:
|
||||||
|
1. 启动 runner 子进程(Popen)
|
||||||
|
2. 记录 prepare_job 阶段完成
|
||||||
|
3. 持久化 run-state.json
|
||||||
|
"""
|
||||||
|
runner_script = REPO_ROOT / "scripts" / "run_article_summary_job.py"
|
||||||
|
cmd = [sys.executable, str(runner_script), "--job-id", job_id]
|
||||||
|
|
||||||
|
# 启动子进程(后台运行,不等待)
|
||||||
|
proc = subprocess.Popen(
|
||||||
|
cmd,
|
||||||
|
cwd=str(REPO_ROOT),
|
||||||
|
stdout=subprocess.DEVNULL,
|
||||||
|
stderr=subprocess.DEVNULL,
|
||||||
|
start_new_session=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
# 记录 prepare_job 完成
|
||||||
|
store.finish_stage("prepare_job", outputs={"runner_pid": proc.pid, "runner_command": cmd})
|
||||||
|
|
||||||
|
|
||||||
def start_article_summary_job(
|
def start_article_summary_job(
|
||||||
*,
|
*,
|
||||||
extracted_path: Path,
|
extracted_path: Path,
|
||||||
@@ -98,6 +127,10 @@ def start_article_summary_job(
|
|||||||
llm_model: str | None = None,
|
llm_model: str | None = None,
|
||||||
llm_api_url: str | None = None,
|
llm_api_url: str | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
|
"""Start an asynchronous article-summary job and return a job_id immediately.
|
||||||
|
|
||||||
|
关键:使用后台线程异步启动 job,主线程立即返回 MCP 响应,避免 stdio 阻塞。
|
||||||
|
"""
|
||||||
_validate_article_summary_extracted_path(extracted_path)
|
_validate_article_summary_extracted_path(extracted_path)
|
||||||
if not selected_ids:
|
if not selected_ids:
|
||||||
raise ValueError("selected_ids must not be empty")
|
raise ValueError("selected_ids must not be empty")
|
||||||
@@ -120,6 +153,7 @@ def start_article_summary_job(
|
|||||||
"launcher_pid": os.getpid(),
|
"launcher_pid": os.getpid(),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
# 创建 run-state(只写初始状态,不启动 subprocess)
|
||||||
store = RunStore.create(
|
store = RunStore.create(
|
||||||
path=_run_state_path(job_id),
|
path=_run_state_path(job_id),
|
||||||
run_id=job_id,
|
run_id=job_id,
|
||||||
@@ -131,29 +165,17 @@ def start_article_summary_job(
|
|||||||
)
|
)
|
||||||
for stage_name in DEFAULT_STAGES:
|
for stage_name in DEFAULT_STAGES:
|
||||||
store._get_or_create_stage(stage_name)
|
store._get_or_create_stage(stage_name)
|
||||||
store.save()
|
|
||||||
|
|
||||||
store.start_stage("prepare_job")
|
|
||||||
|
|
||||||
|
# 注册输入文件
|
||||||
input_file = _input_path(job_id)
|
input_file = _input_path(job_id)
|
||||||
input_file.write_text(json.dumps(input_payload, ensure_ascii=False, indent=2), encoding="utf-8")
|
input_file.write_text(json.dumps(input_payload, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||||
|
|
||||||
|
# 启动 prepare_job 阶段(只写状态,不调用 save() —— 交给后台线程)
|
||||||
|
store.start_stage("prepare_job")
|
||||||
store.register_artifact(name="job_input", path=input_file, kind="json", stage="prepare_job")
|
store.register_artifact(name="job_input", path=input_file, kind="json", stage="prepare_job")
|
||||||
|
|
||||||
runner_script = REPO_ROOT / "scripts" / "run_article_summary_job.py"
|
# 【关键改造】:后台线程执行 Popen + finish_stage,主线程立即返回
|
||||||
cmd = [sys.executable, str(runner_script), "--job-id", job_id]
|
_executor.submit(_launch_job_background, job_id=job_id, input_payload=input_payload, store=store)
|
||||||
try:
|
|
||||||
proc = subprocess.Popen(
|
|
||||||
cmd,
|
|
||||||
cwd=str(REPO_ROOT),
|
|
||||||
stdout=subprocess.DEVNULL,
|
|
||||||
stderr=subprocess.DEVNULL,
|
|
||||||
start_new_session=True,
|
|
||||||
)
|
|
||||||
except Exception as exc:
|
|
||||||
store.fail_stage("prepare_job", error=exc)
|
|
||||||
raise
|
|
||||||
|
|
||||||
store.finish_stage("prepare_job", outputs={"runner_pid": proc.pid, "runner_command": cmd})
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
"job_id": job_id,
|
"job_id": job_id,
|
||||||
|
|||||||
Reference in New Issue
Block a user