首页
/ OpenViking 后台任务 API 实战指南:任务状态查询、协作取消与列表过滤

OpenViking 后台任务 API 实战指南:任务状态查询、协作取消与列表过滤

2026-09-09 20:24:07作者:田桥桑Industrious

后台任务是 OpenViking 将耗时操作(资源导入、会话提交、重建索引、快照恢复等)与请求处理解耦的核心机制:API 调用方不再同步等待结果,而是立即拿到一个 task_id,随后通过 Task API 轮询任务状态、读取结果或请求取消。本指南以 docs/en/api/17-tasks.md 为骨架,结合 openviking/server/routers/tasks.pyopenviking/service/task_tracker.pyopenviking/service/task_store.py 及 Python / TypeScript / Go 三种 SDK 的源码实现,系统讲解 get_task()cancel_task()list_tasks() 三个接口的参数、用法、响应结构与底层状态机,帮助你在实际项目中正确消费后台任务。

一、后台任务机制概览

Task API 用于跟踪异步资源导入(add_resource)、会话提交(session_commit)、重建索引(admin_reindex)、快照恢复重建索引(snapshot_restore_reindex)等操作。凡是会返回 task_id 的接口,调用方都可以通过本 API 查询其最终结果或错误。

任务生命周期由 6 个状态构成,定义于 task_tracker.pyTaskStatus 枚举:

状态 含义
pending 任务等待执行
running 任务执行中
cancelling 已请求取消,正在等待任务持有的持久化队列消息与进程内工作进行收敛
completed 任务成功完成
failed 任务失败
cancelled 任务已取消

其中 completedfailedcancelled 是终态(_TERMINAL_STATUSES),pendingrunningcancelling 是非终态(_ACTIVE_STATUSES)。任务的每条记录都包含 task_idtask_typestatusresource_idcreated_atupdated_atresulterrorstage 等字段(见 TaskRecord 数据类),并通过 TaskRecord.to_dict() 序列化为 JSON 响应,额外附带 created_at_isoupdated_at_iso 两个 ISO 8601 时间戳字段。

任务记录持久化在 AGFS(Async Git File System)中,由 task_store.pyPersistentTaskStore 负责写入账号作用域的系统任务目录。这意味着服务重启后任务记录仍然可以查询,只要尚未被 TTL 清理机制清除。TaskTracker 的清理策略(task_tracker.py)如下:

  • MAX_TASKS = 10_000:进程内缓存的任务记录上限,超出后按 created_at 最早淘汰;
  • TTL_COMPLETED = 86_400completed / cancelled 任务保留 24 小时;
  • TTL_FAILED = 604_800failed 任务保留 7 天;
  • CLEANUP_INTERVAL = 300:后台清理协程每 5 分钟执行一次过期驱逐。

因此查询任务时若遇到 404 NOT_FOUND,除了"任务不存在/不属于当前用户",也可能是"任务已过期被清理"。

二、get_task():查询单个任务状态

1. 接口实现介绍

get_task() 用于查询返回了 task_id 的 API 对应的后台任务状态,覆盖会话提交(session_commit)、资源导入(add_resource)和管理员重建索引(admin_reindex)等场景。

HTTP 路由实现位于 tasks.py:路由会通过 get_request_context 解析当前请求身份。对于 ROOT 身份,会额外尝试在系统任务账号(SYSTEM_TASK_ACCOUNT_ID / SYSTEM_TASK_USER_ID)下查询;普通用户则严格按 account_id + user_id 归属过滤。查询不到任务时抛出 NOT_FOUND 错误(错误详情中包含 resourcetype: "task" 字段)。

从服务端实现看,TaskTracker.get() 先查进程内缓存(_cached_task),未命中时再回源 AGFS(_load_from_store),因此轮询查询本身开销极小,适合高频调用。

2. 接口参数

参数 类型 必填 默认值 说明
task_id str - 后台 API 返回的任务 ID

3. 使用示例

HTTP API

GET /api/v1/tasks/{task_id}
curl -X GET http://localhost:1933/api/v1/tasks/uuid-xxx \
  -H "X-API-Key: your-key"

Python SDK

Python SDK 的实现位于 sdk/python/openviking_sdk/client.pyget_task() 在收到 404 时返回 None(而非抛异常),便于调用方直接判空处理"任务不存在或已过期":

from openviking_sdk import AsyncHTTPClient

client = AsyncHTTPClient(url="http://localhost:1933", api_key="your-key")
await client.initialize()

task = await client.get_task(task_id="uuid-xxx")
print(f"Status: {task['status']}")
await client.close()

TypeScript SDK

TypeScript SDK 的 getTask() 位于 sdk/typescript/src/client.ts,同样在捕获 NOT_FOUND / 404 时返回 null

console.log(await client.getTask("task-id"));

Go SDK

Go SDK 的 GetTask() 位于 sdk/go/sessions.goNOT_FOUND 时返回 (nil, nil)

task, err := client.GetTask(ctx, "uuid-xxx")
if err != nil {
    return err
}
if task != nil {
    fmt.Println(task["status"])
}

CLI

Rust CLI 的实现位于 crates/ov_cli/src/commands/task.rs,底层调用 HTTP GET /api/v1/tasks/{task_id}

ov task status uuid-xxx

4. 响应示例与字段语义

资源导入进行中(add_resource

{
  "status": "ok",
  "result": {
    "task_id": "uuid-xxx",
    "task_type": "add_resource",
    "status": "running",
    "resource_id": "viking://resources/guide",
    "stage": "processing_queue"
  }
}

字段说明

  • stage 字段可空(nullable)。Git 仓库资源导入任务可能上报 queuedfetchingparsingfinalizingprocessing_queue 等阶段;其他任务类型可能保持 null
  • 实时队列计数故意不包含在任务状态中。如需实时计数请使用 observer 队列相关 API;或者等任务完成后读取 result.queue_status

已完成(session_commit

{
  "status": "ok",
  "result": {
    "task_id": "uuid-xxx",
    "task_type": "session_commit",
    "status": "completed",
    "result": {
      "session_id": "a1b2c3d4",
      "archive_uri": "viking://user/alice/sessions/a1b2c3d4/history/archive_001",
      "memory_diff_uri": "viking://user/alice/sessions/a1b2c3d4/history/archive_001/memory_diff.json",
      "memories_extracted": {
        "profile": 1,
        "preferences": 2,
        "entities": 1,
        "cases": 1
      },
      "active_count_updated": 2,
      "token_usage": {
        "llm": {
          "prompt_tokens": 5200,
          "completion_tokens": 1800,
          "total_tokens": 7000
        },
        "embedding": {
          "total_tokens": 1500
        },
        "total": {
          "total_tokens": 8500
        }
      }
    }
  }
}

memories_extracted 在已完成任务的 result 中报告的是**本次提交(this commit)**的各分类提取数量,需要统计本次提交总数时请对各项求和。这类结构化结果由 TaskTracker.complete() 记录,并在 to_dict() 序列化时通过 _sanitize_task_result 过滤敏感字段(如 user_key)。

三、cancel_task():请求协作取消

1. 接口实现介绍

cancel_task() 请求协作式(cooperative)取消一个后台任务。其语义非常明确:

  • 操作会立即阻止任务创建新的 QueueFS 工作,并取消其正在进行的进程内工作;
  • 已完成的写入不会被回滚(取消不提供事务回滚能力);
  • 如果仍有持久化消息或进程内工作未收敛,接口首先返回 cancelling
  • 只有任务拥有的全部工作收敛后,任务才会变成 cancelled

对处于 cancellingcancelled 状态的任务重复取消是幂等的。

支持取消的任务类型(对应 task_tracker.py_CANCELLABLE_TASK_TYPES 集合):

  • add_resource
  • session_commit
  • admin_reindex
  • snapshot_restore_reindex

(注:源码中该集合还包含 compile,取消支持以服务端实际配置为准。)

代码入口

  • openviking/server/routers/tasks.py:cancel_task() —— HTTP 路由(tasks.py),ROOT 身份直接抛出 PERMISSION_DENIED
  • openviking/service/task_tracker.py:TaskTracker.cancel() —— 任务生命周期核心(task_tracker.py
  • crates/ov_cli/src/commands/task.rs:cancel() —— CLI 命令(task.rs

2. 接口参数

参数 类型 必填 默认值 说明
task_id str - 要取消的后台任务 ID

权限约束:只有拥有该任务的当前用户才能取消任务;ROOT 身份不能取消任务(路由层在 tasks.py 直接拒绝)。

3. 使用示例

HTTP API

POST /api/v1/tasks/{task_id}/cancel
curl -X POST http://localhost:1933/api/v1/tasks/uuid-xxx/cancel \
  -H "X-API-Key: your-key"

Python SDK

task = await client.cancel_task(task_id="uuid-xxx")
print(task["status"])

TypeScript SDK

const task = await client.cancelTask("uuid-xxx");
console.log(task.status);

Go SDK

task, err := client.CancelTask(ctx, "uuid-xxx")
if err != nil {
    return err
}
fmt.Println(task["status"])

CLI

ov task cancel uuid-xxx

4. 响应示例

{
  "status": "ok",
  "result": {
    "task_id": "uuid-xxx",
    "task_type": "add_resource",
    "status": "cancelling",
    "resource_id": "viking://resources/guide",
    "stage": "processing_queue",
    "result": null,
    "error": null
  }
}

如果任务已没有剩余工作,响应中的 status 可能立即为 cancelled;否则应继续用 get_task() 轮询,直到状态变为 cancelled

5. 错误处理

错误码 状态码 触发条件
NOT_FOUND 404 任务不存在、已过期,或属于其他用户
PERMISSION_DENIED 403 ROOT 身份尝试取消任务
FAILED_PRECONDITION 412 任务类型不支持取消,或任务已处于 completed / failed

服务端实现中,路由层将 TaskTracker.cancel() 抛出的 ValueError(如"任务类型不支持取消"、"任务已是终态")包装为 FAILED_PRECONDITION;底层取消流程通过 _work_index.cancel_active(task_id) 中断活跃的进程内工作,再经 _finalize_task_on_owner() 在任务持有的全部工作收敛后落定 cancelled 终态,从而保证"先停止产生新工作,再等待收尾"的协作式取消语义。

四、list_tasks():按条件列出任务

1. 接口实现介绍

list_tasks() 列出当前调用方可见的后台任务,支持按类型、状态、资源过滤。默认只返回用户可见任务;排查 Connector 导入问题时,可传 include_internal=true 以包含其内部 add_resource 子任务。

代码入口

  • openviking/server/routers/tasks.py:list_tasks() —— HTTP 路由(tasks.py
  • sdk/python/openviking_sdk/client.py:AsyncHTTPClient.list_tasks() —— Python SDK(client.py

从路由实现看,普通用户按自身 account_id + user_id 过滤;ROOT 身份则合并查询系统任务(SYSTEM_TASK_* 账号)与进程内缓存任务,并按 created_at 降序去重后截断到 limit。排序与截断在 TaskTracker._list_tasks_on_owner() 中完成(task_tracker.py):先按归属过滤,再过滤 internal 标记、类型、状态、资源,最后按 created_at 倒序返回前 limit 条。

2. 接口参数

参数 类型 必填 默认值 说明
task_type str None 按任务类型过滤,例如 session_commit
status str None 按任务状态过滤:pendingrunningcancellingcompletedfailedcancelled
resource_id str None 按任务资源 ID 过滤,例如一个会话 ID
include_internal bool false 是否包含 Connector 导入创建的内部子任务
limit int 50 返回的最大任务记录数(HTTP 层约束 le=200

默认情况下只返回用户可见任务。诊断 Connector 导入问题时传 include_internal=true,即可看到其内部创建的 add_resource 子任务(源码中通过 meta.get("internal") is True 标记过滤,见 task_tracker.py)。

3. 使用示例

HTTP API

GET /api/v1/tasks?task_type=session_commit&status=running&limit=20
curl -X GET "http://localhost:1933/api/v1/tasks?task_type=session_commit&status=running&limit=20" \
  -H "X-API-Key: your-key"

Python SDK

from openviking_sdk import AsyncHTTPClient

client = AsyncHTTPClient(url="http://localhost:1933", api_key="your-key")
await client.initialize()

tasks = await client.list_tasks(
    task_type="session_commit",
    status="running",
    limit=20,
)
for task in tasks:
    print(task["task_id"], task["status"])
await client.close()

TypeScript SDK

console.log(await client.listTasks());

listTasks() 支持 TaskListOptions,可传 taskTypestatusresourceIdlimit,见 sdk/typescript/src/client.ts。)

Go SDK

tasks, err := client.ListTasks(ctx, &openviking.ListTasksOptions{
    TaskType: "session_commit",
    Status:   "running",
    Limit:    20,
})
if err != nil {
    return err
}
for _, task := range tasks {
    fmt.Println(task)
}

CLI

# 列出任务
ov task list

# 按任务类型和状态过滤
ov task list --task-type session_commit --status running

4. 响应示例

{
  "status": "ok",
  "result": [
    {
      "task_id": "uuid-xxx",
      "task_type": "session_commit",
      "status": "running",
      "resource_id": "a1b2c3d4",
      "created_at": 1770000000.0,
      "updated_at": 1770000005.0,
      "result": null,
      "error": null,
      "stage": null
    }
  ]
}

五、源码级原理:TaskTracker 与持久化

理解 Task API 的可靠性,需要认识其背后的两层设计:

1. 进程内缓存 + AGFS 持久化的双写模型

TaskTracker 维护一个进程本地缓存(self._tasks)作为只读快照,所有变更先经 OwnerLoopDispatcher 派发到单一 owner 事件循环,再通过 KeyedAsyncLockPool 按 task_id 加锁串行化,最后经 StoreIOLimiter(默认最大并发存储 IO 为 8)写入 PersistentTaskStore。这样既保证了同一任务状态变更的原子性,也避免了高频轮询直读 AGFS 的性能开销。

2. 取消与终态落定以"工作索引"为准

任务在途工作时,QueueFS 的持久化消息与进程内活跃任务都会被登记到 TaskWorkIndex_finalize_task_on_owner()task_tracker.py)只有在 _work_index.has_work(task_id) 返回 false(即所有工作收敛)时才会把 cancelling 落定为 cancelled、把带 error 的记录落定为 failed、把带 result 的记录落定为 completed。这解释了为什么取消请求返回后通常需要轮询一段时间才能看到终态。

3. 安全与隐私

  • 错误消息在持久化前会经过 _sanitize_error() 脱敏,匹配 sk-cr_ghp_ntn_xox*Bearer 等敏感模式并替换为 [REDACTED],且截断到 500 字符(task_tracker.py);
  • 任务快照序列化时会丢弃 authaccount_iduser_id 字段,并从 result / meta 中过滤 user_key 等敏感键(task_tracker.py);
  • 所有查询与取消都严格按请求身份的归属过滤,跨用户访问一律返回 NOT_FOUND

六、典型使用流程

结合上述三个接口,一个标准的后台任务消费流程是:

  1. 调用会产生后台任务的 API(如会话提交时传 wait=falseadd_resource 或管理端重建索引),获得 task_id
  2. get_task()(HTTP / Python / TypeScript / Go SDK 或 ov task status)轮询状态;
  3. 若任务长时间未完成且不再需要,调用 cancel_task() 请求协作取消,然后继续轮询直到状态变为 cancelled(若响应已直接为 cancelled 则无需再轮询);
  4. 任务进入终态后,从 result 中读取结构化产出(如 session_commitarchive_urimemory_diff_urimemories_extractedtoken_usage),从 error 中读取脱敏后的失败原因;
  5. 需要排查 Connector 导入或批量排查任务时,用 list_tasks()task_type / status / resource_id 过滤,必要时传 include_internal=true 查看内部子任务。

相关文档

  • Sessions —— 会话提交任务(session_commit
  • Resources —— 资源导入任务(add_resource
  • Content —— 内容异步重建索引任务
登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395