任务管理
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| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
identifier | str | 必填 | 元任务的 identifier 字段 |
name | str | None | 自动生成 | 任务实例显示名称 |
params | dict | None | {} | 传给任务的参数 |
wait | bool | False | True 时阻塞等待完成并返回任务 dict |
timeout | float | 300 | 等待超时秒数 |
poll_interval | float | 2 | 轮询间隔(秒) |
on_complete | callable | None | None | 任务完成后的后台回调,签名 (task_dict) -> None |
on_error | callable | None | None | 任务失败的回调,签名 (exception) -> None |
on_feedback | callable | None | None | 任务进入 pending_approval 时的回调 |
show_in_kanban | bool | None | 后台服务默认 False,普通脚本默认 True | 是否在看板显示 |
workspace_slug | str | None | None | 指定工作区 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_id | str | 必填 | create_task() 返回的 UUID |
timeout | float | 300 | 超时秒数,超时抛 RuntimeError |
poll_interval | float | 2 | 轮询间隔(秒) |
on_feedback | callable | None | None | pending_approval 时的回调;不提供则抛 RuntimeError |
返回值: 任务 dict,包含 status、result 等字段。任务失败时抛 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) -> Nonepython
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 前调用