From 5d855e9b92ec8854d70b16ccb6737d6862fc49da Mon Sep 17 00:00:00 2001 From: tiammomo <26957354+tiammomo@users.noreply.github.com> Date: Mon, 14 Sep 2026 14:09:27 +0800 Subject: [PATCH] feat(scheduled-tasks): filter run history by occurrence status (#5384) * feat(scheduled-tasks): filter run history by occurrence status Signed-off-by: tiammomo <26957354+tiammomo@users.noreply.github.com> * fix(scheduled-tasks): share occurrence status contract --------- Signed-off-by: tiammomo <26957354+tiammomo@users.noreply.github.com> Co-authored-by: Willem Jiang --- README.md | 6 + README_zh.md | 6 + backend/AGENTS.md | 3 +- .../app/gateway/routers/scheduled_tasks.py | 4 +- .../persistence/scheduled_task_runs/sql.py | 30 ++-- .../persistence/scheduled_tasks/model.py | 30 +++- .../test_scheduled_task_run_status_filter.py | 169 ++++++++++++++++++ 7 files changed, 232 insertions(+), 16 deletions(-) create mode 100644 backend/tests/test_scheduled_task_run_status_filter.py diff --git a/README.md b/README.md index 10faed9e0..8efe75ecb 100644 --- a/README.md +++ b/README.md @@ -1642,6 +1642,12 @@ Current MVP capabilities: - Execute scheduled work through the normal DeerFlow run lifecycle - Browse execution history in pages of 50; older pages pause automatic refresh, with an explicit return to the latest runs. Counts appear only after a successful read; loading and failed reads are not reported as zero runs. +**Filter execution history through the API** + +To inspect failures without downloading every successful occurrence, authenticated clients with `threads:read` can request `GET /api/scheduled-tasks/{task_id}/runs?status=failed&limit=50&offset=0` for an owned task. The optional `status` accepts `queued`, `launching`, `running`, `success`, `failed`, `skipped`, or `interrupted`; these are occurrence statuses, so task statuses such as `completed` are invalid (422). + +Filtering happens before pagination. `limit` (1–200, default 50) and `offset` (nonnegative, default 0) apply to matching records, ordered by creation time then ID, both descending. Omitting `status` preserves the existing mixed-history array response; no matches return `[]`. The API does not change task execution, and the workspace history UI remains unfiltered. + Current MVP limits: - No conversation-created `schedule_task` tool yet diff --git a/README_zh.md b/README_zh.md index 0e8c0e108..0e6da345c 100644 --- a/README_zh.md +++ b/README_zh.md @@ -840,6 +840,12 @@ DeerFlow 现在在 workspace 里内置了一个一等的定时任务(scheduled - 定时任务通过正常的 DeerFlow run 生命周期执行 - 按每页 50 条浏览执行历史;历史页暂停自动刷新,可随时返回最新记录。 仅在读取成功后显示条数,加载中或失败不会误显示为零条。 +**通过 API 筛选执行历史** + +排查失败记录时,无需先下载所有成功记录。已认证且具有 `threads:read` 权限的客户端,可以针对自己的任务请求 `GET /api/scheduled-tasks/{task_id}/runs?status=failed&limit=50&offset=0`。可选的 `status` 支持 `queued`、`launching`、`running`、`success`、`failed`、`skipped`、`interrupted`;这些是执行记录的状态,`completed` 等任务状态会被拒绝(422)。 + +筛选先于分页执行。`limit`(1–200,默认 50)和 `offset`(非负整数,默认 0)作用于匹配记录,按创建时间、ID 依次降序排列。不传 `status` 时保留原有的混合历史数组,无匹配项返回 `[]`。此 API 不改变任务执行行为,workspace 历史界面仍展示未筛选的记录。 + 当前 MVP 限制: - 暂时还没有可在对话中创建任务的 `schedule_task` 工具 diff --git a/backend/AGENTS.md b/backend/AGENTS.md index 4b7de7560..cbc892709 100644 --- a/backend/AGENTS.md +++ b/backend/AGENTS.md @@ -15,7 +15,8 @@ The backend runs a LangGraph-based super agent with sandbox execution, persisten - Gateway streams `write_file` and `str_replace` argument deltas in bounded batches for multi-mode `messages-tuple` consumers; single-mode message consumers retain the original per-chunk contract. Non-message frames flush pending batches, and `values` remains an optional complete-state snapshot rather than a prerequisite for batching. - With `stream_subgraphs`, subgraph frames keep their namespace in the SSE event name (`values|`, LangGraph Platform style) instead of impersonating root frames — a delegated subagent inherits the parent checkpoint namespace, so publishing its `values` snapshot as bare `values` replaces the whole thread view in SDK clients (#4399). Root-only consumers (file-tool chunk batcher, subagent event persistence, LLM error-fallback detection) ignore namespaced frames. The web frontend does not request subgraph streaming; subtask progress rides root-namespace `task_*` custom events. - Background subagent identity is deliberately split: the provider `tool_call_id` remains the correlation key for `ToolMessage`, `task_*` SSE events, persisted lifecycle events, frontend cards, and the public `ExtensionData.scope_id` contract (stored as `SubagentResult.external_task_id`), while `SubagentExecutor.execute_async()` generates a full server-side `execution_id` for `SubagentResult.task_id`, the process-wide registry, polling, cancellation, timeout handling, and cleanup. Provider IDs are not globally unique across parent runs, so they must never become registry ownership keys; scheduler closures retain their own `SubagentResult` rather than resolving ownership again through the mutable registry. Terminal subagent token usage travels in the current run's `ToolMessage.additional_kwargs` and is attributed from message state, never through a process-global provider-ID cache. -- Scheduled-task executions must reuse that same Gateway run lifecycle. The scheduler may decide *when* work runs, but it must dispatch through the existing run path rather than introducing a parallel execution stack. Scheduled launches pass `scheduler.recursion_limit` (default 1000, matching the web UI's `recursion_limit: 1000`, clamped by `max_recursion_limit`) via `launch_scheduled_thread_run`; the value is read from `get_app_config()` at dispatch, so a YAML edit applies to the next scheduled run without a Gateway restart. +- Scheduled tasks dispatch through the normal Gateway run path. `launch_scheduled_thread_run` reads `get_app_config()` at dispatch and passes `scheduler.recursion_limit` (default 1000, matching the web UI; clamped by `max_recursion_limit`), so YAML changes apply on the next run without restarting Gateway. +- Run-history `status` filters are occurrence states, not task states. `ScheduledTaskRunStatus` in `persistence/scheduled_tasks/model.py` is the shared API/repository vocabulary and must match the active and terminal occurrence-status sets. Keep owner lookup before reading history, and apply SQL task/status predicates before pagination; omitted status preserves the existing response. - The background scheduler is single-instance by default. `scheduler.multi_instance=true` opts into lease-aware recovery across Gateway instances and requires shared Postgres, `run_ownership.heartbeat_enabled=true`, and `run_events.backend=db`; otherwise startup rejects the configuration. Live scheduled runs are preserved when a peer starts; expired launch claims return to the durable queue, expired run leases are atomically taken over, stale launch writes are fenced by lease ownership, and the Postgres advisory-locked budget makes `max_concurrent_runs` a shared global cap for `launching`/`running` rows. - Long-running MCP work uses a separate durable task runtime (`McpTaskService` + `mcp_tasks`, lease-based recovery) rather than keeping remote task IDs or status polling inside the Agent loop; only submit remains Agent-visible, the database is the source of truth, and `ThreadState` receives only a bounded current-thread projection. Full contract (leases, cancellation fencing, delivery idempotency, management-tool exposure): [packages/harness/deerflow/mcp/AGENTS.md](packages/harness/deerflow/mcp/AGENTS.md). - MCP task notification retries, dead-lettering, and the cancel endpoint's worker-stopped 503 are part of that same contract — see [packages/harness/deerflow/mcp/AGENTS.md](packages/harness/deerflow/mcp/AGENTS.md). diff --git a/backend/app/gateway/routers/scheduled_tasks.py b/backend/app/gateway/routers/scheduled_tasks.py index 677f2e040..206e092aa 100644 --- a/backend/app/gateway/routers/scheduled_tasks.py +++ b/backend/app/gateway/routers/scheduled_tasks.py @@ -20,6 +20,7 @@ from app.gateway.deps import ( ) from deerflow.config.agents_config import AGENT_NAME_PATTERN, load_agent_config from deerflow.persistence.scheduled_tasks import ActiveScheduledTaskMutationConflict +from deerflow.persistence.scheduled_tasks.model import ScheduledTaskRunStatus from deerflow.scheduler.schedules import ( MAX_INTERVAL_SECONDS, normalize_cron_expression, @@ -486,6 +487,7 @@ async def list_scheduled_task_runs( request: Request, limit: int = Query(default=50, ge=1, le=200), offset: int = Query(default=0, ge=0), + status: ScheduledTaskRunStatus | None = None, ): task_repo = get_scheduled_task_repo(request) run_repo = get_scheduled_task_run_repo(request) @@ -495,7 +497,7 @@ async def list_scheduled_task_runs( task = await task_repo.get(task_id, user_id=str(user.id)) if task is None: raise HTTPException(status_code=404, detail="Scheduled task not found") - return await run_repo.list_by_task(task_id, limit=limit, offset=offset) + return await run_repo.list_by_task(task_id, limit=limit, offset=offset, status=status) @router.get("/threads/{thread_id}/scheduled-tasks") diff --git a/backend/packages/harness/deerflow/persistence/scheduled_task_runs/sql.py b/backend/packages/harness/deerflow/persistence/scheduled_task_runs/sql.py index 190d6364c..f9220bdb3 100644 --- a/backend/packages/harness/deerflow/persistence/scheduled_task_runs/sql.py +++ b/backend/packages/harness/deerflow/persistence/scheduled_task_runs/sql.py @@ -12,7 +12,12 @@ from deerflow.persistence.run import RunRepository from deerflow.persistence.run.model import RunRow from deerflow.persistence.scheduled_task_runs.model import ScheduledTaskRunRow from deerflow.persistence.scheduled_task_runs.projection import account_launch, can_project -from deerflow.persistence.scheduled_tasks.model import ACTIVE_RUN_STATUSES, TERMINAL_RUN_STATUSES, ScheduledTaskRow +from deerflow.persistence.scheduled_tasks.model import ( + ACTIVE_RUN_STATUSES, + TERMINAL_RUN_STATUSES, + ScheduledTaskRow, + ScheduledTaskRunStatus, +) from deerflow.scheduler.schedules import next_run_at as compute_next_run_at from deerflow.utils.time import coerce_iso @@ -224,17 +229,18 @@ class ScheduledTaskRunRepository: await session.refresh(row) return self._row_to_dict(row) - async def list_by_task(self, task_id: str, *, limit: int = 50, offset: int = 0) -> list[dict[str, Any]]: - stmt = ( - select(ScheduledTaskRunRow) - .where(ScheduledTaskRunRow.task_id == task_id) - .order_by( - ScheduledTaskRunRow.created_at.desc(), - ScheduledTaskRunRow.id.desc(), - ) - .limit(limit) - .offset(offset) - ) + async def list_by_task( + self, + task_id: str, + *, + limit: int = 50, + offset: int = 0, + status: ScheduledTaskRunStatus | None = None, + ) -> list[dict[str, Any]]: + stmt = select(ScheduledTaskRunRow).where(ScheduledTaskRunRow.task_id == task_id) + if status is not None: + stmt = stmt.where(ScheduledTaskRunRow.status == status) + stmt = stmt.order_by(ScheduledTaskRunRow.created_at.desc(), ScheduledTaskRunRow.id.desc()).limit(limit).offset(offset) async with self._sf() as session: result = await session.execute(stmt) return [self._row_to_dict(row) for row in result.scalars()] diff --git a/backend/packages/harness/deerflow/persistence/scheduled_tasks/model.py b/backend/packages/harness/deerflow/persistence/scheduled_tasks/model.py index 48b9f8f7b..feaa088ee 100644 --- a/backend/packages/harness/deerflow/persistence/scheduled_tasks/model.py +++ b/backend/packages/harness/deerflow/persistence/scheduled_tasks/model.py @@ -1,6 +1,7 @@ from __future__ import annotations from datetime import UTC, datetime +from enum import StrEnum from sqlalchemy import JSON, BigInteger, DateTime, Integer, String, Text from sqlalchemy.orm import Mapped, mapped_column @@ -40,11 +41,36 @@ class ScheduledTaskRow(Base): ) +class ScheduledTaskRunStatus(StrEnum): + """Canonical status vocabulary for scheduled-task occurrences.""" + + QUEUED = "queued" + LAUNCHING = "launching" + RUNNING = "running" + SUCCESS = "success" + FAILED = "failed" + SKIPPED = "skipped" + INTERRUPTED = "interrupted" + + # Status constants - shared between scheduled_tasks and scheduled_task_runs # to avoid circular import and ensure consistency. # Import these from deerflow.persistence.scheduled_tasks.model in both modules. -TERMINAL_RUN_STATUSES: frozenset[str] = frozenset({"success", "failed", "skipped", "interrupted"}) -ACTIVE_RUN_STATUSES: frozenset[str] = frozenset({"queued", "launching", "running"}) +TERMINAL_RUN_STATUSES: frozenset[str] = frozenset( + { + ScheduledTaskRunStatus.SUCCESS, + ScheduledTaskRunStatus.FAILED, + ScheduledTaskRunStatus.SKIPPED, + ScheduledTaskRunStatus.INTERRUPTED, + } +) +ACTIVE_RUN_STATUSES: frozenset[str] = frozenset( + { + ScheduledTaskRunStatus.QUEUED, + ScheduledTaskRunStatus.LAUNCHING, + ScheduledTaskRunStatus.RUNNING, + } +) # Parent ``once`` task status projected from a terminal occurrence status. # Shared by the completion path and both recovery paths so the mapping diff --git a/backend/tests/test_scheduled_task_run_status_filter.py b/backend/tests/test_scheduled_task_run_status_filter.py new file mode 100644 index 000000000..92b9464cd --- /dev/null +++ b/backend/tests/test_scheduled_task_run_status_filter.py @@ -0,0 +1,169 @@ +from datetime import UTC, datetime, timedelta +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import httpx +import pytest +import pytest_asyncio +from fastapi import FastAPI +from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine + +from app.gateway.authz import AuthContext +from app.gateway.routers import scheduled_tasks +from deerflow.persistence.scheduled_task_runs import ScheduledTaskRunRepository +from deerflow.persistence.scheduled_task_runs.model import ScheduledTaskRunRow +from deerflow.persistence.scheduled_tasks import ScheduledTaskRepository +from deerflow.persistence.scheduled_tasks.model import ( + ACTIVE_RUN_STATUSES, + TERMINAL_RUN_STATUSES, + ScheduledTaskRow, + ScheduledTaskRunStatus, +) + +URL = "/api/scheduled-tasks/task-1/runs" +NOW = datetime(2026, 9, 12, tzinfo=UTC) + + +def occurrence(record_id, status, *, task_id="task-1", created_at=NOW): + return ScheduledTaskRunRow(id=record_id, task_id=task_id, thread_id="thread-1", scheduled_for=created_at, trigger="scheduled", status=status, created_at=created_at) + + +@pytest_asyncio.fixture +async def history(monkeypatch): + engine = create_async_engine("sqlite+aiosqlite:///:memory:") + try: + async with engine.begin() as connection: + await connection.run_sync(ScheduledTaskRow.__table__.create) + await connection.run_sync(ScheduledTaskRunRow.__table__.create) + sf = async_sessionmaker(engine, expire_on_commit=False) + async with sf() as session: + for task_id, owner in [("task-1", "user-1"), ("task-other", "user-1"), ("task-foreign", "user-2"), ("task-empty", "user-1")]: + session.add(ScheduledTaskRow(id=task_id, user_id=owner, title="History example", prompt="Example", schedule_type="cron", schedule_spec={"cron": "0 9 * * *"}, timezone="UTC")) + session.add_all([occurrence(f"success-{i:02}", "success", created_at=NOW + timedelta(hours=1, seconds=i)) for i in range(60)]) + session.add_all([occurrence(f"failed-{letter}", "failed") for letter in "abc"]) + session.add_all( + [ + occurrence("skipped-old", "skipped", created_at=NOW - timedelta(days=1)), + occurrence("interrupted-old", "interrupted", created_at=NOW - timedelta(days=2)), + occurrence("other-failure", "failed", task_id="task-other", created_at=NOW + timedelta(days=1)), + occurrence("foreign-failure", "failed", task_id="task-foreign", created_at=NOW + timedelta(days=1)), + ] + ) + await session.commit() + task_repo = ScheduledTaskRepository(sf) + run_repo = ScheduledTaskRunRepository(sf) + spy = AsyncMock(wraps=run_repo.list_by_task) + monkeypatch.setattr(run_repo, "list_by_task", spy) + monkeypatch.setattr(scheduled_tasks, "get_scheduled_task_repo", lambda request: task_repo) + monkeypatch.setattr(scheduled_tasks, "get_scheduled_task_run_repo", lambda request: run_repo) + + async def user_from_request(request): + return request.state.auth.user + + monkeypatch.setattr(scheduled_tasks, "get_optional_user_from_request", user_from_request) + app = FastAPI() + app.include_router(scheduled_tasks.router) + + @app.middleware("http") + async def authenticate(request, call_next): + user = None if request.headers.get("x-test-auth") == "anonymous" else SimpleNamespace(id="user-1") + permissions = [] if request.headers.get("x-test-auth") == "denied" else ["threads:read"] + request.state.auth = AuthContext(user=user, permissions=permissions) + return await call_next(request) + + async with httpx.AsyncClient(transport=httpx.ASGITransport(app=app), base_url="http://test") as client: + yield SimpleNamespace(client=client, sf=sf, repo=run_repo, spy=spy) + finally: + await engine.dispose() + + +@pytest.mark.asyncio +async def test_filter_finds_older_failures_and_paginates_matching_rows(history): + response = await history.client.get(URL, params={"status": "failed", "limit": 2}) + assert response.status_code == 200 + assert [row["id"] for row in response.json()] == ["failed-c", "failed-b"] + response = await history.client.get(URL, params={"status": "failed", "limit": 2, "offset": 2}) + assert [row["id"] for row in response.json()] == ["failed-a"] + response = await history.client.get(URL, params={"status": "failed", "offset": 3}) + assert response.json() == [] + + +@pytest.mark.asyncio +async def test_repository_applies_status_before_limit_and_offset(history): + rows = await history.repo.list_by_task("task-1", status="failed", limit=1, offset=1) + assert [row["id"] for row in rows] == ["failed-b"] + + +@pytest.mark.asyncio +async def test_omitted_status_preserves_mixed_history_and_pagination(history): + response = await history.client.get(URL, params={"limit": 200}) + assert response.status_code == 200 + rows = response.json() + assert len(rows) == 65 + assert [row["id"] for row in rows[:2]] == ["success-59", "success-58"] + assert {row["status"] for row in rows} == {"success", "failed", "skipped", "interrupted"} + assert all(row["task_id"] == "task-1" for row in rows) + response = await history.client.get(URL, params={"limit": 2, "offset": 60}) + assert [row["id"] for row in response.json()] == ["failed-c", "failed-b"] + assert len((await history.client.get(URL)).json()) == 50 + + +@pytest.mark.asyncio +@pytest.mark.parametrize("status", list(ScheduledTaskRunStatus)) +async def test_each_occurrence_status_is_supported(history, status): + async with history.sf() as session: + session.add(occurrence("new-occurrence", status, created_at=NOW + timedelta(days=2))) + await session.commit() + response = await history.client.get(URL, params={"status": status, "limit": 200}) + assert response.status_code == 200 + rows = response.json() + assert rows[0]["id"] == "new-occurrence" + assert all(row["status"] == status and row["task_id"] == "task-1" for row in rows) + + +@pytest.mark.asyncio +async def test_occurrence_status_contract_is_shared_with_openapi(history): + statuses = {status.value for status in ScheduledTaskRunStatus} + assert statuses == ACTIVE_RUN_STATUSES | TERMINAL_RUN_STATUSES + + response = await history.client.get("/openapi.json") + assert response.status_code == 200 + assert set(response.json()["components"]["schemas"]["ScheduledTaskRunStatus"]["enum"]) == statuses + + +@pytest.mark.asyncio +@pytest.mark.parametrize("status", ["", "completed", "error", "FAILED", "failed,success"]) +async def test_unknown_status_is_rejected_before_reading_history(history, status): + response = await history.client.get(URL, params={"status": status}) + assert response.status_code == 422 + history.spy.assert_not_awaited() + + +@pytest.mark.asyncio +@pytest.mark.parametrize("task_id", ["task-foreign", "missing-task"]) +async def test_filter_cannot_read_another_owner_or_missing_task(history, task_id): + response = await history.client.get(f"/api/scheduled-tasks/{task_id}/runs", params={"status": "failed"}) + assert response.status_code == 404 + history.spy.assert_not_awaited() + + +@pytest.mark.asyncio +@pytest.mark.parametrize(("auth", "status"), [("anonymous", 401), ("denied", 403)]) +async def test_filter_keeps_authentication_and_read_permission_checks(history, auth, status): + response = await history.client.get(URL, params={"status": "failed"}, headers={"x-test-auth": auth}) + assert response.status_code == status + history.spy.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_empty_or_unmatched_history_returns_empty_array(history): + assert (await history.client.get("/api/scheduled-tasks/task-empty/runs", params={"status": "failed"})).json() == [] + assert (await history.client.get(URL, params={"status": "running"})).json() == [] + + +@pytest.mark.asyncio +@pytest.mark.parametrize("params", [{"limit": 0}, {"limit": 201}, {"offset": -1}]) +async def test_filter_preserves_pagination_bounds(history, params): + response = await history.client.get(URL, params={"status": "failed", **params}) + assert response.status_code == 422 + history.spy.assert_not_awaited()