Skip to content

任务管理

create_task()

在运行中的脚本内创建一个新任务实例。

python
def create_task(
    identifier: str,
    name: str | None = None,
    params: dict | None = None,
    wait: bool = False,
    timeout: float = 300,
    poll_interval: float = 2,
    on_complete=None,
    on_error=None,
    on_feedback=None,
    show_in_kanban: bool | None = None,
    workspace_slug: str | None = None,
) -> str | dict
参数类型默认值说明
identifierstr必填元任务的 identifier 字段
namestr | None自动生成任务实例显示名称
paramsdict | None{}传给任务的参数
waitboolFalseTrue 时阻塞等待完成并返回任务 dict
timeoutfloat300等待超时秒数
poll_intervalfloat2轮询间隔(秒)
on_completecallable | NoneNone任务完成后的后台回调,签名 (task_dict) -> None
on_errorcallable | NoneNone任务失败的回调,签名 (exception) -> None
on_feedbackcallable | NoneNone任务进入 pending_approval 时的回调
show_in_kanbanbool | None后台服务默认 False,普通脚本默认 True是否在看板显示
workspace_slugstr | NoneNone指定工作区 slug,覆盖元任务默认工作区

返回值:

  • wait=False 或使用 on_complete:返回新任务的 UUID 字符串
  • wait=True:返回任务完成后的 dict,包含 "result" 字段

示例:阻塞等待

python
def run(params: dict, reporter) -> None:
    reporter.set_phase("运行子任务")
    result = create_task(
        "data-processor",
        params={"source": params["source"]},
        wait=True,
        timeout=120,
    )
    print(result["result"])  # 子任务的返回值

示例:后台并行

python
from workflow import create_task, flush_tasks

results = []

def run(params: dict, reporter) -> None:
    for url in params["urls"]:
        create_task(
            "fetch-url",
            params={"url": url},
            on_complete=lambda r: results.append(r["result"]),
            on_error=lambda e: print(f"失败: {e}"),
            show_in_kanban=False,
        )
    flush_tasks()  # 等待所有回调完成
    reporter.set_progress(100)
    return results

示例:审批流程

python
def handle_feedback(task: dict) -> None:
    # 用户在 UI 触发审批后,自动调用此回调
    feedback(task["id"], message="已审核,继续执行")

def run(params: dict, reporter) -> None:
    result = create_task(
        "review-step",
        params={"content": params["draft"]},
        wait=True,
        on_feedback=handle_feedback,
    )

wait_task()

等待一个任务实例完成,返回其结果。

python
def wait_task(
    task_id: str,
    timeout: float = 300,
    poll_interval: float = 2,
    on_feedback=None,
) -> dict
参数类型默认值说明
task_idstr必填create_task() 返回的 UUID
timeoutfloat300超时秒数,超时抛 RuntimeError
poll_intervalfloat2轮询间隔(秒)
on_feedbackcallable | NoneNonepending_approval 时的回调;不提供则抛 RuntimeError

返回值: 任务 dict,包含 statusresult 等字段。任务失败时抛 RuntimeError

python
task_id = create_task("summariser", params={"url": "https://example.com"})
# ... 做其他工作 ...
task = wait_task(task_id, timeout=120)
print(task["result"])

feedback()

向处于 pending_approval 状态的子任务提供反馈。

python
def feedback(task_id: str, message: str = "") -> None
  • Agent 任务(有关联 Hermes session):message 作为下一轮用户输入发给 Agent,函数阻塞直到 Agent 完成该轮。
  • 脚本任务(无 session):直接审批,message 忽略。
python
def handle_approval(task: dict) -> None:
    feedback(task["id"], message="内容已确认,请继续输出下一段")

result = create_task("draft-writer", wait=True, on_feedback=handle_approval)

flush_tasks()

等待所有通过 on_complete 启动的后台线程完成。在 run() 函数末尾调用,确保所有回调在容器退出前执行完毕。

python
def flush_tasks(timeout: float = 300) -> None
python
def run(params: dict, reporter) -> None:
    for item in params["items"]:
        create_task("process-item", params={"item": item},
                    on_complete=lambda r: save(r))
    flush_tasks()   # ← 必须在 return 前调用

Built with VitePress