AI终端后台执行架构:任务状态机与异步调度实践
发布时间:2026/9/14 20:16:36来源:尧图网络
在命令行里跑 AI 编码助手最理想的状态是什么你敲下一条指令界面立刻有反馈任务在背后慢慢跑进度随时能看随时能用 CtrlC 断掉。可大多数初版实现做出来的效果是指令一敲整个终端像被掐住脖子等它跑完界面才开始刷新。我在这套“从零构建 AI Code 终端系统”写到第 9 期遇到的最关键架构转折点就是把“启动”和“等结果”拆开。这篇文章不聊 Agent 提示词也不聊模型选型只讲后台执行这条主链路任务怎么提交、状态怎么跟踪、结果怎么回收、进程怎么清理。顺带拆一下 Continue 和 OpenCode 这类开源项目为什么也走这条路。如果你正在写终端里的 AI 编码工具或者只是想搞清楚异步任务调度该怎么设计这篇应该对你有用。1. 后台执行真正要解决的事UI 不能陪着任务一起阻塞1.1 最初把启动和等待写在一起翻车是迟早的事我做这个终端系统的早期版本所有耗时命令的调用方式都长这样async def run_and_wait(cmd): proc await asyncio.create_subprocess_shell( cmd, stdoutasyncio.subprocess.PIPE, stderrasyncio.subprocess.PIPE, ) stdout, stderr await proc.communicate() return stdout, stderr, proc.returncode看着没什么问题一个函数把“启动进程”和“等结果”全包了调用方省事。但真正的使用场景很快就让它露馅AI Agent 在修改完代码之后需要跑测试验证。测试一跑就是几十秒这几十秒里await proc.communicate()会把当前协程牢牢挂住。我的 TUI 渲染循环和键盘事件循环跑在同一个 asyncio loop 里主协程被挂住整个界面就再也不刷新了输入框失去响应连终端尺寸变化都处理不了。更难受的是这种写法天生做不了并行。后来我加了一个“同时跑 lint 和单元测试”的功能如果用run_and_wait串行等第二个任务的等待全部浪费在第一个任务上。如果改用asyncio.gather并发又要维护一堆零散的协程句柄一个任务出异常整批任务的排查体验都很差。1.2 “启动”和“等结果”明明是两次不同的操作问题根源不是 asyncio 不行是我把两件不同生命周期的事情焊死在了一个函数里。启动这个动作应该快速返回。你把任务丢给后台拿到一个任务 ID然后就可以继续处理 UI 事件、接受用户输入、展示“任务已提交”的反馈。等结果是另一回事你可以选择阻塞等待也可以选择轮询状态还可以选择订阅事件流这些方式都是围绕任务 ID 展开的。拆开之后我得到的好处非常明确提交任务后能立刻给用户反馈而不是让终端看起来像死机。多个后台任务可以并行跑互不阻塞。每个任务有独立状态和结果存储随时可查。取消、超时、重试这些操作有了明确的挂载点。所以我在这一期做的第一件事不是写 UI而是先把任务状态机定义清楚。2. 任务状态机先给执行现场一个身份才能谈启动和等待2.1 状态定义与流转规则后台执行的本质是“任务生命周期管理”。如果连状态都没有定义清楚后面无论是 UI 展示还是 Agent 内部逻辑都会陷入一团乱麻。我的任务状态机长这样状态含义进入时机PENDING已提交排队等待 worker 调度start()返回时RUNNINGworker 取出任务开始执行出队的那一刻COMPLETED正常执行完成执行函数正常返回FAILED执行过程抛异常捕获到非取消类异常CANCELLED用户主动取消或超时取消任务被取消流转规则也很简单PENDING 只能进 RUNNINGRUNNING 之后只能进 COMPLETED、FAILED 或 CANCELLED。没有反向流转也不允许从 COMPLETED 再跳回 RUNNING。这套状态机看起来简单但它是整个后台执行模块的地基。定义好状态之后接口设计就清晰了start()只负责把状态变成 PENDINGworker 出队时把状态改成 RUNNING执行完按结果改成终态。2.2 TaskRecord 记录什么状态只是枚举真正要能在任意时刻把任务的执行现场捞出来需要一个结构体。我用的 TaskRecord 长这样from dataclasses import dataclass, field from enum import Enum, auto import time class TaskStatus(Enum): PENDING auto() RUNNING auto() COMPLETED auto() FAILED auto() CANCELLED auto() dataclass class TaskRecord: task_id: str name: str status: TaskStatus TaskStatus.PENDING created_at: float field(default_factorytime.time) started_at: float | None None finished_at: float | None None result: object None error: str | None None property def elapsed(self) - float: end self.finished_at or time.time() begin self.started_at or self.created_at return round(end - begin, 3)为什么要把这些字段全部放在一个外部结构里而不是塞在协程的局部变量中原因很简单当你把任务丢到后台之后调用start()的那个协程早就返回了执行现场不在任何栈上。所有后续的 UI 渲染、日志记录、Agent 决策都需要靠 task_id 去一个统一的地方把记录捞出来。2.3 状态对象就是前后端的契约在我的终端系统里TaskRecord 不只是内部数据结构它直接就是 UI 层和 Agent 层之间的契约。TUI 的渲染循环每帧会读一次所有活跃任务的 TaskRecord根据 status 决定该显示 spinner、进度条还是结果面板。Agent 内部某个步骤需要等待前一个步骤完成时它await manager.wait(task_id)拿到的也是一个 TaskRecord而不是散落的返回值。这个契约让后端实现随便换UI 不用改。今天用进程内字典存记录明天换成 SQLite只要 TaskRecord 结构不变展示层完全无感。3. 最小可用的 TaskManager把两段式模型落实成代码3.1 调度器骨架定义好状态机之后核心调度器就可以写了。我在这一期实现的管理器核心就四个东西任务字典、Future 字典、任务队列、一组 worker。import asyncio import uuid from collections.abc import Awaitable, Callable class TaskManager: def __init__(self, max_workers: int 4): self._max_workers max_workers self._records: dict[str, TaskRecord] {} self._futures: dict[str, asyncio.Future] {} self._running: dict[str, asyncio.Task] {} self._queue: asyncio.Queue[tuple[str, Callable[[], Awaitable]]] asyncio.Queue() self._workers: list[asyncio.Task] [] async def start_workers(self): for _ in range(self._max_workers): self._workers.append(asyncio.create_task(self._worker())) async def start(self, name: str, runner: Callable[[], Awaitable]) - str: task_id uuid.uuid4().hex record TaskRecord(task_idtask_id, namename) loop asyncio.get_running_loop() done loop.create_future() self._records[task_id] record self._futures[task_id] done await self._queue.put((task_id, runner)) return task_id async def wait(self, task_id: str, timeout: float | None None) - TaskRecord: future self._futures.get(task_id) if future is None: raise KeyError(ftask {task_id} not found) try: await asyncio.wait_for(future, timeout) except asyncio.TimeoutError: pass return self._records[task_id] def get(self, task_id: str) - TaskRecord: return self._records[task_id] async def cancel(self, task_id: str): record self._records.get(task_id) if record is None or record.status in ( TaskStatus.COMPLETED, TaskStatus.FAILED, TaskStatus.CANCELLED, ): return if record.status TaskStatus.PENDING: record.status TaskStatus.CANCELLED record.finished_at time.time() return task self._running.get(task_id) if task: task.cancel()这里有一个关键设计_futures[task_id]里的 Future 只承担“完成通知”这个职责不承担搬运结果和异常的职责。结果放在 TaskRecord 的result字段里异常放在error字段里。这样能避免一个很隐蔽的坑——asyncio 的 Future 如果设置了异常又没人去future.result()取GC 的时候会打出一条 “Future exception was never retrieved” 的警告。对于长跑不退出的事件循环这种警告会越积越多而且很难排查。3.2 Worker 消费模型与取消的传播worker 的实现是理解整套模型的关键。它从队列里取任务然后立刻把任务包装成一个独立的 asyncio.Task 来执行而不是直接await runner()。这一步的目的是把“可取消”的能力保留下来。async def _worker(self): while True: task_id, runner await self._queue.get() record self._records[task_id] # 处理在队列里就被取消的任务 if record.status TaskStatus.CANCELLED: record.finished_at time.time() self._queue.task_done() continue record.status TaskStatus.RUNNING record.started_at time.time() running_task asyncio.create_task(runner()) self._running[task_id] running_task done self._futures[task_id] try: result await running_task record.result result record.status TaskStatus.COMPLETED except asyncio.CancelledError: record.status TaskStatus.CANCELLED except Exception as exc: record.status TaskStatus.FAILED record.error f{type(exc).__name__}: {exc} finally: record.finished_at time.time() self._running.pop(task_id, None) if not done.done(): done.set_result(None) self._queue.task_done()为什么不能直接await runner()因为直接 await 时如果调用task.cancel()取消的是 worker 本身而不是具体某个任务。worker 一旦被取消整条消费队列就断了后面的任务全都没人处理。所以 worker 必须常驻循环取到任务后再包一层create_task(runner())并把 running_task 存到self._running里。取消单个任务时直接_running[task_id].cancel()worker 本身不受影响。3.3 超时和取消要注意的先后顺序使用wait(task_id, timeout120)有一个常见误区超时只代表你的等待被打断了并不代表后台任务真的停了。asyncio 的wait_for超时后取消的是 Future 对象而 Future 背后的 running_task 还在继续跑。如果你想“超时即取消”需要显式补一个取消动作try: record await manager.wait(task_id, timeout120) except asyncio.TimeoutError: await manager.cancel(task_id) record await manager.wait(task_id)这个顺序我踩过不少次。系统的第一个版本里超时之后只提示了一句“timeout”结果后台的测试进程还在跑占着文件锁下一次测试直接报错。后来我统一改成超时后主动调用cancel()并把取消链路一路传到子进程那一层问题才消失。4. 结果回收的三种姿势Future、事件回调、轮询后台任务执行完之后结果怎么送回来是整个系统体验的分水岭。我实际用下来有三种姿势都值得掌握各自有明确的适用场景。4.1 Future 等待最符合直觉的主流程等待在 Agent 内部逻辑里经常需要严格等待一个任务完成再继续下一步。比如“跑完测试再决定要不要改代码”这种情况下最顺手的写法就是await manager.wait(task_id)。这种方式非常像同步编程里的函数调用。调用方不关心任务怎么被调度、怎么被并发它只想知道“那件事办完没有”。TaskManager 内部的 Future 只负责“办完了”这个信号调用方拿到 TaskRecord 之后自己决定是看result还是看status。需要提醒的是在 asyncio 的单线程模型里await manager.wait(task_id)不会真的阻塞操作系统线程。如果有人把 UI 渲染和 Agent 的逻辑写在一个协程里然后又去等待一个很重的任务那整个 UI 还是会卡。正确做法是Agent 的执行逻辑本身就是后台任务跑在 worker 里UI 永远只消费 TaskRecord 和事件不直接等待执行结果。4.2 事件回调流式输出的命脉TaskRecord 适合拿来做“最终结果”但 AI Code 终端系统里大量场景需要“中间过程”。比如 Agent 正在跑一个测试你希望终端里像本地命令一样实时滚出日志。这种流式场景Future 显然不合适。Future 只有完成和未完成两种状态没有“跑到一半”的通道。我给 TaskManager 配了一个简单的 EventBusfrom collections import defaultdict from collections.abc import Callable import time class EventBus: def __init__(self): self._listeners: dict[str, list[Callable]] defaultdict(list) def subscribe(self, event_type: str, callback: Callable): self._listeners[event_type].append(callback) async def publish(self, event_type: str, payload: dict): for callback in self._listeners.get(event_type, []): callback(payload)事件类型我会按任务维度打散比如task.stdout、task.stderr、task.statuspayload 里带上 task_id 和内容行。终端 UI 订阅这些事件以后把输出行追加到对应的任务面板里。这里有一个性能上的细节终端渲染如果每收到一行就刷新一次输出量一大就会疯狂闪烁。我在 UI 层对输出做了合并和节流事件回调只负责把数据塞到渲染缓冲里真正的重绘动作按帧合并执行。4.3 轮询简单但别随便用进程内通信我基本不用轮询原因很直接没有跨进程边界事件和 Future 都够用轮询除了浪费 CPU 没有额外价值。轮询真正发挥作用的场景是跨进程或跨节点。比如 TaskRecord 落到了 SQLite 里另一个守护进程需要通过 HTTP 查询任务状态再比如分布式场景下多个执行机共享一张任务表。这时候没有注册回调的通道只能定时去查询任务表。如果确实要用轮询我建议把频率压到 250ms 到 500ms 一档。对终端 UI 来说超过 500ms 的刷新间隔能感受到明显延迟250ms 以下绝大多数场景根本感知不到提升。4.4 三种方式的横向选择方式核心价值适用场景主要缺点Future 等待按任务聚合最终结果Agent 步骤内部依赖拿不到中间输出事件回调实时、流式、低延迟终端日志展示、进度更新需要自己做背压和节流轮询跨进程、无状态任务存数据库后的外部查询有延迟且浪费资源实际项目里我的默认组合是“Future 等待 事件回调”。先通过事件把过程实时刷给 UI再通过 Future/TaskRecord 让 Agent 逻辑可以优雅地拿到结果。轮询只有在跨进程场景才启用。5. 换成真正命令行执行进程级后台执行还有这些细节5.1 create_subprocess_exec 天然就是两段式如果你后台跑的不是一个 Python 函数而是一条真正的 shell 命令那“启动和等结果拆开”这件事就得更认真对待。因为 asyncio 的进程创建接口本身就是两段式的proc await asyncio.create_subprocess_exec( *cmd, stdoutasyncio.subprocess.PIPE, stderrasyncio.subprocess.STDOUT, )这一行调用做完子进程已经启动了但还没等它退出。proc是操作系统层面的进程句柄。之后你可以选择await proc.wait()只等退出码。通过async for逐行读它的输出。先不等待让它在后台继续跑自己去做别的事。这个拆分的天然性和我上面说的 TaskManager 模型完全同构start返回 task_id等价于create_subprocess_exec返回 proc 句柄wait等待 future等价于proc.wait()。5.2 读输出与等退出码必须并发很多人在子进程这里重新踩了坑就是重新把输出读取和退出码等待写成了串行# 错误示例 return_code await proc.wait() stdout await proc.stdout.read()如果子进程输出量很大这个写法会死锁。原理是管道缓冲区是有限的子进程写满缓冲区之后写不进去就会阻塞等待父进程来读而父进程还卡在proc.wait()上谁都没法先走一步。正确姿势是把输出读取和退出码等待做成并发的两个协程最后再一起收口async def _consume(stream: asyncio.StreamReader, on_line: Callable[[str], None]): async for raw in stream: line raw.decode(errorsreplace).rstrip() on_line(line) async def run_process(cmd: list[str], on_line: Callable[[str], None]) - int: proc await asyncio.create_subprocess_exec( *cmd, stdoutasyncio.subprocess.PIPE, stderrasyncio.subprocess.STDOUT, ) reader_task asyncio.create_task(_consume(proc.stdout, on_line)) return_code await proc.wait() await reader_task return return_code如果希望 stdout 和 stderr 分开展示那就要开两个读取任务分别消费两个管道。这里一个关键判断是是否把 stderr 合并进 stdout。对于终端 UI 展示而言合并往往更符合用户对“命令输出”的直觉只有当你需要分别标记错误流时才考虑拆开。5.3 取消子进程时记得把清理链路走到位进程级任务的取消比普通函数任务复杂一截。普通协程取消后Python 的异步栈会自动清理。可子进程是操作系统资源你取消协程不会自动把子进程杀掉。如果处理不好终端里会不断累积僵尸进程。CancelledError 处理函数必须主动 terminateasync def _runner(): proc await asyncio.create_subprocess_exec( *cmd, stdoutasyncio.subprocess.PIPE, stderrasyncio.subprocess.STDOUT, ) try: return_code await proc.wait() return return_code except asyncio.CancelledError: proc.terminate() try: await asyncio.wait_for(proc.wait(), timeout3) except asyncio.TimeoutError: proc.kill() await proc.wait() raise这里有个细节值得展开先terminate()发送 SIGTERM 给进程给它一点时间做清理3 秒后还没退再kill()强杀。直接 kill 有概率让子进程来不及保存状态间接造成文件写一半或临时文件残留。5.4 并发限制要限制到“共享资源”这一层TaskManager 的 max_workers 控制的是任务执行协程的数量但真实世界里并发瓶颈往往不在任务引擎而在共享资源。比如两个测试任务同时跑进同一个 Git 仓库目录互相之间可能因为索引锁冲突直接失败。我给不同类型的任务配了不同的信号量。跑 Agent 自主修改代码的任务同时最多只允许一个跑只读检查的任务lint、type check可以放宽到两个跑 AI 模型请求的任务单独按 API 限流来控制。TaskManager 只保证调度层面不阻塞具体任务之间的资源冲突要由任务类型自己约束。6. 从 OpenCode 和 Continue 回看这个设计6.1 开源 AI Code Agent 的套路也长这样“command code ai 对比 opencode”这个话题最近很热很多人关心这些工具谁比谁强。但落到工程架构上看无论上层界面怎么花哨底层面对的问题都一模一样一个 Agent 执行步骤可能耗时几十秒甚至几分钟你不能让整个终端终端陪着它等。Continue 这类开源 AI 编码助手核心循环里每一步都可以是一个独立步骤步骤之间有依赖关系。它界面上能逐步展示“读取文件”“修改代码”“执行测试”“分析结果”本质就是每个步骤都作为一个可观测的任务在跑。界面只负责订阅步骤状态变化事件不会像个傻瓜一样同步等待整条链路完成后才刷新。OpenCode 这种终端形态的 AI agent给我的感觉是对任务流的拆分更细。整个 Agent 的执行过程会铺设成一条事件流提示、工具调用、输出都以事件形式投递终端组件消费这些事件做增量渲染。这和我前面的 EventBus 设计思路是同一个套路执行方制造事件展示方消费事件两边互不阻塞。6.2 从这个对比里能提炼的三条原则第一UI 和任务执行不能共享同一个“等待栈”。谁负责渲染谁就老老实实消费状态和事件不要亲手去等某个耗时任务结束。第二状态必须外置事件必须单向流动。TaskRecord 里的 status 可以被多个协程并发读取但写状态的入口只能有一个就是 worker 执行链路的内部逻辑。第三超时、取消、错误不能只出现在日志里必须变成可见状态。我在任务面板上长年保留一个“已取消”的展示态它不只是一个 UI 细节还是调试异步系统的眼睛。没有这个状态你根本分不清任务是“还没跑完”还是“永远跑不完”。6.3 我们这套设计和它们差在哪差距主要不在任务调度模型而在持久化和恢复。OpenCode 这类工具会对任务状态做持久化让用户在中途退出之后还能恢复现场。我的 TaskManager 目前还是纯内存结构服务一停任务记录全丢。后续演进方向很明确把 TaskRecord 落成 SQLite 表Future 仍然做内存内的完成通知记录本身支持跨会话还原。两段式接口在这个模型里依然成立start返回任务 IDwait等待完成信号只是状态的存储换了层皮。如果你也在写类似的终端工具我建议你先把启动和等待拆开再往上堆功能。这个改动看起来只是把一段代码截成两段但拆完之后UI 响应、并行能力、取消机制、可观测性这些硬骨头会自然地长出挂载点。我在拆这一步之前已经被卡死的界面坑了太多次拆完之后后面几期的功能才敢放心往上加。
网站建设高端定制企业官网