把定时 Agent 做成可恢复任务系统:状态机、队列与模型调用边界

📅 2026/8/3 13:28:16
把定时 Agent 做成可恢复任务系统:状态机、队列与模型调用边界
很多团队已经不满足于“打开聊天窗口让模型回答一次”。真实研发场景里任务往往需要跨越较长时间凌晨扫描仓库依赖、每天汇总告警、定时检查工单状态、失败后隔一段时间重试或者在外部条件满足后继续执行下一步。所谓“任务自动醒来继续工作”本质不是让模型拥有魔法般的记忆而是把 Agent 放进一个可靠的任务系统里任务状态可保存、执行过程可恢复、外部调用可追踪、失败结果可解释。如果只用一个 cron 脚本直接调用大模型很快会遇到几个问题进程崩溃后不知道执行到哪一步同一任务可能被重复触发模型输出不稳定导致后续动作误执行API 超时后不清楚该重试还是终止多人协作时缺少审计记录。因此定时 Agent 的关键不是“定时”二字而是围绕任务生命周期建立工程边界。本文选择从任务队列、状态机和模型调用封装三个层面实现一个最小可落地版本。示例使用 Node.js、PostgreSQL 和普通 HTTP 接口便于迁移到现有后端。模型服务可以是自建网关、云厂商接口或兼容 OpenAI 风格的中转接口如果团队需要统一管理多模型入口也可以根据官方文档评估 HaerAPIhttps://www.haerapi.com这类服务。核心原理一个可恢复 Agent 任务系统至少包含四个对象任务、运行记录、步骤状态和外部调用记录。任务定义回答“要做什么”例如每天 9 点扫描仓库 issue 并生成摘要。运行记录回答“这一次执行到哪里”例如 2026-08-02 09:00 的任务已经完成了拉取 issue正在等待模型总结。步骤状态回答“当前阶段是否成功”例如FETCH_INPUT、CALL_MODEL、WRITE_RESULT。外部调用记录回答“对外部系统做了什么”例如请求了哪个接口、输入摘要哈希是什么、响应是否通过结构化校验。状态机通常比自由文本日志更可靠。一个简单任务可以定义如下流转PENDING - RUNNING - WAITING_RETRY - RUNNING - SUCCEEDED | | v v FAILED FAILED其中PENDING表示等待执行RUNNING表示已有 worker 获得执行权WAITING_RETRY表示失败但仍可重试SUCCEEDED和FAILED是终态。为了避免多个 worker 同时执行同一条任务需要使用租约字段例如locked_until。worker 只领取未加锁或锁已过期的任务领取时在数据库事务中更新锁这比在内存里维护状态更容易恢复。Agent 的“智能”部分应被限制在明确边界内。模型可以生成摘要、分类、建议动作或结构化 JSON但不应该绕过状态机直接修改生产系统。实践中建议采用三层防线第一提示词要求输出固定 JSON第二服务端用 schema 校验字段第三只有通过业务规则校验的结果才能进入写入步骤。数据表设计下面是最小表结构保留了任务状态、重试次数、租约时间和审计信息。生产环境可以继续增加租户、权限、输入快照、输出版本等字段。CREATETABLEagent_jobs(id BIGSERIALPRIMARYKEY,job_typeTEXTNOTNULL,statusTEXTNOTNULLCHECK(statusIN(PENDING,RUNNING,WAITING_RETRY,SUCCEEDED,FAILED)),payload JSONBNOTNULLDEFAULT{},result JSONB,attempt_countINTNOTNULLDEFAULT0,max_attemptsINTNOTNULLDEFAULT3,run_after TIMESTAMPTZNOTNULLDEFAULTnow(),locked_until TIMESTAMPTZ,last_errorTEXT,created_at TIMESTAMPTZNOTNULLDEFAULTnow(),updated_at TIMESTAMPTZNOTNULLDEFAULTnow());CREATEINDEXidx_agent_jobs_pickupONagent_jobs(status,run_after,locked_until);任务入队时不要直接运行模型只插入一条PENDING记录。定时触发器、Webhook、人工操作都可以复用这条入口。INSERTINTOagent_jobs(job_type,status,payload,run_after)VALUES(daily_issue_digest,PENDING,{repo:example/service-api,since:2026-08-01}::jsonb,now());可执行步骤第一步准备环境变量。密钥必须由环境注入不能写入代码或配置仓库。exportDATABASE_URLpostgres://user:passwordlocalhost:5432/appexportMODEL_BASE_URLhttps://api.example.com/v1exportMODEL_API_KEYreplace-by-secret-managerexportMODEL_NAMEyour-model-name第二步安装依赖。npminit-ynpminstallpg zod node-fetch第三步实现任务领取逻辑。这里的关键是FOR UPDATE SKIP LOCKED它允许多个 worker 并发抢任务但同一行只会被一个事务锁定。importpgfrompg;constpoolnewpg.Pool({connectionString:process.env.DATABASE_URL});exportasyncfunctionpickupJob(){constclientawaitpool.connect();try{awaitclient.query(BEGIN);const{rows}awaitclient.query(SELECT * FROM agent_jobs WHERE status IN (PENDING, WAITING_RETRY) AND run_after now() AND (locked_until IS NULL OR locked_until now()) ORDER BY run_after ASC, id ASC LIMIT 1 FOR UPDATE SKIP LOCKED);if(rows.length0){awaitclient.query(COMMIT);returnnull;}constjobrows[0];constupdatedawaitclient.query(UPDATE agent_jobs SET status RUNNING, locked_until now() interval 5 minutes, attempt_count attempt_count 1, updated_at now() WHERE id $1 RETURNING *,[job.id]);awaitclient.query(COMMIT);returnupdated.rows[0];}catch(err){awaitclient.query(ROLLBACK);throwerr;}finally{client.release();}}第四步封装模型调用并强制结构化输出。示例假设接口兼容常见的 chat completions 形式具体字段以实际服务文档为准。importfetchfromnode-fetch;import{z}fromzod;constDigestSchemaz.object({title:z.string().min(1).max(80),summary:z.string().min(1).max(1000),risk_level:z.enum([low,medium,high]),actions:z.array(z.string().min(1)).max(5)});exportasyncfunctioncallModelForDigest(input){constresawaitfetch(${process.env.MODEL_BASE_URL}/chat/completions,{method:POST,headers:{Authorization:Bearer${process.env.MODEL_API_KEY},Content-Type:application/json},body:JSON.stringify({model:process.env.MODEL_NAME,messages:[{role:system,content:你是研发助理。只输出 JSON不要输出 Markdown。},{role:user,content:请总结以下 issue 列表并给出风险等级和动作${JSON.stringify(input)}}]})});if(!res.ok){thrownewError(model_http_${res.status});}constdataawaitres.json();consttextdata.choices?.[0]?.message?.content;if(!text){thrownewError(model_empty_content);}letparsed;try{parsedJSON.parse(text);}catch{thrownewError(model_invalid_json);}returnDigestSchema.parse(parsed);}第五步实现 worker 主循环。失败时不要简单退出而是根据次数进入延迟重试或终态失败。import{pickupJob}from./pickup.js;import{callModelForDigest}from./model.js;importpgfrompg;constpoolnewpg.Pool({connectionString:process.env.DATABASE_URL});asyncfunctionmarkSucceeded(id,result){awaitpool.query(UPDATE agent_jobs SET status SUCCEEDED, result $2, locked_until NULL, updated_at now() WHERE id $1,[id,result]);}asyncfunctionmarkFailedOrRetry(job,err){constcanRetryjob.attempt_countjob.max_attempts;awaitpool.query(UPDATE agent_jobs SET status $2, last_error $3, locked_until NULL, run_after CASE WHEN $2 WAITING_RETRY THEN now() (($4 * 30) || seconds)::interval ELSE run_after END, updated_at now() WHERE id $1,[job.id,canRetry?WAITING_RETRY:FAILED,String(err.message||err),Math.max(1,job.attempt_count)]);}asyncfunctionrunOnce(){constjobawaitpickupJob();if(!job)returnfalse;try{if(job.job_type!daily_issue_digest){thrownewError(unsupported_job_type:${job.job_type});}constissueInputawaitloadIssues(job.payload);constdigestawaitcallModelForDigest(issueInput);awaitwriteDigest(job.payload.repo,digest);awaitmarkSucceeded(job.id,digest);}catch(err){awaitmarkFailedOrRetry(job,err);}returntrue;}asyncfunctionmain(){while(true){constworkedawaitrunOnce();if(!worked)awaitnewPromise(rsetTimeout(r,3000));}}asyncfunctionloadIssues(payload){return[{id:1,title:检查${payload.repo}的待处理事项}];}asyncfunctionwriteDigest(repo,digest){console.log(digest ready,repo,digest.title);}main().catch(err{console.error(err);process.exit(1);});这个示例没有编造任何性能指标也不承诺某个模型一定能产出稳定 JSON。实际项目中结构化输出质量取决于模型能力、提示词、输入复杂度和服务端校验策略。对于关键动作例如发版、删除数据、付款、修改权限建议加入人工审批或二次确认。部署与运行开发环境可以直接运行 workernodeworker.js生产环境建议把 worker 作为长期进程部署并配合进程管理器或容器编排平台。下面是一个 systemd 示例只展示关键配置[Unit] DescriptionAgent Job Worker Afternetwork.target [Service] WorkingDirectory/opt/agent-worker ExecStart/usr/bin/node worker.js Restartalways RestartSec5 EnvironmentNODE_ENVproduction EnvironmentFile/etc/agent-worker.env [Install] WantedBymulti-user.target/etc/agent-worker.env应由运维或密钥系统生成并限制文件权限DATABASE_URLpostgres://user:passworddb:5432/appMODEL_BASE_URLhttps://api.example.com/v1MODEL_API_KEYfrom-secret-managerMODEL_NAMEyour-model-name定时触发可以使用系统 cron、应用内调度器或云平台计划任务。无论哪种方式触发器都只负责插入任务不直接执行业务逻辑。这样即使触发器重复运行也可以通过业务唯一键避免重复任务。例如每天每个仓库只允许一条摘要任务ALTERTABLEagent_jobsADDCOLUMNdedupe_keyTEXT;CREATEUNIQUEINDEXuniq_agent_jobs_dedupeONagent_jobs(dedupe_key)WHEREstatusIN(PENDING,RUNNING,WAITING_RETRY,SUCCEEDED);常见问题1. 为什么不用纯 cron 加脚本cron 适合触发不适合表达复杂生命周期。只要任务可能失败重试、跨步骤恢复、多人审计或并发执行就应该把状态放到数据库或队列系统中而不是只依赖进程日志。2. 模型返回了非 JSON 怎么办不要在业务代码里猜测修复。可以先重试一次也可以把任务转入WAITING_RETRY。如果多次失败应记录原始错误摘要并进入FAILED由人工查看提示词、输入长度或模型能力是否匹配。3. worker 崩溃会不会导致任务永远卡在 RUNNING如果使用locked_until锁过期后任务可以重新被领取。但要注意下游写入必须幂等。例如写摘要时使用任务 id 或业务日期作为唯一键避免崩溃前已经写入、重试后再次写入。4. 如何处理模型供应商切换业务代码不应散落供应商字段。建议把模型调用封装成callModelForDigest这类函数并把MODEL_BASE_URL、MODEL_NAME、鉴权方式放入配置。切换到任何兼容接口或中转服务前都要用少量真实脱敏样例验证响应结构、错误码、超时和速率限制行为。5. Agent 能不能自动执行修复动作可以但要分级。低风险动作如生成草稿、添加标签、创建待办可以在 schema 校验后自动执行中风险动作如改配置、提交 PR建议保留审查高风险动作如删除数据、调整权限、触发财务流程应默认要求人工批准。总结定时 Agent 的工程重点不是让模型“记住上次聊到哪”而是让系统可靠地保存任务状态、恢复执行进度、约束模型输出并审计外部动作。最小可行架构可以从 PostgreSQL 状态机、租约式任务领取、结构化模型输出和幂等写入开始。等任务量增加后再引入专门队列、分布式追踪、人工审批台和更细的权限模型。只要把 Agent 当作普通后端任务系统中的一个可替换能力而不是不可控的黑箱就能在保留自动化收益的同时降低误执行和不可恢复失败的风险。本文包含 HaerAPI 的推广信息是否采用应根据其当前文档、数据处理条款、可用模型和自身合规要求独立判断。