{"version":3,"sources":["../../../../src/adapters/a2a/serve/agent_executor.ts"],"names":["logger","Logger","root","child","name","BaseA2AAgentExecutor","abortControllers","Map","agent","execute","requestContext","eventBus","userMessage","existingTask","task","kind","id","uuidv4","contextId","status","state","timestamp","Date","toISOString","history","metadata","taskId","publish","workingStatusUpdate","final","abortController","AbortController","set","clone","memory","reset","addMany","map","convertA2AMessageToFrameworkMessage","response","run","signal","observe","emitter","process_events","agentMessage","role","messageId","parts","text","result","push","finalUpdate","message","finished","error","errorMessage","Error","String","errorUpdate","_emitter","_requestContext","_eventBus","cancelTask","has","get","cancelledUpdate","abort","ToolCallingAgentExecutor","lastMsg","undefined","processEvent","messages","at","index","lastIndexOf","slice","update","JSON","stringify","toPlain","on","ReActAgentExecutor","updateEvent","value"],"mappings":";;;;;;;;AAiBA,MAAMA,MAAAA,GAASC,iBAAAA,CAAOC,IAAAA,CAAKC,KAAAA,CAAM;EAC/BC,IAAAA,EAAM;AACR,CAAA,CAAA;AAEO,MAAeC,oBAAAA,CAAAA;EArBtB;;;;AAsBqBC,EAAAA,gBAAAA,uBAAuBC,GAAAA,EAAAA;AAE1C,EAAA,WAAA,CAAsBC,KAAAA,EAAiB;SAAjBA,KAAAA,GAAAA,KAAAA;AAAkB,EAAA;EAExC,MAAMC,OAAAA,CAAQC,gBAAgCC,QAAAA,EAA4C;AACxF,IAAA,MAAMC,cAAcF,cAAAA,CAAeE,WAAAA;AACnC,IAAA,MAAMC,YAAAA,GAAeH,eAAeI,IAAAA,IAAQ;MAC1CC,IAAAA,EAAM,MAAA;AACNC,MAAAA,EAAAA,EAAIC,OAAAA,EAAAA;MACJC,SAAAA,EAAWN,WAAAA,CAAYM,aAAaD,OAAAA,EAAAA;MACpCE,MAAAA,EAAQ;QACNC,KAAAA,EAAO,WAAA;QACPC,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,OAAA;MACAC,OAAAA,EAAS;AAACZ,QAAAA;;AACVa,MAAAA,QAAAA,EAAUb,WAAAA,CAAYa;AACxB,KAAA;AAEA,IAAA,MAAMC,SAASb,YAAAA,CAAaG,EAAAA;AAC5B,IAAA,MAAME,YAAYL,YAAAA,CAAaK,SAAAA;AAG/BP,IAAAA,QAAAA,CAASgB,QAAQd,YAAAA,CAAAA;AAGjB,IAAA,MAAMe,mBAAAA,GAA6C;MACjDb,IAAAA,EAAM,eAAA;AACNW,MAAAA,MAAAA;AACAR,MAAAA,SAAAA;MACAC,MAAAA,EAAQ;QACNC,KAAAA,EAAO,SAAA;QACPC,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,OAAA;MACAM,KAAAA,EAAO;AACT,KAAA;AACAlB,IAAAA,QAAAA,CAASgB,QAAQC,mBAAAA,CAAAA;AAEjB,IAAA,MAAME,eAAAA,GAAkB,IAAIC,eAAAA,EAAAA;AAC5B,IAAA,IAAA,CAAKzB,gBAAAA,CAAiB0B,IAAIN,MAAAA,EAAQ;AAACR,MAAAA,SAAAA;AAAWY,MAAAA;AAAgB,KAAA,CAAA;AAE9D,IAAA,MAAMtB,KAAAA,GAAQ,MAAM,IAAA,CAAKA,KAAAA,CAAMyB,KAAAA,EAAK;AACpCzB,IAAAA,KAAAA,CAAM0B,MAAAA,GAAS,MAAM1B,KAAAA,CAAM0B,MAAAA,CAAOD,KAAAA,EAAK;AACvCzB,IAAAA,KAAAA,CAAM0B,OAAOC,KAAAA,EAAK;AAClB,IAAA,MAAM3B,KAAAA,CAAM0B,MAAAA,CAAOE,OAAAA,CAAAA,CAChBvB,YAAAA,CAAaW,OAAAA,IAAW;MAACd,cAAAA,CAAeE;OAAcyB,GAAAA,CACrDC,6CAAAA,CAAAA,IACG,EAAE,CAAA;AAGT,IAAA,IAAI;AAEF,MAAA,MAAMC,QAAAA,GAAW,MAAM/B,KAAAA,CACpBgC,GAAAA,CACC,EAAC,EACD;AACEC,QAAAA,MAAAA,EAAQX,eAAAA,CAAgBW;OAC1B,CAAA,CAEDC,OAAAA,CAAQ,OAAOC,OAAAA,KAAAA;AACd,QAAA,MAAM,IAAA,CAAKC,cAAAA,CAAeD,OAAAA,EAASjC,cAAAA,EAAgBC,QAAAA,CAAAA;MACrD,CAAA,CAAA;AAEF,MAAA,MAAMkC,YAAAA,GAAwB;QAC5B9B,IAAAA,EAAM,SAAA;QACN+B,IAAAA,EAAM,OAAA;AACNC,QAAAA,SAAAA,EAAW9B,OAAAA,EAAAA;QACX+B,KAAAA,EAAO;AAAC,UAAA;YAAEjC,IAAAA,EAAM,MAAA;AAAQkC,YAAAA,IAAAA,EAAMV,SAASW,MAAAA,CAAOD;AAAK;;AACnDvB,QAAAA,MAAAA;AACAR,QAAAA;AACF,OAAA;AAGA,MAAA,IAAI,CAACL,aAAaW,OAAAA,EAAS;AACzBX,QAAAA,YAAAA,CAAaW,UAAU,EAAA;AACzB,MAAA;AACAX,MAAAA,YAAAA,CAAaW,OAAAA,CAAQ2B,KAAKN,YAAAA,CAAAA;AAC1BlC,MAAAA,QAAAA,CAASgB,QAAQd,YAAAA,CAAAA;AAGjB,MAAA,MAAMuC,WAAAA,GAAqC;QACzCrC,IAAAA,EAAM,eAAA;AACNW,QAAAA,MAAAA;AACAR,QAAAA,SAAAA;QACAC,MAAAA,EAAQ;UACNC,KAAAA,EAAO,WAAA;UACPiC,OAAAA,EAASR,YAAAA;UACTxB,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,SAAA;QACAM,KAAAA,EAAO;AACT,OAAA;AACAlB,MAAAA,QAAAA,CAASgB,QAAQyB,WAAAA,CAAAA;AAEjBzC,MAAAA,QAAAA,CAAS2C,QAAAA,EAAQ;AACnB,IAAA,CAAA,CAAA,OAASC,KAAAA,EAAO;AACdvD,MAAAA,MAAAA,CAAOuD,KAAAA,CAAMA,OAAO,uBAAA,CAAA;AACpB,MAAA,MAAMC,eAAeD,KAAAA,YAAiBE,KAAAA,GAAQF,KAAAA,CAAMF,OAAAA,GAAUK,OAAOH,KAAAA,CAAAA;AAErE,MAAA,MAAMI,WAAAA,GAAqC;QACzC5C,IAAAA,EAAM,eAAA;AACNW,QAAAA,MAAAA;AACAR,QAAAA,SAAAA;QACAC,MAAAA,EAAQ;UACNC,KAAAA,EAAO,QAAA;UACPiC,OAAAA,EAAS;YACPtC,IAAAA,EAAM,SAAA;YACN+B,IAAAA,EAAM,OAAA;AACNC,YAAAA,SAAAA,EAAW9B,OAAAA,EAAAA;YACX+B,KAAAA,EAAO;AAAC,cAAA;gBAAEjC,IAAAA,EAAM,MAAA;AAAQkC,gBAAAA,IAAAA,EAAM,gBAAgBO,YAAAA,CAAAA;AAAe;;AAC7D9B,YAAAA,MAAAA;AACAR,YAAAA;AACF,WAAA;UACAG,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,SAAA;QACAM,KAAAA,EAAO;AACT,OAAA;AACAlB,MAAAA,QAAAA,CAASgB,QAAQgC,WAAAA,CAAAA;AACjBhD,MAAAA,QAAAA,CAAS2C,QAAAA,EAAQ;AACnB,IAAA;AACF,EAAA;EAEA,MAAMV,cAAAA,CACJgB,QAAAA,EACAC,eAAAA,EACAC,SAAAA,EACe;AACf,IAAA;AACF,EAAA;EAEOC,UAAAA,mBAAa,MAAA,CAAA,OAAOrC,QAAgBf,QAAAA,KAAAA;AACzC,IAAA,IAAI,IAAA,CAAKL,gBAAAA,CAAiB0D,GAAAA,CAAItC,MAAAA,CAAAA,EAAS;AACrC,MAAA,MAAM,CAACR,SAAAA,EAAWY,eAAAA,IAAmB,IAAA,CAAKxB,gBAAAA,CAAiB2D,IAAIvC,MAAAA,CAAAA;AAC/D,MAAA,MAAMwC,eAAAA,GAAyC;QAC7CnD,IAAAA,EAAM,eAAA;AACNW,QAAAA,MAAAA;AACAR,QAAAA,SAAAA;QACAC,MAAAA,EAAQ;UACNC,KAAAA,EAAO,UAAA;UACPC,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,SAAA;QACAM,KAAAA,EAAO;AACT,OAAA;AACAlB,MAAAA,QAAAA,CAASgB,QAAQuC,eAAAA,CAAAA;AACjBpC,MAAAA,eAAAA,CAAgBqC,KAAAA,EAAK;AACvB,IAAA;EACF,CAAA,EAhBoB,YAAA,CAAA;AAiBtB;AAEO,MAAMC,iCAAiC/D,oBAAAA,CAAAA;EAzK9C;;;EA0KE,MAAMuC,cAAAA,CACJD,OAAAA,EACAjC,cAAAA,EACAC,QAAAA,EACe;AACf,IAAA,IAAI0D,OAAAA,GAAwCC,MAAAA;AAE5C,IAAA,MAAMC,YAAAA,mBAA8D,MAAA,CAAA,OAAO,EAAEnD,KAAAA,EAAK,KAAE;AAClF,MAAA,MAAMoD,QAAAA,GAAWpD,MAAMc,MAAAA,CAAOsC,QAAAA;AAC9B,MAAA,IAAIH,YAAYC,MAAAA,EAAW;AACzBD,QAAAA,OAAAA,GAAUG,QAAAA,CAASC,GAAG,EAAC,CAAA;AACzB,MAAA;AAEA,MAAA,MAAMC,KAAAA,GAAQL,OAAAA,GAAUG,QAAAA,CAASG,WAAAA,CAAYN,OAAAA,CAAAA,GAAW,EAAA;AACxD,MAAA,KAAA,MAAWhB,OAAAA,IAAWmB,QAAAA,CAASI,KAAAA,CAAMF,KAAAA,GAAQ,CAAA,CAAA,EAAI;AAC/C,QAAA,MAAMG,MAAAA,GAAgC;UACpC9D,IAAAA,EAAM,eAAA;AACNW,UAAAA,MAAAA,EAAQhB,cAAAA,CAAegB,MAAAA;AACvBR,UAAAA,SAAAA,EAAWR,cAAAA,CAAeQ,SAAAA;UAC1BC,MAAAA,EAAQ;YACNC,KAAAA,EAAO,SAAA;YACPiC,OAAAA,EAAS;cACPtC,IAAAA,EAAM,SAAA;cACN+B,IAAAA,EAAM,OAAA;AACNC,cAAAA,SAAAA,EAAW9B,OAAAA,EAAAA;cACX+B,KAAAA,EAAO;AAAC,gBAAA;kBAAEjC,IAAAA,EAAM,MAAA;AAAQkC,kBAAAA,IAAAA,EAAM6B,IAAAA,CAAKC,SAAAA,CAAU1B,OAAAA,CAAQ2B,OAAAA,EAAO;AAAI;;AAChEtD,cAAAA,MAAAA,EAAQhB,cAAAA,CAAegB,MAAAA;AACvBR,cAAAA,SAAAA,EAAWR,cAAAA,CAAeQ;AAC5B,aAAA;YACAG,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,WAAA;UACAM,KAAAA,EAAO;AACT,SAAA;AACAlB,QAAAA,QAAAA,CAASgB,QAAQkD,MAAAA,CAAAA;AACjBR,QAAAA,OAAAA,GAAUhB,OAAAA;AACZ,MAAA;IACF,CAAA,EA7BoE,cAAA,CAAA;AA8BpEV,IAAAA,OAAAA,CAAQsC,EAAAA,CAAG,SAASV,YAAAA,CAAAA;AACpB5B,IAAAA,OAAAA,CAAQsC,EAAAA,CAAG,WAAWV,YAAAA,CAAAA;AACxB,EAAA;AACF;AAEO,MAAMW,2BAA2B7E,oBAAAA,CAAAA;EApNxC;;;EAqNE,MAAMuC,cAAAA,CACJD,OAAAA,EACAjC,cAAAA,EACAC,QAAAA,EACe;AACfgC,IAAAA,OAAAA,CAAQsC,EAAAA,CAAG,QAAA,EAAU,OAAO,EAAEJ,QAAM,KAAE;AACpC,MAAA,MAAMM,WAAAA,GAAqC;QACzCpE,IAAAA,EAAM,eAAA;AACNW,QAAAA,MAAAA,EAAQhB,cAAAA,CAAegB,MAAAA;AACvBR,QAAAA,SAAAA,EAAWR,cAAAA,CAAeQ,SAAAA;QAC1BC,MAAAA,EAAQ;UACNC,KAAAA,EAAO,SAAA;UACPiC,OAAAA,EAAS;YACPtC,IAAAA,EAAM,SAAA;YACN+B,IAAAA,EAAM,OAAA;AACNC,YAAAA,SAAAA,EAAW9B,OAAAA,EAAAA;YACX+B,KAAAA,EAAO;AAAC,cAAA;gBAAEjC,IAAAA,EAAM,MAAA;AAAQkC,gBAAAA,IAAAA,EAAM4B,MAAAA,CAAOO;AAAM;;AAC3C1D,YAAAA,MAAAA,EAAQhB,cAAAA,CAAegB,MAAAA;AACvBR,YAAAA,SAAAA,EAAWR,cAAAA,CAAeQ;AAC5B,WAAA;UACAG,SAAAA,EAAAA,iBAAW,IAAIC,IAAAA,EAAAA,EAAOC,WAAAA;AACxB,SAAA;QACAM,KAAAA,EAAO;AACT,OAAA;AACAlB,MAAAA,QAAAA,CAASgB,QAAQwD,WAAAA,CAAAA;IACnB,CAAA,CAAA;AACF,EAAA;AACF","file":"agent_executor.cjs","sourcesContent":["/**\n * Copyright 2025 © BeeAI a Series of LF Projects, LLC\n * SPDX-License-Identifier: Apache-2.0\n */\n\nimport { v4 as uuidv4 } from \"uuid\";\n\nimport { AnyAgent } from \"@/agents/types.js\";\nimport { AgentExecutor, RequestContext, ExecutionEventBus } from \"@a2a-js/sdk/server\";\nimport { Message, TaskStatusUpdateEvent } from \"@a2a-js/sdk\";\nimport { Logger } from \"@/logger/logger.js\";\nimport { convertA2AMessageToFrameworkMessage } from \"@/adapters/a2a/agents/utils.js\";\nimport { Callback, Emitter } from \"@/emitter/emitter.js\";\nimport { ToolCallingAgentCallbacks, ToolCallingAgentRunState } from \"@/agents/toolCalling/types.js\";\nimport { ReActAgentCallbacks } from \"@/agents/react/types.js\";\nimport { Message as FrameworkMessage } from \"@/backend/message.js\";\n\nconst logger = Logger.root.child({\n  name: \"A2A server\",\n});\n\nexport abstract class BaseA2AAgentExecutor implements AgentExecutor {\n  protected readonly abortControllers = new Map<string, [string, AbortController]>();\n\n  constructor(protected agent: AnyAgent) {}\n\n  async execute(requestContext: RequestContext, eventBus: ExecutionEventBus): Promise<void> {\n    const userMessage = requestContext.userMessage;\n    const existingTask = requestContext.task || {\n      kind: \"task\",\n      id: uuidv4(),\n      contextId: userMessage.contextId || uuidv4(),\n      status: {\n        state: \"submitted\",\n        timestamp: new Date().toISOString(),\n      },\n      history: [userMessage],\n      metadata: userMessage.metadata,\n    };\n\n    const taskId = existingTask.id;\n    const contextId = existingTask.contextId;\n\n    // Publish initial Task event if it's a new task\n    eventBus.publish(existingTask);\n\n    // Publish \"working\" status update\n    const workingStatusUpdate: TaskStatusUpdateEvent = {\n      kind: \"status-update\",\n      taskId: taskId,\n      contextId: contextId,\n      status: {\n        state: \"working\",\n        timestamp: new Date().toISOString(),\n      },\n      final: false,\n    };\n    eventBus.publish(workingStatusUpdate);\n\n    const abortController = new AbortController();\n    this.abortControllers.set(taskId, [contextId, abortController]);\n\n    const agent = await this.agent.clone();\n    agent.memory = await agent.memory.clone();\n    agent.memory.reset();\n    await agent.memory.addMany(\n      (existingTask.history || [requestContext.userMessage]).map(\n        convertA2AMessageToFrameworkMessage,\n      ) || [],\n    );\n\n    try {\n      // run the agent\n      const response = await agent\n        .run(\n          {},\n          {\n            signal: abortController.signal,\n          },\n        )\n        .observe(async (emitter) => {\n          await this.process_events(emitter, requestContext, eventBus);\n        });\n\n      const agentMessage: Message = {\n        kind: \"message\",\n        role: \"agent\",\n        messageId: uuidv4(),\n        parts: [{ kind: \"text\", text: response.result.text }],\n        taskId: taskId,\n        contextId: contextId,\n      };\n\n      // Append agent message to task history\n      if (!existingTask.history) {\n        existingTask.history = [];\n      }\n      existingTask.history.push(agentMessage);\n      eventBus.publish(existingTask);\n\n      // Publish completed status update with agent message\n      const finalUpdate: TaskStatusUpdateEvent = {\n        kind: \"status-update\",\n        taskId: taskId,\n        contextId: contextId,\n        status: {\n          state: \"completed\",\n          message: agentMessage,\n          timestamp: new Date().toISOString(),\n        },\n        final: true,\n      };\n      eventBus.publish(finalUpdate);\n\n      eventBus.finished();\n    } catch (error) {\n      logger.error(error, \"Agent execution error\");\n      const errorMessage = error instanceof Error ? error.message : String(error);\n      // Publish failed status update\n      const errorUpdate: TaskStatusUpdateEvent = {\n        kind: \"status-update\",\n        taskId: taskId,\n        contextId: contextId,\n        status: {\n          state: \"failed\",\n          message: {\n            kind: \"message\",\n            role: \"agent\",\n            messageId: uuidv4(),\n            parts: [{ kind: \"text\", text: `Agent error: ${errorMessage}` }],\n            taskId: taskId,\n            contextId: contextId,\n          },\n          timestamp: new Date().toISOString(),\n        },\n        final: true,\n      };\n      eventBus.publish(errorUpdate);\n      eventBus.finished();\n    }\n  }\n\n  async process_events(\n    _emitter: Emitter,\n    _requestContext: RequestContext,\n    _eventBus: ExecutionEventBus,\n  ): Promise<void> {\n    return;\n  }\n\n  public cancelTask = async (taskId: string, eventBus: ExecutionEventBus): Promise<void> => {\n    if (this.abortControllers.has(taskId)) {\n      const [contextId, abortController] = this.abortControllers.get(taskId)!;\n      const cancelledUpdate: TaskStatusUpdateEvent = {\n        kind: \"status-update\",\n        taskId: taskId,\n        contextId: contextId,\n        status: {\n          state: \"canceled\",\n          timestamp: new Date().toISOString(),\n        },\n        final: true,\n      };\n      eventBus.publish(cancelledUpdate);\n      abortController.abort();\n    }\n  };\n}\n\nexport class ToolCallingAgentExecutor extends BaseA2AAgentExecutor {\n  async process_events(\n    emitter: Emitter<ToolCallingAgentCallbacks>,\n    requestContext: RequestContext,\n    eventBus: ExecutionEventBus,\n  ): Promise<void> {\n    let lastMsg: FrameworkMessage | undefined = undefined;\n\n    const processEvent: Callback<{ state: ToolCallingAgentRunState }> = async ({ state }) => {\n      const messages = state.memory.messages;\n      if (lastMsg === undefined) {\n        lastMsg = messages.at(-1);\n      }\n\n      const index = lastMsg ? messages.lastIndexOf(lastMsg) : -1;\n      for (const message of messages.slice(index + 1)) {\n        const update: TaskStatusUpdateEvent = {\n          kind: \"status-update\",\n          taskId: requestContext.taskId,\n          contextId: requestContext.contextId,\n          status: {\n            state: \"working\",\n            message: {\n              kind: \"message\",\n              role: \"agent\",\n              messageId: uuidv4(),\n              parts: [{ kind: \"text\", text: JSON.stringify(message.toPlain()) }],\n              taskId: requestContext.taskId,\n              contextId: requestContext.contextId,\n            },\n            timestamp: new Date().toISOString(),\n          },\n          final: false,\n        };\n        eventBus.publish(update);\n        lastMsg = message;\n      }\n    };\n    emitter.on(\"start\", processEvent);\n    emitter.on(\"success\", processEvent);\n  }\n}\n\nexport class ReActAgentExecutor extends BaseA2AAgentExecutor {\n  async process_events(\n    emitter: Emitter<ReActAgentCallbacks>,\n    requestContext: RequestContext,\n    eventBus: ExecutionEventBus,\n  ): Promise<void> {\n    emitter.on(\"update\", async ({ update }) => {\n      const updateEvent: TaskStatusUpdateEvent = {\n        kind: \"status-update\",\n        taskId: requestContext.taskId,\n        contextId: requestContext.contextId,\n        status: {\n          state: \"working\",\n          message: {\n            kind: \"message\",\n            role: \"agent\",\n            messageId: uuidv4(),\n            parts: [{ kind: \"text\", text: update.value }],\n            taskId: requestContext.taskId,\n            contextId: requestContext.contextId,\n          },\n          timestamp: new Date().toISOString(),\n        },\n        final: false,\n      };\n      eventBus.publish(updateEvent);\n    });\n  }\n}\n"]}