欢迎光临
我们一直在努力

AI Agent 编排的声明式配置:像写 K8s YAML 一样定义 Agent

AI Agent 编排的声明式配置:像写 K8s YAML 一样定义 Agent

把 Agent 的行为硬编码在代码里,改一个步骤就要重新编译部署——声明式配置让 Agent 逻辑和运行引擎彻底分离。

一、场景痛点

你的 Agent 系统有三条业务链路:用户问答、文档生成、数据查询。每条链路的步骤数量、工具调用顺序、错误处理策略都不一样。你用 Python 写了三个 orchestrator 类,每个类 200 行,步骤之间的跳转逻辑嵌套在条件判断里。

两周后业务要求:问答链路在"用户情绪低落"时插入一个安抚步骤。你改了问答 orchestrator 的代码,重新部署,其他两条链路不受影响——但你需要重新测试问答链路的所有分支。一个月后,文档生成链路也要加类似步骤,你又改一次代码,又部署一次。

核心矛盾:Agent 的编排逻辑与执行引擎耦合在一起,每次逻辑变更都牵动代码和部署,无法做到"改配置不改代码"。

二、底层机制与原理剖析

2.1 声明式 vs 命令式编排

2.2 配置模型设计

Agent 编排配置至少包含五个维度:

  • steps:步骤列表,每个步骤定义工具名、输入映射、超时时间
  • conditions:条件跳转,步骤间的前置判断
  • error_policy:错误策略——重试次数、降级方案、终止条件
  • dependencies:步骤依赖关系——哪些步骤必须在哪些步骤之后
  • resources:资源约束——token 上限、并发限制、超时阈值
  • 类比 K8s:steps 对应 Pod 的容器列表,conditions 对应 initContainer 的条件,error_policy 对应 Pod 的 restartPolicy,dependencies 对效 Pod 的依赖顺序。

    2.3 配置验证与运行时执行

    声明式配置必须经过两阶段验证:

    • 静态验证(配置加载时):YAML 结构是否合法、步骤名是否唯一、条件表达式是否可解析、循环依赖检测
    • 动态验证(执行时):工具是否存在、参数类型是否匹配、超时阈值是否可达

    三、生产级代码实现

    3.1 Agent 配置 YAML 定义

    # agent_configs/qa_pipeline.yaml —— 问答链路的声明式配置
    # 格式仿照 K8s YAML:apiVersion、kind、metadata、spec
    apiVersion: agent.orb/v1
    kind: AgentPipeline
    metadata:
    name: qa-pipeline
    description: "用户问答链路:意图识别→知识检索→回答生成→情绪安抚"
    labels:
    team: support
    tier: l2

    spec:
    # 全局资源约束
    resources:
    maxTokensPerStep: 4000
    totalTokenLimit: 16000
    maxConcurrentSteps: 3
    globalTimeoutMs: 60000

    # 步骤定义:顺序执行,条件跳转控制分支
    steps:
    – name: intent-parse
    tool: nlu-parser
    version: v2.3
    input:
    # 输入映射:从上下文中提取字段,映射到工具参数
    text: "{{ user_input }}"
    language: "{{ metadata.language | default('zh') }}"
    timeoutMs: 2000
    retryPolicy:
    maxRetries: 1
    backoffMs: 500
    output:
    # 输出绑定:工具返回值映射到上下文变量
    intent: "{{ result.intent }}"
    confidence: "{{ result.confidence }}"
    emotion: "{{ result.emotion | default('neutral') }}"

    – name: knowledge-search
    tool: vector-search
    version: v1.5
    # 条件前置:只有置信度 > 0.6 才执行检索
    condition: "{{ steps.intent-parse.confidence > 0.6 }}"
    input:
    query: "{{ steps.intent-parse.intent }}"
    topK: 5
    timeoutMs: 5000
    retryPolicy:
    maxRetries: 2
    backoffMs: 1000
    output:
    documents: "{{ result.hits }}"

    – name: emotion-check
    # 安抚步骤:用户情绪低落时插入
    tool: emotion-response
    version: v1.0
    condition: "{{ steps.intent-parse.emotion == 'negative' }}"
    input:
    emotion: "{{ steps.intent-parse.emotion }}"
    intent: "{{ steps.intent-parse.intent }}"
    timeoutMs: 3000
    output:
    comfortMessage: "{{ result.message }}"

    – name: answer-generate
    tool: llm-generator
    version: v3.1
    # 没有条件则总是执行(依赖前面步骤的输出)
    input:
    intent: "{{ steps.intent-parse.intent }}"
    documents: "{{ steps.knowledge-search.documents | default([]) }}"
    comfortMessage: "{{ steps.emotion-check.comfortMessage | default(null) }}"
    timeoutMs: 10000
    output:
    answer: "{{ result.text }}"

    # 错误策略:全局兜底
    errorPolicy:
    # 全链路超时:超过 60 秒直接终止
    onGlobalTimeout: abort
    # 单步超时:跳过该步骤,用降级方案继续
    onStepTimeout:
    strategy: fallback
    fallbackTool: fallback-response
    fallbackInput:
    intent: "{{ steps.intent-parse.intent | default('unknown') }}"
    # 工具调用失败:重试后仍失败则降级
    onToolFailure:
    strategy: retry_then_fallback
    maxRetries: 2
    fallbackTool: fallback-response

    3.2 配置验证器

    // config-validator.ts —— Agent 配置静态验证
    import Ajv, { ValidateFunction } from 'ajv';
    import addFormats from 'ajv-formats';

    /** 配置验证结果 */
    export interface ValidationResult {
    valid: boolean;
    errors: string[];
    warnings: string[];
    }

    /** Agent 配置的类型定义(从 YAML schema 转换) */
    interface AgentConfig {
    apiVersion: string;
    kind: string;
    metadata: { name: string; description?: string; labels?: Record<string, string> };
    spec: {
    resources: {
    maxTokensPerStep: number;
    totalTokenLimit: number;
    maxConcurrentSteps: number;
    globalTimeoutMs: number;
    };
    steps: StepConfig[];
    errorPolicy: ErrorPolicyConfig;
    };
    }

    interface StepConfig {
    name: string;
    tool: string;
    version: string;
    condition?: string;
    input: Record<string, string>;
    timeoutMs: number;
    retryPolicy?: { maxRetries: number; backoffMs: number };
    output?: Record<string, string>;
    }

    interface ErrorPolicyConfig {
    onGlobalTimeout: string;
    onStepTimeout?: { strategy: string; fallbackTool?: string; fallbackInput?: Record<string, string> };
    onToolFailure?: { strategy: string; maxRetries?: number; fallbackTool?: string };
    }

    export class ConfigValidator {
    private ajv: Ajv;
    private validateFn: ValidateFunction;

    constructor() {
    this.ajv = new Ajv({ allErrors: true, strict: true });
    addFormats(this.ajv);

    // JSON Schema 定义 Agent 配置的结构约束
    const schema = {
    type: 'object',
    required: ['apiVersion', 'kind', 'metadata', 'spec'],
    properties: {
    apiVersion: { type: 'string', pattern: '^agent\\\\.orb/v[0-9]+$' },
    kind: { type: 'string', enum: ['AgentPipeline'] },
    metadata: {
    type: 'object',
    required: ['name'],
    properties: {
    name: { type: 'string', minLength: 3, maxLength: 64 },
    },
    },
    spec: {
    type: 'object',
    required: ['steps'],
    properties: {
    resources: {
    type: 'object',
    required: ['globalTimeoutMs'],
    properties: {
    globalTimeoutMs: { type: 'number', minimum: 1000 },
    maxConcurrentSteps: { type: 'number', minimum: 1, maximum: 10 },
    },
    },
    steps: {
    type: 'array',
    minItems: 1,
    items: {
    type: 'object',
    required: ['name', 'tool', 'version', 'input', 'timeoutMs'],
    properties: {
    name: { type: 'string', minLength: 2 },
    tool: { type: 'string', minLength: 1 },
    timeoutMs: { type: 'number', minimum: 100, maximum: 120000 },
    },
    },
    },
    },
    },
    },
    };

    this.validateFn = this.ajv.compile(schema);
    }

    /** 验证配置文件的结构合法性 */
    validateStructure(config: AgentConfig): ValidationResult {
    const errors: string[] = [];
    const warnings: string[] = [];

    // Ajv 结构验证
    if (!this.validateFn(config)) {
    for (const err of this.validateFn.errors ?? []) {
    errors.push(`Schema error at ${err.instancePath}: ${err.message}`);
    }
    }

    // 自定义业务规则验证
    const stepNames = config.spec.steps.map((s) => s.name);

    // 步骤名唯一性检测
    const duplicates = stepNames.filter((name, idx) => stepNames.indexOf(name) !== idx);
    if (duplicates.length > 0) {
    errors.push(`Duplicate step names: ${duplicates.join(', ')}`);
    }

    // 循环依赖检测:条件表达式引用了后面的步骤
    for (const step of config.spec.steps) {
    if (step.condition) {
    // 简化检测:条件中引用了不存在的步骤名
    for (const otherStep of config.spec.steps) {
    if (step.condition.includes(otherStep.name) && otherStep !== step) {
    // 检查引用步骤是否在当前步骤之前
    const refIdx = config.spec.steps.indexOf(otherStep);
    const curIdx = config.spec.steps.indexOf(step);
    if (refIdx > curIdx) {
    errors.push(
    `Step "${step.name}" condition references future step "${otherStep.name}" — circular dependency`
    );
    }
    }
    }
    }
    }

    // 超时总和检测:所有步骤超时之和不应超过全局超时
    const totalStepTimeout = config.spec.steps.reduce((sum, s) => sum + s.timeoutMs, 0);
    const globalTimeout = config.spec.resources?.globalTimeoutMs ?? 60000;
    if (totalStepTimeout > globalTimeout) {
    warnings.push(
    `Sum of step timeouts (${totalStepTimeout}ms) exceeds global timeout (${globalTimeout}ms)`
    );
    }

    return {
    valid: errors.length === 0,
    errors,
    warnings,
    };
    }
    }

    3.3 运行时执行引擎

    // pipeline-engine.ts —— 声明式配置的运行时执行引擎
    // 引擎只负责"按配置执行",不包含任何业务逻辑——逻辑全在配置里

    import { ConfigValidator, ValidationResult } from './config-validator';

    /** 工具注册表:name → executor function */
    type ToolExecutor = (params: Record<string, unknown>) => Promise<unknown>;

    export class PipelineEngine {
    private tools: Map<string, ToolExecutor> = new Map();
    private validator: ConfigValidator;

    constructor() {
    this.validator = new ConfigValidator();
    }

    /** 注册工具执行器 */
    registerTool(name: string, executor: ToolExecutor): void {
    this.tools.set(name, executor);
    }

    /** 加载并验证配置,构建执行上下文 */
    async execute(config: AgentConfig, userInput: Record<string, unknown>): Promise<{
    result: unknown;
    trace: Array<{ step: string; status: string; durationMs: number }>;
    }> {
    // 静态验证:配置结构合法性
    const validation: ValidationResult = this.validator.validateStructure(config);
    if (!validation.valid) {
    throw new Error(`Invalid config: ${validation.errors.join('; ')}`);
    }

    // 动态验证:工具是否已注册
    for (const step of config.spec.steps) {
    if (!this.tools.has(step.tool)) {
    throw new Error(`Tool not registered: ${step.tool} (step: ${step.name})`);
    }
    }

    // 初始化执行上下文:存储步骤输出和全局变量
    const context: Record<string, unknown> = {
    user_input: userInput.text,
    metadata: userInput.metadata ?? {},
    steps: {} as Record<string, Record<string, unknown>>,
    };

    const trace: Array<{ step: string; status: string; durationMs: number }> = [];
    const globalStart = Date.now();

    // 按配置顺序执行步骤
    for (const step of config.spec.steps) {
    // 条件前置:不满足则跳过
    if (step.condition) {
    const shouldExecute = this.evaluateCondition(step.condition, context);
    if (!shouldExecute) {
    trace.push({ step: step.name, status: 'skipped', durationMs: 0 });
    continue;
    }
    }

    // 输入映射:从上下文中提取变量值
    const inputParams = this.resolveInput(step.input, context);

    // 执行工具调用,带超时和重试
    const result = await this.executeWithRetry(
    step.tool,
    inputParams,
    step.timeoutMs,
    step.retryPolicy
    );

    // 输出绑定:工具返回值写入上下文
    if (step.output && result.status === 'success') {
    context.steps[step.name] = this.resolveOutput(step.output, result.data);
    }

    trace.push({
    step: step.name,
    status: result.status,
    durationMs: result.durationMs,
    });

    // 全局超时检测
    if (Date.now() – globalStart > config.spec.resources?.globalTimeoutMs ?? 60000) {
    trace.push({ step: '__global_timeout__', status: 'timeout', durationMs: 0 });
    break;
    }
    }

    // 返回最终结果和执行追踪
    const lastStep = config.spec.steps[config.spec.steps.length – 1];
    return {
    result: context.steps[lastStep.name] ?? null,
    trace,
    };
    }

    /** 条件表达式求值:简化版模板引擎 */
    private evaluateCondition(condition: string, context: Record<string, unknown>): boolean {
    try {
    // 将 {{ }} 模板替换为上下文中的实际值
    let expr = condition.replace(/\\{\\{([^}]+)\\}\\}/g, (_, path) => {
    const value = this.getPathValue(path.trim(), context);
    return JSON.stringify(value);
    });
    // 安全求值:仅允许比较表达式,不允许任意 JS 执行
    // 使用 Function 构造器限制作用域
    const fn = new Function('return ' + expr);
    return fn() === true;
    } catch {
    // 条件求值失败:默认不跳过(保守策略)
    return true;
    }
    }

    /** 从上下文中按路径取值 */
    private getPathValue(path: string, context: Record<string, unknown>): unknown {
    const parts = path.split('.');
    let current: unknown = context;
    for (const part of parts) {
    if (current && typeof current === 'object') {
    current = (current as Record<string, unknown>)[part];
    } else {
    return undefined;
    }
    }
    return current;
    }

    /** 带超时和重试的工具执行 */
    private async executeWithRetry(
    toolName: string,
    params: Record<string, unknown>,
    timeoutMs: number,
    retryPolicy?: { maxRetries: number; backoffMs: number }
    ): Promise<{ status: string; data: unknown; durationMs: number }> {
    const executor = this.tools.get(toolName)!;
    const maxRetries = retryPolicy?.maxRetries ?? 0;
    const backoffMs = retryPolicy?.backoffMs ?? 500;

    for (let attempt = 0; attempt <= maxRetries; attempt++) {
    const start = Date.now();
    try {
    const data = await Promise.race([
    executor(params),
    new Promise<never>((_, reject) =>
    setTimeout(() => reject(new Error('Timeout')), timeoutMs)
    ),
    ]);
    return { status: 'success', data, durationMs: Date.now() – start };
    } catch (err) {
    if (attempt < maxRetries) {
    await new Promise((r) => setTimeout(r, backoffMs * (attempt + 1)));
    } else {
    return {
    status: 'failure',
    data: null,
    durationMs: Date.now() – start,
    };
    }
    }
    }
    return { status: 'failure', data: null, durationMs: 0 };
    }

    /** 输入映射解析:{{ }} 模板替换 */
    private resolveInput(
    input: Record<string, string>,
    context: Record<string, unknown>
    ): Record<string, unknown> {
    const resolved: Record<string, unknown> = {};
    for (const [key, template] of Object.entries(input)) {
    resolved[key] = this.resolveTemplate(template, context);
    }
    return resolved;
    }

    /** 模板字符串解析 */
    private resolveTemplate(template: string, context: Record<string, unknown>): unknown {
    if (!template.includes('{{')) return template;

    return template.replace(/\\{\\{([^}]+)\\}\\}/g, (_, path) => {
    const trimmed = path.trim();
    // 支持 default 过滤器:{{ value | default('fallback') }}
    const parts = trimmed.split('|');
    const valuePath = parts[0].trim();
    const value = this.getPathValue(valuePath, context);
    if (value !== undefined && value !== null) return JSON.stringify(value);
    // default 过滤器
    if (parts.length > 1) {
    const defaultExpr = parts[1].trim();
    const defaultMatch = defaultExpr.match(/default\\('([^']*)'\\)/);
    if (defaultMatch) return defaultMatch[1];
    }
    return 'null';
    });
    }
    }

    四、边界分析与架构权衡

    4.1 条件表达式安全性

    模板引擎用 new Function() 执行条件表达式,理论上可以注入恶意代码。虽然模板变量来自上下文而非用户输入,但用户输入最终会写入上下文。

    对策:条件表达式只允许比较运算(>, <, ==, !=),禁止函数调用和赋值。生产环境建议用专门的模板引擎(如 jsonpath-plus),避免 new Function()。

    4.2 配置热更新的原子性

    改配置后立即生效,但如果新配置有错误(比如引用了不存在的工具),正在执行的链路会崩溃。你需要保证"配置更新时,已启动的链路用旧配置执行完毕,新链路用新配置"。

    对策:配置版本化管理——每个配置有版本号,执行引擎启动链路时锁定当前版本,新版本只对后续启动的链路生效。

    4.3 适用边界与禁用场景

    • 适用:多条业务链路共享同一执行引擎、频繁调整步骤顺序和条件、需要可视化展示 Agent 拓扑
    • 禁用:单条简单链路(配置文件比代码还复杂)、条件表达式需要复杂计算(模板引擎表达力不足)、对执行性能要求极高(模板解析增加延迟)

    4.4 与 DAG 编排引擎的对比

    Argo Workflows、Temporal 等 DAG 编排引擎也是声明式配置,但它们侧重于任务调度(并发、依赖、重试),而 Agent 编排侧重于实时决策(条件跳转、降级、上下文传递)。两者可以结合:Agent 编排配置定义决策逻辑,DAG 引擎处理调度和重试。

    五、总结

    声明式配置把 Agent 编排逻辑从代码中分离出来,用 YAML 定义步骤、条件、错误策略,运行引擎只负责"按配置执行"。核心收益:热更新、可视化、逻辑与引擎解耦。代价:配置验证的复杂度、条件表达式的安全性、热更新的原子性保证。配置版本化是解决原子性的关键——旧链路用旧配置,新链路用新配置,不交叉。模板引擎的安全限制是重中之重——只允许比较运算,禁止函数调用。

    赞(0)
    未经允许不得转载:171主机测评 » AI Agent 编排的声明式配置:像写 K8s YAML 一样定义 Agent
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址