{"version":3,"sources":["../../src/middleware/streamToolCall.ts"],"names":["StreamToolCallMiddleware","Middleware","target","key","matchNested","forceStreaming","cleanups","output","ChatModelOutput","buffer","delta","emitter","options","Emitter","root","child","namespace","reset","add","chunk","handleNewToken","value","callbacks","abort","isEmpty","length","bind","ctx","push","instance","match","meta","creator","ChatModel","name","handleStart","createEmitterOptions","handleSuccess","unbind","fn","shift","isBlocking","process","toolName","args","parsedArgs","isString","parseBrokenJson","pair","outputStructured","parse","catch","hasProp","JSON","stringify","slice","emit","data","_meta","input","stream","streamPartialToolCalls","messages","merge","toolCalls","getToolCalls","toolCall","textContent","getTextContent","isPlainObject","parameters"],"mappings":";;;;;;;;;;;AAuEO,MAAMA,iCAA0CC,sBAAAA,CAAAA;EAvEvD;;;AAwEmBC,EAAAA,MAAAA;AACAC,EAAAA,GAAAA;AACAC,EAAAA,WAAAA;AACAC,EAAAA,cAAAA;AACAC,EAAAA,QAAAA,GAA2B,EAAA;EAEpCC,MAAAA,GAAS,IAAIC,wBAAAA,CAAgB,EAAE,CAAA;EAC/BC,MAAAA,GAAS,EAAA;EACTC,KAAAA,GAAQ,EAAA;AAEAC,EAAAA,OAAAA;AAEhB,EAAA,WAAA,CAAYC,OAAAA,EAA0C;AACpD,IAAA,KAAA,EAAK;AAEL,IAAA,IAAA,CAAKV,SAASU,OAAAA,CAAQV,MAAAA;AACtB,IAAA,IAAA,CAAKC,MAAMS,OAAAA,CAAQT,GAAAA;AACnB,IAAA,IAAA,CAAKC,WAAAA,GAAcQ,QAAQR,WAAAA,IAAe,KAAA;AAC1C,IAAA,IAAA,CAAKC,cAAAA,GAAiBO,QAAQP,cAAAA,IAAkB,KAAA;AAEhD,IAAA,IAAA,CAAKM,OAAAA,GAAUE,mBAAAA,CAAQC,IAAAA,CAAKC,KAAAA,CAA4C;MACtEC,SAAAA,EAAW;AAAC,QAAA,YAAA;AAAc,QAAA;;KAC5B,CAAA;AACF,EAAA;EAEAC,KAAAA,GAAQ;AACN,IAAA,IAAA,CAAKV,MAAAA,GAAS,IAAIC,wBAAAA,CAAgB,EAAE,CAAA;AACpC,IAAA,IAAA,CAAKC,MAAAA,GAAS,EAAA;AACd,IAAA,IAAA,CAAKC,KAAAA,GAAQ,EAAA;AACf,EAAA;AAEA,EAAA,MAAMQ,IAAIC,KAAAA,EAAwB;AAChC,IAAA,MAAM,KAAKC,cAAAA,CACT;MACEC,KAAAA,EAAOF,KAAAA;MACPG,SAAAA,EAAW;AACTC,QAAAA,KAAAA,kBAAO,MAAA,CAAA,MAAA;QAAO,CAAA,EAAP,OAAA;AACT;AACF,KAAA,EACA,EAAC,CAAA;AAEL,EAAA;EAEAC,OAAAA,GAAU;AACR,IAAA,OAAO,IAAA,CAAKf,OAAOgB,MAAAA,KAAW,CAAA;AAChC,EAAA;AAEAC,EAAAA,IAAAA,CAAKC,GAAAA,EAAoC;AACvC,IAAA,IAAA,CAAKV,KAAAA,EAAK;AAGV,IAAA,IAAA,CAAKX,QAAAA,CAASsB,KACZD,GAAAA,CAAIE,QAAAA,CAASlB,QAAQmB,KAAAA,CACnB,CAACC,IAAAA,KAASA,IAAAA,CAAKC,OAAAA,YAAmBC,kBAAAA,IAAaF,KAAKG,IAAAA,KAAS,OAAA,EAC7D,KAAKC,WAAAA,CAAYT,IAAAA,CAAK,IAAI,CAAA,EAC1B,IAAA,CAAKU,oBAAAA,EAAoB,CAAA,CAAA;AAK7B,IAAA,IAAA,CAAK9B,QAAAA,CAASsB,KACZD,GAAAA,CAAIE,QAAAA,CAASlB,QAAQmB,KAAAA,CACnB,CAACC,IAAAA,KAASA,IAAAA,CAAKC,OAAAA,YAAmBC,kBAAAA,IAAaF,KAAKG,IAAAA,KAAS,UAAA,EAC7D,KAAKd,cAAAA,CAAeM,IAAAA,CAAK,IAAI,CAAA,EAC7B,IAAA,CAAKU,oBAAAA,EAAoB,CAAA,CAAA;AAK7B,IAAA,IAAA,CAAK9B,QAAAA,CAASsB,KACZD,GAAAA,CAAIE,QAAAA,CAASlB,QAAQmB,KAAAA,CACnB,CAACC,IAAAA,KAASA,IAAAA,CAAKC,OAAAA,YAAmBC,kBAAAA,IAAaF,KAAKG,IAAAA,KAAS,SAAA,EAC7D,KAAKG,aAAAA,CAAcX,IAAAA,CAAK,IAAI,CAAA,EAC5B,IAAA,CAAKU,oBAAAA,EAAoB,CAAA,CAAA;AAG/B,EAAA;EAEAE,MAAAA,GAAe;AAEb,IAAA,OAAO,IAAA,CAAKhC,QAAAA,CAASmB,MAAAA,GAAS,CAAA,EAAG;AAC/B,MAAA,MAAMc,EAAAA,GAAK,IAAA,CAAKjC,QAAAA,CAASkC,KAAAA,EAAK;AAC9BD,MAAAA,EAAAA,EAAAA;AACF,IAAA;AACF,EAAA;EAEUH,oBAAAA,GAAuC;AAC/C,IAAA,OAAO;AACLhC,MAAAA,WAAAA,EAAa,IAAA,CAAKA,WAAAA;MAClBqC,UAAAA,EAAY;AACd,KAAA;AACF,EAAA;EAEA,MAAcC,OAAAA,CAAQC,UAAkBC,IAAAA,EAA0B;AAChE,IAAA,IAAID,QAAAA,KAAa,IAAA,CAAKzC,MAAAA,CAAOgC,IAAAA,EAAM;AACjC,MAAA;AACF,IAAA;AAEA,IAAA,MAAMW,UAAAA,GAAaC,eAAAA,CAASF,IAAAA,CAAAA,GAAQG,2BAAgBH,IAAAA,EAAM;MAAEI,IAAAA,EAAM;AAAC,QAAA,GAAA;AAAK,QAAA;;AAAK,KAAA,CAAA,GAAKJ,IAAAA;AAClF,IAAA,MAAMK,gBAAAA,GAAoB,MAAM,IAAA,CAAK/C,MAAAA,CAAOgD,MAAML,UAAAA,CAAAA,CAAYM,KAAAA,CAAM,MAAM,IAAA,CAAA;AAC1E,IAAA,IAAI,CAACF,gBAAAA,EAAkB;AACrB,MAAA;AACF,IAAA;AAEA,IAAA,IAAI1C,MAAAA,GAAS,EAAA;AACb,IAAA,IAAI6C,kBAAAA,CAAQH,gBAAAA,EAAkB,IAAA,CAAK9C,GAAG,CAAA,EAAc;AAClDI,MAAAA,MAAAA,GAAU0C,gBAAAA,CAAyB,IAAA,CAAK9C,GAAG,CAAA,IAAK,EAAA;AAChD,MAAA,IAAI,CAAC2C,eAAAA,CAASvC,MAAAA,CAAAA,EAAS;AACrBA,QAAAA,MAAAA,GAAS8C,IAAAA,CAAKC,UAAU/C,MAAAA,CAAAA;AAC1B,MAAA;AACA,MAAA,IAAA,CAAKG,KAAAA,GAAQH,MAAAA,CAAOgD,KAAAA,CAAM,IAAA,CAAK9C,OAAOgB,MAAM,CAAA;AAC5C,MAAA,IAAA,CAAKhB,MAAAA,GAASF,MAAAA;AAEd,MAAA,IAAI,CAAC,KAAKG,KAAAA,EAAO;AACf,QAAA;AACF,MAAA;AACF,IAAA;AAEA,IAAA,MAAM,IAAA,CAAKC,OAAAA,CAAQ6C,IAAAA,CAAK,QAAA,EAAU;AAChCP,MAAAA,gBAAAA;AACAvC,MAAAA,KAAAA,EAAO,IAAA,CAAKA,KAAAA;AACZH,MAAAA;KACF,CAAA;AACF,EAAA;EAEA,MAAc4B,WAAAA,CACZsB,MACAC,KAAAA,EACe;AACf,IAAA,IAAI,KAAKrD,cAAAA,EAAgB;AACvBoD,MAAAA,IAAAA,CAAKE,MAAMC,MAAAA,GAAS,IAAA;AACpBH,MAAAA,IAAAA,CAAKE,MAAME,sBAAAA,GAAyB,IAAA;AACtC,IAAA;AACF,EAAA;EAEA,MAAcxB,aAAAA,CACZoB,MACAC,KAAAA,EACe;AAEf,IAAA,IAAI,IAAA,CAAKnD,MAAAA,CAAOuD,QAAAA,CAASrC,MAAAA,KAAW,CAAA,EAAG;AACrC,MAAA,MAAM,KAAKL,cAAAA,CAAe;AAAEC,QAAAA,KAAAA,EAAOoC,IAAAA,CAAKpC,KAAAA;QAAOC,SAAAA,EAAW;AAAEC,UAAAA,KAAAA,kBAAO,MAAA,CAAA,MAAA;UAAO,CAAA,EAAP,OAAA;AAAS;AAAE,OAAA,EAAGmC,KAAAA,CAAAA;AACnF,IAAA;AACF,EAAA;EAEA,MAActC,cAAAA,CACZqC,MACAC,KAAAA,EACe;AACf,IAAA,MAAM,IAAA,CAAKnD,MAAAA,CAAOwD,KAAAA,CAAMN,IAAAA,CAAKpC,KAAK,CAAA;AAElC,IAAA,MAAM2C,SAAAA,GAAY,IAAA,CAAKzD,MAAAA,CAAO0D,YAAAA,EAAY;AAC1C,IAAA,IAAID,SAAAA,CAAUvC,SAAS,CAAA,EAAG;AACxB,MAAA,KAAA,MAAWyC,aAAYF,SAAAA,EAAW;AAChC,QAAA,MAAM,IAAA,CAAKtB,OAAAA,CAAQwB,SAAAA,CAASvB,QAAAA,EAAUuB,UAASP,KAAK,CAAA;AACtD,MAAA;AACA,MAAA;AACF,IAAA;AAGA,IAAA,MAAMQ,WAAAA,GAAc,IAAA,CAAK5D,MAAAA,CAAO6D,cAAAA,EAAc;AAC9C,IAAA,MAAMF,QAAAA,GAAWnB,2BAAgBoB,WAAAA,EAAa;MAAEnB,IAAAA,EAAM;AAAC,QAAA,GAAA;AAAK,QAAA;;KAAK,CAAA;AAEjE,IAAA,IAAI,CAACkB,QAAAA,IAAY,CAACG,oBAAAA,CAAcH,QAAAA,CAAAA,EAAW;AACzC,MAAA;AACF,IAAA;AAEA,IAAA,MAAM,IAAA,CAAKxB,OAAAA,CAAQwB,QAAAA,CAAShC,IAAAA,EAAgBgC,SAASI,UAAU,CAAA;AACjE,EAAA;AACF","file":"streamToolCall.cjs","sourcesContent":["/**\n * Copyright 2025 © BeeAI a Series of LF Projects, LLC\n * SPDX-License-Identifier: Apache-2.0\n */\n\nimport { Middleware, RunContext, RunInstance } from \"@/context.js\";\nimport { Callback, Emitter, EventMeta } from \"@/emitter/emitter.js\";\nimport { ChatModel, ChatModelEvents } from \"@/backend/chat.js\";\nimport { ChatModelOutput } from \"@/backend/chat.js\";\nimport { Tool } from \"@/tools/base.js\";\nimport { parseBrokenJson } from \"@/internals/helpers/schema.js\";\nimport { isPlainObject, isString } from \"remeda\";\nimport { hasProp } from \"@/internals/helpers/object.js\";\nimport { EmitterOptions, InferCallbackValue } from \"@/emitter/types.js\";\n\n/**\n * Event emitted when the middleware detects an update to the target tool's arguments\n */\nexport interface StreamToolCallMiddlewareUpdateEvent<T = any> {\n  /** The validated and structured tool input */\n  outputStructured: T | null;\n  /** The current value of the target field */\n  output: string;\n  /** The incremental change since the last update */\n  delta: string;\n}\n\n/**\n * Callbacks for StreamToolCallMiddleware events\n */\nexport interface StreamToolCallMiddlewareCallbacks<T = any> {\n  update: Callback<StreamToolCallMiddlewareUpdateEvent<T>>;\n}\n\n/**\n * Options for configuring StreamToolCallMiddleware\n */\nexport interface StreamToolCallMiddlewareOptions {\n  /** The tool to monitor for streaming updates */\n  target: Tool<any>;\n  /** The field name in the tool's input schema to stream */\n  key: string;\n  /** Whether to apply middleware to nested run contexts */\n  matchNested?: boolean;\n  /** Whether to force streaming on the ChatModel */\n  forceStreaming?: boolean;\n}\n\n/**\n * Middleware for handling streaming tool calls in a ChatModel.\n *\n * This middleware observes and listens to ChatModel stream updates and parses\n * the tool calls on demand so that they can be consumed as soon as possible.\n *\n * @example\n * ```typescript\n * const middleware = new StreamToolCallMiddleware({\n *   target: thinkTool,\n *   key: \"thoughts\",\n *   matchNested: false,\n *   forceStreaming: true,\n * });\n *\n * middleware.emitter.on(\"update\", (event) => {\n *   console.log(\"Delta:\", event.delta);\n *   console.log(\"Structured:\", event.outputStructured);\n * });\n *\n * await llm.run(messages, { tools: [thinkTool] }).middleware(middleware);\n * ```\n */\nexport class StreamToolCallMiddleware<T = any> extends Middleware<RunInstance> {\n  private readonly target: Tool<any>;\n  private readonly key: string;\n  private readonly matchNested: boolean;\n  private readonly forceStreaming: boolean;\n  private readonly cleanups: (() => void)[] = [];\n\n  private output = new ChatModelOutput([]);\n  private buffer = \"\";\n  private delta = \"\";\n\n  public readonly emitter: Emitter<StreamToolCallMiddlewareCallbacks<T>>;\n\n  constructor(options: StreamToolCallMiddlewareOptions) {\n    super();\n\n    this.target = options.target;\n    this.key = options.key;\n    this.matchNested = options.matchNested ?? false;\n    this.forceStreaming = options.forceStreaming ?? false;\n\n    this.emitter = Emitter.root.child<StreamToolCallMiddlewareCallbacks<T>>({\n      namespace: [\"middleware\", \"streamToolCall\"],\n    });\n  }\n\n  reset() {\n    this.output = new ChatModelOutput([]);\n    this.buffer = \"\";\n    this.delta = \"\";\n  }\n\n  async add(chunk: ChatModelOutput) {\n    await this.handleNewToken(\n      {\n        value: chunk,\n        callbacks: {\n          abort: () => {},\n        },\n      },\n      {} as EventMeta,\n    );\n  }\n\n  isEmpty() {\n    return this.buffer.length === 0;\n  }\n\n  bind(ctx: RunContext<RunInstance>): void {\n    this.reset();\n\n    // Listen to ChatModel start event\n    this.cleanups.push(\n      ctx.instance.emitter.match(\n        (meta) => meta.creator instanceof ChatModel && meta.name === \"start\",\n        this.handleStart.bind(this),\n        this.createEmitterOptions(),\n      ),\n    );\n\n    // Listen to ChatModel newToken event\n    this.cleanups.push(\n      ctx.instance.emitter.match(\n        (meta) => meta.creator instanceof ChatModel && meta.name === \"newToken\",\n        this.handleNewToken.bind(this),\n        this.createEmitterOptions(),\n      ),\n    );\n\n    // Listen to ChatModel success event\n    this.cleanups.push(\n      ctx.instance.emitter.match(\n        (meta) => meta.creator instanceof ChatModel && meta.name === \"success\",\n        this.handleSuccess.bind(this),\n        this.createEmitterOptions(),\n      ),\n    );\n  }\n\n  unbind(): void {\n    // Clean up previous bindings\n    while (this.cleanups.length > 0) {\n      const fn = this.cleanups.shift()!;\n      fn();\n    }\n  }\n\n  protected createEmitterOptions(): EmitterOptions {\n    return {\n      matchNested: this.matchNested,\n      isBlocking: true,\n    };\n  }\n\n  private async process(toolName: string, args: any): Promise<void> {\n    if (toolName !== this.target.name) {\n      return;\n    }\n\n    const parsedArgs = isString(args) ? parseBrokenJson(args, { pair: [\"{\", \"}\"] }) : args;\n    const outputStructured = (await this.target.parse(parsedArgs).catch(() => null)) as T;\n    if (!outputStructured) {\n      return;\n    }\n\n    let output = \"\";\n    if (hasProp(outputStructured, this.key as keyof T)) {\n      output = (outputStructured as any)[this.key] || \"\";\n      if (!isString(output)) {\n        output = JSON.stringify(output);\n      }\n      this.delta = output.slice(this.buffer.length);\n      this.buffer = output;\n\n      if (!this.delta) {\n        return;\n      }\n    }\n\n    await this.emitter.emit(\"update\", {\n      outputStructured,\n      delta: this.delta,\n      output,\n    });\n  }\n\n  private async handleStart(\n    data: InferCallbackValue<ChatModelEvents[\"start\"]>,\n    _meta: EventMeta,\n  ): Promise<void> {\n    if (this.forceStreaming) {\n      data.input.stream = true;\n      data.input.streamPartialToolCalls = true;\n    }\n  }\n\n  private async handleSuccess(\n    data: InferCallbackValue<ChatModelEvents[\"success\"]>,\n    _meta: EventMeta,\n  ): Promise<void> {\n    // If we haven't received any tokens yet, process the final output\n    if (this.output.messages.length === 0) {\n      await this.handleNewToken({ value: data.value, callbacks: { abort: () => {} } }, _meta);\n    }\n  }\n\n  private async handleNewToken(\n    data: InferCallbackValue<ChatModelEvents[\"newToken\"]>,\n    _meta: EventMeta,\n  ): Promise<void> {\n    await this.output.merge(data.value);\n\n    const toolCalls = this.output.getToolCalls();\n    if (toolCalls.length > 0) {\n      for (const toolCall of toolCalls) {\n        await this.process(toolCall.toolName, toolCall.input);\n      }\n      return;\n    }\n\n    // Try to parse text content as a tool call\n    const textContent = this.output.getTextContent();\n    const toolCall = parseBrokenJson(textContent, { pair: [\"{\", \"}\"] });\n\n    if (!toolCall || !isPlainObject(toolCall)) {\n      return;\n    }\n\n    await this.process(toolCall.name as string, toolCall.parameters as any);\n  }\n}\n"]}