{"version":3,"file":"stream_events.cjs","names":[],"sources":["../../src/utils/stream_events.ts"],"sourcesContent":["/**\n * Converts Gemini-style stream responses into LangChain ChatModelStreamEvents.\n *\n * @module\n */\n\nimport { finalizeContentBlock } from \"@langchain/core/language_models/compat\";\nimport type {\n  ChatModelStreamEvent,\n  FinishReason,\n} from \"@langchain/core/language_models/event\";\nimport type { ContentBlock, UsageMetadata } from \"@langchain/core/messages\";\n\n// oxlint-disable-next-line @typescript-eslint/no-explicit-any\nexport type GeminiStreamResponse = Record<string, any>;\n\nexport interface ConvertGoogleGeminiStreamOptions {\n  streamUsage?: boolean;\n}\n\ntype BlockKey = \"text\" | \"reasoning\" | `tool:${number}`;\n\nexport async function* convertGoogleGeminiStream(\n  source: AsyncIterable<GeminiStreamResponse>,\n  options: ConvertGoogleGeminiStreamOptions = {}\n): AsyncGenerator<ChatModelStreamEvent> {\n  const shouldStreamUsage = options.streamUsage ?? true;\n  const blockAccumulators = new Map<\n    number,\n    // oxlint-disable-next-line @typescript-eslint/no-explicit-any\n    Record<string, any>\n  >();\n  const blockKeyToIndex = new Map<BlockKey, number>();\n  let nextBlockIndex = 0;\n  let messageStarted = false;\n  let usageSnapshot: UsageMetadata | undefined;\n  let finishReason: FinishReason = \"stop\";\n\n  const getOrCreateBlockIndex = (\n    key: BlockKey,\n    initial: Record<string, unknown>\n  ): { index: number; isNew: boolean } => {\n    const existing = blockKeyToIndex.get(key);\n    if (existing !== undefined) {\n      return { index: existing, isNew: false };\n    }\n    const index = nextBlockIndex++;\n    blockKeyToIndex.set(key, index);\n    blockAccumulators.set(index, { ...initial });\n    return { index, isNew: true };\n  };\n\n  for await (const response of source) {\n    if (!messageStarted) {\n      messageStarted = true;\n      yield { event: \"message-start\" as const };\n    }\n\n    const usageMetadata = response.usageMetadata ?? response.usage_metadata;\n    if (shouldStreamUsage && usageMetadata) {\n      const input = usageMetadata.promptTokenCount ?? 0;\n      const output = usageMetadata.candidatesTokenCount ?? 0;\n      usageSnapshot = {\n        input_tokens: input,\n        output_tokens: output,\n        total_tokens: usageMetadata.totalTokenCount ?? input + output,\n      };\n      yield { event: \"usage\" as const, usage: usageSnapshot };\n    }\n\n    const candidate = response.candidates?.[0];\n    if (candidate?.finishReason) {\n      finishReason = mapGeminiFinishReason(candidate.finishReason);\n    }\n\n    const parts = candidate?.content?.parts;\n    if (!parts) continue;\n\n    let toolIdx = 0;\n    for (const part of parts) {\n      if (part.text) {\n        if (part.thought) {\n          const key: BlockKey = \"reasoning\";\n          const { index, isNew } = getOrCreateBlockIndex(key, {\n            type: \"reasoning\",\n            reasoning: \"\",\n          });\n          if (isNew) {\n            yield {\n              event: \"content-block-start\" as const,\n              index,\n              content: { type: \"reasoning\", reasoning: \"\" } as ContentBlock,\n            };\n          }\n          const acc = blockAccumulators.get(index)!;\n          acc.reasoning = (acc.reasoning ?? \"\") + part.text;\n          yield {\n            event: \"content-block-delta\" as const,\n            index,\n            delta: { type: \"reasoning-delta\" as const, reasoning: part.text },\n          };\n        } else {\n          const key: BlockKey = \"text\";\n          const { index, isNew } = getOrCreateBlockIndex(key, {\n            type: \"text\",\n            text: \"\",\n          });\n          if (isNew) {\n            yield {\n              event: \"content-block-start\" as const,\n              index,\n              content: { type: \"text\", text: \"\" } as ContentBlock,\n            };\n          }\n          const acc = blockAccumulators.get(index)!;\n          acc.text = (acc.text ?? \"\") + part.text;\n          yield {\n            event: \"content-block-delta\" as const,\n            index,\n            delta: { type: \"text-delta\" as const, text: part.text },\n          };\n        }\n      } else if (part.functionCall) {\n        const key: BlockKey = `tool:${toolIdx}`;\n        const args = JSON.stringify(part.functionCall.args ?? {});\n        const { index, isNew } = getOrCreateBlockIndex(key, {\n          type: \"tool_call_chunk\",\n          name: part.functionCall.name,\n          args: \"\",\n          index: toolIdx,\n        });\n        if (isNew) {\n          yield {\n            event: \"content-block-start\" as const,\n            index,\n            content: {\n              type: \"tool_call_chunk\",\n              name: part.functionCall.name,\n              args: \"\",\n              index: toolIdx,\n            } as ContentBlock,\n          };\n        }\n        const acc = blockAccumulators.get(index)!;\n        acc.args = args;\n        yield {\n          event: \"content-block-delta\" as const,\n          index,\n          delta: {\n            type: \"block-delta\" as const,\n            fields: {\n              type: \"tool_call_chunk\",\n              name: acc.name,\n              args: acc.args,\n            },\n          },\n        };\n        toolIdx += 1;\n      }\n    }\n  }\n\n  for (const [index, acc] of blockAccumulators) {\n    yield {\n      event: \"content-block-finish\" as const,\n      index,\n      content: finalizeContentBlock(acc as ContentBlock),\n    };\n  }\n\n  yield {\n    event: \"message-finish\" as const,\n    reason: finishReason,\n    ...(usageSnapshot ? { usage: usageSnapshot } : {}),\n    responseMetadata: { model_provider: \"google\" },\n  };\n}\n\nfunction mapGeminiFinishReason(reason: string): FinishReason {\n  switch (reason.toLowerCase()) {\n    case \"max_tokens\":\n    case \"max-token\":\n    case \"max_token\":\n      return \"length\";\n    case \"safety\":\n    case \"recitation\":\n    case \"language\":\n    case \"blocklist\":\n    case \"prohibited_content\":\n    case \"prohibited-content\":\n    case \"spii\":\n    case \"image_safety\":\n    case \"image-safety\":\n    case \"image_prohibited_content\":\n    case \"image-prohibited-content\":\n    case \"image_recitation\":\n    case \"image-recitation\":\n      return \"content_filter\";\n    default:\n      return \"stop\";\n  }\n}\n"],"mappings":";;;;;;;AAsBA,gBAAuB,0BACrB,QACA,UAA4C,EAAE,EACR;CACtC,MAAM,oBAAoB,QAAQ,eAAe;CACjD,MAAM,oCAAoB,IAAI,KAI3B;CACH,MAAM,kCAAkB,IAAI,KAAuB;CACnD,IAAI,iBAAiB;CACrB,IAAI,iBAAiB;CACrB,IAAI;CACJ,IAAI,eAA6B;CAEjC,MAAM,yBACJ,KACA,YACsC;EACtC,MAAM,WAAW,gBAAgB,IAAI,IAAI;AACzC,MAAI,aAAa,KAAA,EACf,QAAO;GAAE,OAAO;GAAU,OAAO;GAAO;EAE1C,MAAM,QAAQ;AACd,kBAAgB,IAAI,KAAK,MAAM;AAC/B,oBAAkB,IAAI,OAAO,EAAE,GAAG,SAAS,CAAC;AAC5C,SAAO;GAAE;GAAO,OAAO;GAAM;;AAG/B,YAAW,MAAM,YAAY,QAAQ;AACnC,MAAI,CAAC,gBAAgB;AACnB,oBAAiB;AACjB,SAAM,EAAE,OAAO,iBAA0B;;EAG3C,MAAM,gBAAgB,SAAS,iBAAiB,SAAS;AACzD,MAAI,qBAAqB,eAAe;GACtC,MAAM,QAAQ,cAAc,oBAAoB;GAChD,MAAM,SAAS,cAAc,wBAAwB;AACrD,mBAAgB;IACd,cAAc;IACd,eAAe;IACf,cAAc,cAAc,mBAAmB,QAAQ;IACxD;AACD,SAAM;IAAE,OAAO;IAAkB,OAAO;IAAe;;EAGzD,MAAM,YAAY,SAAS,aAAa;AACxC,MAAI,WAAW,aACb,gBAAe,sBAAsB,UAAU,aAAa;EAG9D,MAAM,QAAQ,WAAW,SAAS;AAClC,MAAI,CAAC,MAAO;EAEZ,IAAI,UAAU;AACd,OAAK,MAAM,QAAQ,MACjB,KAAI,KAAK,KACP,KAAI,KAAK,SAAS;GAEhB,MAAM,EAAE,OAAO,UAAU,sBAAsB,aAAK;IAClD,MAAM;IACN,WAAW;IACZ,CAAC;AACF,OAAI,MACF,OAAM;IACJ,OAAO;IACP;IACA,SAAS;KAAE,MAAM;KAAa,WAAW;KAAI;IAC9C;GAEH,MAAM,MAAM,kBAAkB,IAAI,MAAM;AACxC,OAAI,aAAa,IAAI,aAAa,MAAM,KAAK;AAC7C,SAAM;IACJ,OAAO;IACP;IACA,OAAO;KAAE,MAAM;KAA4B,WAAW,KAAK;KAAM;IAClE;SACI;GAEL,MAAM,EAAE,OAAO,UAAU,sBAAsB,QAAK;IAClD,MAAM;IACN,MAAM;IACP,CAAC;AACF,OAAI,MACF,OAAM;IACJ,OAAO;IACP;IACA,SAAS;KAAE,MAAM;KAAQ,MAAM;KAAI;IACpC;GAEH,MAAM,MAAM,kBAAkB,IAAI,MAAM;AACxC,OAAI,QAAQ,IAAI,QAAQ,MAAM,KAAK;AACnC,SAAM;IACJ,OAAO;IACP;IACA,OAAO;KAAE,MAAM;KAAuB,MAAM,KAAK;KAAM;IACxD;;WAEM,KAAK,cAAc;GAC5B,MAAM,MAAgB,QAAQ;GAC9B,MAAM,OAAO,KAAK,UAAU,KAAK,aAAa,QAAQ,EAAE,CAAC;GACzD,MAAM,EAAE,OAAO,UAAU,sBAAsB,KAAK;IAClD,MAAM;IACN,MAAM,KAAK,aAAa;IACxB,MAAM;IACN,OAAO;IACR,CAAC;AACF,OAAI,MACF,OAAM;IACJ,OAAO;IACP;IACA,SAAS;KACP,MAAM;KACN,MAAM,KAAK,aAAa;KACxB,MAAM;KACN,OAAO;KACR;IACF;GAEH,MAAM,MAAM,kBAAkB,IAAI,MAAM;AACxC,OAAI,OAAO;AACX,SAAM;IACJ,OAAO;IACP;IACA,OAAO;KACL,MAAM;KACN,QAAQ;MACN,MAAM;MACN,MAAM,IAAI;MACV,MAAM,IAAI;MACX;KACF;IACF;AACD,cAAW;;;AAKjB,MAAK,MAAM,CAAC,OAAO,QAAQ,kBACzB,OAAM;EACJ,OAAO;EACP;EACA,UAAA,GAAA,uCAAA,sBAA8B,IAAoB;EACnD;AAGH,OAAM;EACJ,OAAO;EACP,QAAQ;EACR,GAAI,gBAAgB,EAAE,OAAO,eAAe,GAAG,EAAE;EACjD,kBAAkB,EAAE,gBAAgB,UAAU;EAC/C;;AAGH,SAAS,sBAAsB,QAA8B;AAC3D,SAAQ,OAAO,aAAa,EAA5B;EACE,KAAK;EACL,KAAK;EACL,KAAK,YACH,QAAO;EACT,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK;EACL,KAAK,mBACH,QAAO;EACT,QACE,QAAO"}