419 lines
16 KiB
Python
419 lines
16 KiB
Python
"""Public (unauthenticated) API for scheduled tasks.
|
|
|
|
Mounted under the ``/api/public/`` prefix so the auth middleware lets these
|
|
routes through without a session cookie. Intended for embedding scheduled
|
|
task results into external dashboards / documentation / chat systems that
|
|
shouldn't need a DeerFlow login.
|
|
|
|
Surface area:
|
|
|
|
- ``GET /api/public/scheduled-tasks/`` — list every task (no prompt body)
|
|
- ``GET /api/public/scheduled-tasks/{task_id}`` — task metadata
|
|
- ``GET /api/public/scheduled-tasks/{task_id}/runs`` — runs for a task (latest first)
|
|
- ``GET /api/public/scheduled-tasks/runs/{run_id}`` — single run (result included)
|
|
- ``GET /api/public/scheduled-tasks/runs/{run_id}/files/{i}``— inline file content / preview
|
|
- ``PUT /api/public/scheduled-tasks/runs/{run_id}/files/{i}``— overwrite an editable result file
|
|
- ``PUT /api/public/scheduled-tasks/runs/{run_id}/body`` — overwrite the inline Markdown body
|
|
- ``GET /api/public/scheduled-tasks/runs/{run_id}/download`` — zip / md download of result
|
|
- ``DELETE /api/public/scheduled-tasks/runs/{run_id}`` — remove a single run (no auth)
|
|
|
|
The PUT endpoints intentionally skip ownership checks so the anonymous public
|
|
viewer at ``#/public/scheduled-tasks/<task_id>`` can fix typos in result MD
|
|
files / body content. Tighten this if the deployment requires read-only
|
|
sharing.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import mimetypes
|
|
import re
|
|
import zipfile
|
|
from io import BytesIO
|
|
from pathlib import Path
|
|
from typing import Any
|
|
from urllib.parse import quote
|
|
|
|
from fastapi import APIRouter, HTTPException, Query
|
|
from fastapi.responses import FileResponse, Response, StreamingResponse
|
|
from pydantic import BaseModel, Field
|
|
|
|
from deerflow.runtime.scheduler import get_scheduled_task_service
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Whitelisted under ``auth_middleware._PUBLIC_PATH_PREFIXES`` (``/api/public/``).
|
|
router = APIRouter(prefix="/api/public/scheduled-tasks", tags=["public-scheduled-tasks"])
|
|
|
|
|
|
_MAX_PUBLIC_EDIT_BYTES = 5 * 1024 * 1024 # 5 MB — mirrors the authenticated edit cap
|
|
|
|
|
|
class PublicRunFileUpdateRequest(BaseModel):
|
|
content: str = Field(..., description="New text content to write into the run file")
|
|
|
|
|
|
class PublicRunBodyUpdateRequest(BaseModel):
|
|
content: str = Field(..., description="New Markdown body for the run's result")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sanitisation helpers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# Fields stripped from every task before it leaves the public boundary.
|
|
# Keeps owner ids, raw cron, and the (often verbose) author prompt out of
|
|
# anonymous responses.
|
|
_TASK_PRIVATE_FIELDS: frozenset[str] = frozenset(
|
|
{
|
|
"user_id",
|
|
"owner_user_id",
|
|
"creator_user_id",
|
|
"creator_email",
|
|
"prompt",
|
|
"execution_context",
|
|
"_subscription",
|
|
"role",
|
|
"scheduler_thread_id",
|
|
}
|
|
)
|
|
|
|
# Fields stripped from every run before it leaves the public boundary.
|
|
# We keep the result body (the user explicitly wants the MD content surfaced)
|
|
# but redact any path that points at a server-side filesystem location.
|
|
_RUN_PRIVATE_FIELDS: frozenset[str] = frozenset(
|
|
{
|
|
"user_id",
|
|
"owner_user_id",
|
|
"scheduler_thread_id",
|
|
"agent_run_id",
|
|
"can_edit",
|
|
}
|
|
)
|
|
|
|
|
|
def _redact_task(task: dict[str, Any]) -> dict[str, Any]:
|
|
return {k: v for k, v in task.items() if k not in _TASK_PRIVATE_FIELDS}
|
|
|
|
|
|
def _redact_run(run: dict[str, Any]) -> dict[str, Any]:
|
|
cleaned = {k: v for k, v in run.items() if k not in _RUN_PRIVATE_FIELDS}
|
|
result = cleaned.get("result")
|
|
if isinstance(result, dict):
|
|
files = result.get("files")
|
|
if isinstance(files, list):
|
|
cleaned_files = []
|
|
for idx, item in enumerate(files):
|
|
if not isinstance(item, dict):
|
|
continue
|
|
# Strip absolute filesystem path; clients use the file index
|
|
# to fetch content via the public file endpoint.
|
|
public_file = {k: v for k, v in item.items() if k != "path"}
|
|
public_file["index"] = idx
|
|
# Mark whether the underlying archive is still on disk. The
|
|
# public viewer filters out tabs whose source is gone so users
|
|
# never see the raw "Result file is no longer available" 404.
|
|
path_value = item.get("path")
|
|
public_file["available"] = (
|
|
isinstance(path_value, str)
|
|
and bool(path_value)
|
|
and Path(path_value).is_file()
|
|
)
|
|
cleaned_files.append(public_file)
|
|
result = {**result, "files": cleaned_files}
|
|
cleaned["result"] = result
|
|
return cleaned
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Internal helpers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _safe_filename(value: str, default: str = "scheduled-task-result") -> str:
|
|
name = re.sub(r'[\\/:*?"<>|\r\n]+', "-", value).strip(" .")
|
|
return name or default
|
|
|
|
|
|
def _result_markdown(result: dict[str, Any]) -> str:
|
|
title = str(result.get("title") or "Scheduled Task Result")
|
|
content = str(result.get("content") or "")
|
|
return f"# {title}\n\n{content}".strip() + "\n"
|
|
|
|
|
|
def _utf8_sig_bytes(content: str) -> bytes:
|
|
return content.encode("utf-8-sig")
|
|
|
|
|
|
def _content_disposition(filename: str) -> str:
|
|
quoted = quote(filename)
|
|
return f"attachment; filename*=UTF-8''{quoted}"
|
|
|
|
|
|
def _run_file_path(run: dict[str, Any], file_index: int) -> tuple[dict[str, Any], Path]:
|
|
result = run.get("result") or {}
|
|
files = result.get("files") if isinstance(result, dict) and isinstance(result.get("files"), list) else []
|
|
if file_index < 0 or file_index >= len(files):
|
|
raise HTTPException(status_code=404, detail="Result file not found")
|
|
item = files[file_index]
|
|
if not isinstance(item, dict) or not isinstance(item.get("path"), str):
|
|
raise HTTPException(status_code=404, detail="Result file not found")
|
|
path = Path(item["path"])
|
|
if not path.exists() or not path.is_file():
|
|
raise HTTPException(status_code=404, detail="Result file is no longer available")
|
|
return item, path
|
|
|
|
|
|
def _is_text_previewable(path: Path, mime_type: str | None) -> bool:
|
|
if mime_type and (mime_type.startswith("text/") or mime_type in {"application/json", "application/xml"}):
|
|
return True
|
|
return path.suffix.lower() in {".md", ".markdown", ".txt", ".json", ".csv", ".xml", ".yaml", ".yml", ".log"}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Routes
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@router.get("")
|
|
async def public_list_tasks(
|
|
search: str | None = Query(default=None, description="Case-insensitive fuzzy match on task name"),
|
|
limit: int = Query(default=200, ge=1, le=500),
|
|
offset: int = Query(default=0, ge=0),
|
|
) -> dict[str, Any]:
|
|
"""Return every scheduled task across all users.
|
|
|
|
Sensitive fields (user_id, raw prompt, scheduler thread id, …) are stripped
|
|
before serialisation. The response shape is intentionally stable so external
|
|
integrations can pin to it without surprises.
|
|
"""
|
|
service = get_scheduled_task_service()
|
|
tasks = await service.admin_list_tasks()
|
|
|
|
if search:
|
|
needle = search.strip().lower()
|
|
tasks = [t for t in tasks if needle in str(t.get("name", "")).lower()]
|
|
|
|
total = len(tasks)
|
|
sliced = tasks[offset : offset + limit]
|
|
return {
|
|
"total": total,
|
|
"limit": limit,
|
|
"offset": offset,
|
|
"tasks": [_redact_task(task) for task in sliced],
|
|
}
|
|
|
|
|
|
@router.get("/{task_id}")
|
|
async def public_get_task(task_id: str) -> dict[str, Any]:
|
|
"""Return metadata for a single task (without prompt body)."""
|
|
service = get_scheduled_task_service()
|
|
task = await service.store.get_task_any(task_id)
|
|
if task is None:
|
|
raise HTTPException(status_code=404, detail=f"Task {task_id} not found")
|
|
return _redact_task(task)
|
|
|
|
|
|
@router.get("/{task_id}/runs")
|
|
async def public_list_task_runs(
|
|
task_id: str,
|
|
limit: int = Query(default=50, ge=1, le=200),
|
|
) -> dict[str, Any]:
|
|
"""Return runs for a task, latest first."""
|
|
service = get_scheduled_task_service()
|
|
task = await service.store.get_task_any(task_id)
|
|
if task is None:
|
|
raise HTTPException(status_code=404, detail=f"Task {task_id} not found")
|
|
runs = await service.admin_list_runs(task_id, limit=limit)
|
|
return {
|
|
"task": _redact_task(task),
|
|
"total": len(runs),
|
|
"runs": [_redact_run(run) for run in runs],
|
|
}
|
|
|
|
|
|
@router.get("/runs/{run_id}")
|
|
async def public_get_run(run_id: str) -> dict[str, Any]:
|
|
"""Return a single run with its full result body (Markdown included)."""
|
|
service = get_scheduled_task_service()
|
|
run = await service.admin_get_run(run_id)
|
|
if run is None:
|
|
raise HTTPException(status_code=404, detail=f"Run {run_id} not found")
|
|
return _redact_run(run)
|
|
|
|
|
|
def _html_payload(run_id: str, page: dict[str, Any]) -> dict[str, Any]:
|
|
return {
|
|
"run_id": run_id,
|
|
"task_id": page.get("task_id"),
|
|
"title": page.get("title"),
|
|
"html": page.get("html_content"),
|
|
"style_summary": page.get("style_summary"),
|
|
"layout_summary": page.get("layout_summary"),
|
|
"page_id": page.get("id"),
|
|
"created_at": page.get("created_at"),
|
|
}
|
|
|
|
|
|
@router.get("/runs/{run_id}/html")
|
|
async def public_get_run_html(run_id: str) -> dict[str, Any]:
|
|
"""Return the generated HTML document for an ``html_page`` run (anonymous)."""
|
|
service = get_scheduled_task_service()
|
|
page = await service.get_html_page_for_run(run_id)
|
|
if page is None:
|
|
raise HTTPException(status_code=404, detail="该执行记录没有 HTML 页面")
|
|
return _html_payload(run_id, page)
|
|
|
|
|
|
@router.get("/{task_id}/latest-html")
|
|
async def public_get_latest_html(task_id: str) -> dict[str, Any]:
|
|
"""Return the most recent generated HTML page for a task (anonymous)."""
|
|
service = get_scheduled_task_service()
|
|
page = await service.get_latest_html_page_for_task(task_id)
|
|
if page is None:
|
|
raise HTTPException(status_code=404, detail="该任务暂无 HTML 页面")
|
|
return _html_payload(str(page.get("run_id") or ""), page)
|
|
|
|
|
|
@router.get("/runs/{run_id}/files/{file_index}")
|
|
async def public_get_run_file(run_id: str, file_index: int) -> dict[str, Any]:
|
|
"""Return a single result-file entry with its inline text content where
|
|
previewable. Binary files come back with ``previewable: false`` and the
|
|
consumer should fetch the download endpoint instead."""
|
|
service = get_scheduled_task_service()
|
|
run = await service.admin_get_run(run_id)
|
|
if run is None:
|
|
raise HTTPException(status_code=404, detail=f"Run {run_id} not found")
|
|
item, path = _run_file_path(run, file_index)
|
|
mime_type = item.get("mime_type") or mimetypes.guess_type(path.name)[0]
|
|
previewable = _is_text_previewable(path, mime_type)
|
|
content: str | None = None
|
|
if previewable:
|
|
try:
|
|
content = path.read_text(encoding="utf-8")
|
|
except UnicodeDecodeError:
|
|
content = path.read_text(encoding="utf-8", errors="replace")
|
|
return {
|
|
"index": file_index,
|
|
"name": item.get("name") or path.name,
|
|
"mime_type": mime_type,
|
|
"kind": item.get("kind"),
|
|
"size": path.stat().st_size,
|
|
"previewable": previewable,
|
|
"content": content,
|
|
}
|
|
|
|
|
|
@router.get("/runs/{run_id}/files/{file_index}/download")
|
|
async def public_download_run_file(run_id: str, file_index: int):
|
|
"""Stream a single result file as a download."""
|
|
service = get_scheduled_task_service()
|
|
run = await service.admin_get_run(run_id)
|
|
if run is None:
|
|
raise HTTPException(status_code=404, detail=f"Run {run_id} not found")
|
|
item, path = _run_file_path(run, file_index)
|
|
filename = _safe_filename(str(item.get("name") or path.name), path.name)
|
|
media = item.get("mime_type") or mimetypes.guess_type(path.name)[0] or "application/octet-stream"
|
|
return FileResponse(path, filename=filename, media_type=media)
|
|
|
|
|
|
@router.get("/runs/{run_id}/download")
|
|
async def public_download_run(run_id: str):
|
|
"""Download the full result: a zipped bundle of result.md + every file,
|
|
or a single ``.md`` when no auxiliary files were produced."""
|
|
service = get_scheduled_task_service()
|
|
run = await service.admin_get_run(run_id)
|
|
if run is None:
|
|
raise HTTPException(status_code=404, detail=f"Run {run_id} not found")
|
|
|
|
result = run.get("result") or {}
|
|
if not isinstance(result, dict):
|
|
result = {}
|
|
files = result.get("files") if isinstance(result.get("files"), list) else []
|
|
accessible_files: list[tuple[str, Path]] = []
|
|
for item in files:
|
|
if not isinstance(item, dict):
|
|
continue
|
|
path_value = item.get("path")
|
|
if not isinstance(path_value, str) or not path_value:
|
|
continue
|
|
path = Path(path_value)
|
|
if path.exists() and path.is_file():
|
|
accessible_files.append((str(item.get("name") or path.name), path))
|
|
|
|
base_name = _safe_filename(str(result.get("title") or "scheduled-task-result"))
|
|
if not accessible_files:
|
|
return Response(
|
|
content=_utf8_sig_bytes(_result_markdown(result)),
|
|
media_type="text/markdown; charset=utf-8",
|
|
headers={"Content-Disposition": _content_disposition(f"{base_name}.md")},
|
|
)
|
|
|
|
archive = BytesIO()
|
|
with zipfile.ZipFile(archive, mode="w", compression=zipfile.ZIP_DEFLATED) as zf:
|
|
zf.writestr("result.md", _utf8_sig_bytes(_result_markdown(result)))
|
|
for name, path in accessible_files:
|
|
arcname = _safe_filename(name, path.name)
|
|
zf.write(path, arcname=arcname)
|
|
archive.seek(0)
|
|
return StreamingResponse(
|
|
archive,
|
|
media_type="application/zip",
|
|
headers={"Content-Disposition": _content_disposition(f"{base_name}.zip")},
|
|
)
|
|
|
|
|
|
@router.put("/runs/{run_id}/files/{file_index}")
|
|
async def public_update_run_file(
|
|
run_id: str,
|
|
file_index: int,
|
|
body: PublicRunFileUpdateRequest,
|
|
) -> dict[str, Any]:
|
|
"""Overwrite an MD/text result file on disk and refresh its size in the run row.
|
|
|
|
Anonymous — no ownership check. See module docstring for the rationale.
|
|
"""
|
|
if len(body.content.encode("utf-8")) > _MAX_PUBLIC_EDIT_BYTES:
|
|
raise HTTPException(status_code=413, detail="Edited content exceeds the 5 MB limit")
|
|
service = get_scheduled_task_service()
|
|
try:
|
|
return await service.admin_update_run_file(run_id, file_index, body.content)
|
|
except FileNotFoundError as exc:
|
|
raise HTTPException(status_code=404, detail=str(exc))
|
|
except PermissionError as exc:
|
|
# Surfaces archive-path escape attempts. Treat as 400 since the caller is
|
|
# already trusted to be reaching a real run row.
|
|
raise HTTPException(status_code=400, detail=str(exc))
|
|
except ValueError as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc))
|
|
|
|
|
|
@router.delete("/runs/{run_id}", status_code=204)
|
|
async def public_delete_run(run_id: str) -> Response:
|
|
"""Remove a single run regardless of owner.
|
|
|
|
Anonymous — no ownership check, mirroring the existing public PUT endpoints.
|
|
Tighten this if the deployment requires read-only public sharing.
|
|
"""
|
|
service = get_scheduled_task_service()
|
|
deleted = await service.admin_delete_run(run_id)
|
|
if not deleted:
|
|
raise HTTPException(status_code=404, detail=f"Run {run_id} not found")
|
|
return Response(status_code=204)
|
|
|
|
|
|
@router.put("/runs/{run_id}/body")
|
|
async def public_update_run_body(run_id: str, body: PublicRunBodyUpdateRequest) -> dict[str, Any]:
|
|
"""Overwrite the inline Markdown body (``result.content``) of a run.
|
|
|
|
Anonymous — no ownership check. Used by the public viewer's "正文" tab.
|
|
"""
|
|
if len(body.content.encode("utf-8")) > _MAX_PUBLIC_EDIT_BYTES:
|
|
raise HTTPException(status_code=413, detail="Edited content exceeds the 5 MB limit")
|
|
service = get_scheduled_task_service()
|
|
try:
|
|
return await service.admin_update_run_body(run_id, body.content)
|
|
except FileNotFoundError as exc:
|
|
raise HTTPException(status_code=404, detail=str(exc))
|