{"version":3,"sources":["/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-P2AFKIPT.cjs","../../src/services/intents/fulfill.service.ts"],"names":[],"mappings":"AAAA;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACA;ACtBA,gCAA2B;AAsBpB,IAAM,qBAAA,EAAN,MAA2B;AAAA,EACzB;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EAER,WAAA,CACC,WAAA,EACA,YAAA,EACA,kBAAA,EACA,gBAAA,EACA,UAAA,EACA,wBAAA,EACC;AACD,IAAA,IAAA,CAAK,YAAA,EAAc,WAAA;AACnB,IAAA,IAAA,CAAK,aAAA,EAAe,YAAA;AACpB,IAAA,IAAA,CAAK,mBAAA,EAAqB,kBAAA;AAC1B,IAAA,IAAA,CAAK,iBAAA,EAAmB,gBAAA;AACxB,IAAA,IAAA,CAAK,iBAAA,EAAmB,IAAI,uCAAA,CAAiB,WAAA,EAAa,YAAY,CAAA;AACtE,IAAA,IAAA,CAAK,WAAA,EAAa,UAAA;AAClB,IAAA,IAAA,CAAK,yBAAA,EAA2B,wBAAA;AAAA,EACjC;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAkBA,MAAA,CAAe,gBAAA,CACd,KAAA,EACA,MAAA,EACA,MAAA,EAC8B;AAC9B,IAAA,MAAM,OAAA,EAAS,MAAM,+CAAA,IAAc,CAAK,YAAA,EAAc,MAAM,CAAA;AAE5D,IAAA,MAAM,cAAA,EAAgB,IAAA,CAAK,WAAA,CAAY,QAAA,CAAS,CAAA;AAChD,IAAA,MAAM,SAAA,EAAW,aAAA,CAAc,gBAAA,CAAiB;AAAA,MAC/C,KAAA;AAAA,MACA,MAAA;AAAA,MACA,YAAA,EAAc,MAAA,CAAO,IAAA,CAAK;AAAA,IAC3B,CAAC,CAAA;AAED,IAAA,yBAAA,CAAQ,MAAA,CAAO,KAAA,CAAM,0BAAA,EAA4B;AAAA,MAChD,QAAA,EAAU,MAAA,CAAO,QAAA;AAAA,MACjB;AAAA,IACD,CAAC,CAAA;AAED,IAAA,MAAM,WAAA,EAAa,MAAM,mDAAA,IAAiB,CAAK,YAAY,CAAA;AAC3D,IAAA,MAAM,MAAA,EAAQ,MAAM,IAAA,CAAK,kBAAA,CAAmB,QAAA,CAAS;AAAA,MACpD,UAAA;AAAA,MACA,IAAA,EAAM;AAAA,IACP,CAAC,CAAA;AACD,IAAA,MAAM,OAAA,EAAS,IAAA,CAAK,kBAAA,CAAmB,GAAA,CAAI;AAAA,MAC1C,QAAA;AAAA,MACA,KAAA;AAAA,MACA,KAAA;AAAA,MACA,MAAA;AAAA,MACA,UAAA,kBAAY,MAAA,2BAAQ,aAAA,IAAe,WAAA,EAAa,WAAA,EAAa;AAAA,IAC9D,CAAC,CAAA;AAED,IAAA,IAAI,OAAA,EAAS,MAAM,MAAA,CAAO,IAAA,CAAK,CAAA;AAC/B,IAAA,MAAA,CAAO,CAAC,MAAA,CAAO,IAAA,EAAM;AACpB,MAAA,MAAM,MAAA,CAAO,KAAA;AACb,MAAA,OAAA,EAAS,MAAM,MAAA,CAAO,IAAA,CAAK,CAAA;AAAA,IAC5B;AAEA,IAAA,yBAAA,CAAQ,MAAA,CAAO,KAAA,CAAM,8BAAA,EAAgC;AAAA,MACpD,QAAA,EAAU,MAAA,CAAO,QAAA;AAAA,MACjB,iBAAA,EAAmB,MAAA,CAAO,KAAA,CAAM,iBAAA;AAAA,MAChC,UAAA,kBAAY,MAAA,6BAAQ;AAAA,IACrB,CAAC,CAAA;AAAA,EACF;AAAA;AAAA;AAAA;AAAA;AAAA,EAMQ,eAAA,CACP,eAAA,EACA,MAAA,EAC0C;AAC1C,IAAA,MAAM,EAAE,SAAA,EAAW,EAAA,EAAI,OAAO,EAAA,EAAI,eAAA;AAElC,IAAA,GAAA,CAAI,CAAC,OAAA,GAAU,IAAA,CAAK,gBAAA,EAAkB;AACrC,MAAA,yBAAA,CAAQ,MAAA,CAAO,IAAA,CAAK,6CAA6C,CAAA;AACjE,MAAA,MAAM,eAAA,EAAiB,IAAA,CAAK,gBAAA,CAAiB;AAAA,QAC5C,eAAA;AAAA,QACA;AAAA,MACD,CAAC,CAAA;AACD,MAAA,GAAA,CAAI,eAAA,IAAmB,KAAA,CAAA,EAAW;AACjC,QAAA,OAAO,cAAA;AAAA,MACR;AAAA,IACD;AAEA,IAAA,GAAA,iBAAI,MAAA,6BAAQ,aAAA,GAAc,IAAA,CAAK,wBAAA,EAA0B;AACxD,MAAA,OAAO,IAAA,CAAK,wBAAA,CAAyB,eAAA,EAAiB,MAAM,CAAA;AAAA,IAC7D;AAEA,IAAA,OAAO,IAAA,CAAK,gBAAA,CAAiB,QAAA,EAAU,MAAA,EAAQ,MAAM,CAAA;AAAA,EACtD;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EASA,MAAA,CAAe,wBAAA,CACd,eAAA,EACA,MAAA,EAC8B;AAC9B,IAAA,MAAM,EAAE,SAAA,EAAW,EAAA,EAAI,OAAO,EAAA,EAAI,eAAA;AAClC,IAAA,MAAM,WAAA,kBAAa,MAAA,6BAAQ,YAAA;AAC3B,IAAA,GAAA,CAAI,CAAC,WAAA,GAAc,CAAC,IAAA,CAAK,wBAAA,EAA0B;AAClD,MAAA,KAAA,EAAO,IAAA,CAAK,gBAAA,CAAiB,QAAA,EAAU,MAAA,EAAQ,MAAM,CAAA;AACrD,MAAA,MAAA;AAAA,IACD;AAEA,IAAA,yBAAA,CAAQ,MAAA,CAAO,IAAA,CAAK,uCAAA,EAAyC;AAAA,MAC5D,UAAA,kBAAY,MAAA,6BAAQ,MAAA;AAAA,MACpB;AAAA,IACD,CAAC,CAAA;AAED,IAAA,IAAI,QAAA,EAAU,KAAA;AACd,IAAA,IAAI;AACH,MAAA,MAAM,OAAA,EAAS,IAAA,CAAK,wBAAA,CAAyB,2BAAA;AAAA,QAC5C,UAAA;AAAA,QACA,MAAA;AAAA,QACA;AAAA,MACD,CAAA;AACA,MAAA,IAAA,MAAA,CAAA,MAAiB,MAAA,GAAS,MAAA,EAAQ;AACjC,QAAA,QAAA,EAAU,IAAA;AACV,QAAA,MAAM,KAAA;AAAA,MACP;AACA,MAAA,MAAA;AAAA,IACD,EAAA,MAAA,CAAS,KAAA,EAAO;AACf,MAAA,GAAA,CAAI,OAAA,EAAS;AACZ,QAAA,MAAM,KAAA;AAAA,MACP;AACA,MAAA,yBAAA,CAAQ,MAAA,CAAO,IAAA;AAAA,QACd,iEAAA;AAAA,QACA,EAAE,UAAA,kBAAY,MAAA,6BAAQ,MAAA,EAAM,UAAA,EAAY,MAAM;AAAA,MAC/C,CAAA;AAAA,IACD;AAEA,IAAA,KAAA,EAAO,IAAA,CAAK,gBAAA,CAAiB,QAAA,EAAU,MAAA,EAAQ,MAAM,CAAA;AAAA,EACtD;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAiBA,MAAA,CAAc,aAAA,CACb,OAAA,EACA,MAAA,EACA,aAAA,EACA,gBAAA,EAC8B;AAC9B,IAAA,MAAM,gBAAA,EAAkB,IAAA,CAAK,GAAA,CAAI,CAAA;AACjC,IAAA,yBAAA,CAAQ,YAAA,CAAa,IAAA,CAAK,wBAAA,EAA0B;AAAA,MACnD,QAAA,EAAU,MAAA,CAAO,QAAA;AAAA,MACjB,WAAA,EAAa,OAAA,CAAQ,MAAA;AAAA,MACrB,gBAAA;AAAA,MACA,SAAA,EAAW,IAAI,IAAA,CAAK,eAAe,CAAA,CAAE,WAAA,CAAY;AAAA,IAClD,CAAC,CAAA;AAED,IAAA,IAAI,kBAAA,EAAoB,EAAA;AACxB,IAAA,IAAI,cAAA;AAEJ,IAAA,GAAA,CAAI,OAAA,CAAQ,OAAA,GAAU,CAAA,EAAG;AAExB,MAAA,MAAM,gBAAA,EAAkB,OAAA,CAAQ,CAAC,CAAA;AACjC,MAAA,GAAA,CAAI,CAAC,eAAA,EAAiB;AACrB,QAAA,MAAA;AAAA,MACD;AAEA,MAAA,MAAM,EAAE,SAAA,EAAW,EAAA,EAAI,MAAA,EAAQ,WAAW,EAAA,EAAI,eAAA;AAC9C,MAAA,yBAAA,CAAQ,MAAA,CAAO,IAAA;AAAA,QACd,CAAA,uBAAA,EAA0B,QAAQ,CAAA,EAAA,kBAAK,MAAA,6BAAQ,MAAI,CAAA;AAAA,MAAA;AAEpD,MAAA;AAEA,MAAA;AACA,MAAA;AACC,QAAA;AAAwD,MAAA;AAIzD,MAAA;AACA,MAAA;AACC,QAAA;AAAA,MAAA;AAID,MAAA;AACC,QAAA;AACC,UAAA;AAAgC,QAAA;AAEhC,UAAA;AAA4B,QAAA;AAE7B,QAAA;AAAM,MAAA;AACP,IAAA;AAGA,MAAA;AACC,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AAEA,QAAA;AAEA,QAAA;AACA,QAAA;AACC,UAAA;AAAwD,QAAA;AAIzD,QAAA;AACA,QAAA;AACC,UAAA;AAAA,QAAA;AAGD,QAAA;AAEC,UAAA;AACC,YAAA;AACC,cAAA;AAAgC,YAAA;AAEhC,cAAA;AAA4B,YAAA;AAE7B,YAAA;AAAM,UAAA;AACP,QAAA;AAGA,UAAA;AACA,UAAA;AACC,YAAA;AACC,cAAA;AAA2B,YAAA;AAE3B,cAAA;AAA4B,YAAA;AAG5B,cAAA;AAAM,YAAA;AACP,UAAA;AAGD,UAAA;AAAqB,YAAA;AACE,YAAA;AACtB,YAAA;AACoB,YAAA;AAC2B,YAAA;AAClB,UAAA;AAC7B,QAAA;AACF,MAAA;AACD,IAAA;AAGA,MAAA;AAEA,MAAA;AACC,QAAA;AACA,QAAA;AACA,QAAA;AACA,QAAA;AAGA,QAAA;AACC,UAAA;AACA,UAAA;AAAqB,YAAA;AACE,YAAA;AACtB,YAAA;AACoB,YAAA;AACkC,YAAA;AACzB,UAAA;AAC7B,QAAA;AAGF,QAAA;AACA,QAAA;AACC,UAAA;AAAwD,QAAA;AAIzD,QAAA;AACA,QAAA;AACC,UAAA;AAAA,QAAA;AAID,QAAA;AACA,QAAA;AACC,UAAA;AACC,YAAA;AAA2B,UAAA;AAE3B,YAAA;AAA4B,UAAA;AAG5B,YAAA;AAAM,UAAA;AACP,QAAA;AAGD,QAAA;AAAwB,UAAA;AACvB,UAAA;AACA,UAAA;AACA,UAAA;AACU,QAAA;AACV,MAAA;AAIF,MAAA;AAA8C,QAAA;AAC7C,QAAA;AACA,MAAA;AAGD,MAAA;AACC,QAAA;AACC,UAAA;AAAgC,QAAA;AAEjC,QAAA;AAAM,MAAA;AACP,IAAA;AAID,IAAA;AACC,MAAA;AAAsE,IAAA;AAIvE,IAAA;AACC,MAAA;AAAM,QAAA;AACA,QAAA;AACL,QAAA;AAAA,QAAA;AAEA,QAAA;AACsC,MAAA;AACvC,IAAA;AAEA,MAAA;AAAkE,IAAA;AAGnE,IAAA;AACA,IAAA;AAEA,IAAA;AAAsD,MAAA;AACpC,MAAA;AACU,MAAA;AACkB,IAAA;AAC7C,EAAA;AACF,EAAA;AAKC,IAAA;AACA,IAAA;AACC,MAAA;AAAO,IAAA;AAER,IAAA;AAA4B,MAAA;AAC+B,MAAA;AAC/B,IAAA;AAC3B,EAAA;AAEH;AD7EA;AACA;AACA;AACA","file":"/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-P2AFKIPT.cjs","sourcesContent":[null,"import { randomUUID } from \"node:crypto\";\nimport { getManifest } from \"@/config/manifest\";\nimport type { MemoryModule, ModelModule } from \"@/modules\";\nimport type { OnIntentFallback } from \"@/types/agent\";\nimport {\n\ttype FulfillmentResult,\n\ttype Intent,\n\tMessageRole,\n\ttype ThreadObject,\n\ttype TriggeredIntent,\n} from \"@/types/memory\";\nimport type { StreamEvent } from \"@/types/stream\";\nimport { loggers } from \"@/utils/logger\";\nimport { appendTextMessageToThread } from \"@/utils/thread-messages\";\nimport { sanitizeThinkingData } from \"@/utils/tool-args\";\nimport { PIIFilterMode, type PIIService } from \"../pii.service\";\nimport fulfillPrompt from \"../prompts/fulfill\";\nimport toolSelectPrompt from \"../prompts/tool-select\";\nimport type { ToolCallingService } from \"../tool-calling.service\";\nimport type { WorkflowExecutionService } from \"../workflow-execution.service\";\nimport { AggregateService } from \"./aggregate.service\";\n\nexport class IntentFulfillService {\n\tprivate modelModule: ModelModule;\n\tprivate memoryModule: MemoryModule;\n\tprivate onIntentFallback?: OnIntentFallback;\n\tprivate aggregateService: AggregateService;\n\tprivate piiService?: PIIService;\n\tprivate toolCallingService: ToolCallingService;\n\tprivate workflowExecutionService?: WorkflowExecutionService;\n\n\tconstructor(\n\t\tmodelModule: ModelModule,\n\t\tmemoryModule: MemoryModule,\n\t\ttoolCallingService: ToolCallingService,\n\t\tonIntentFallback?: OnIntentFallback,\n\t\tpiiService?: PIIService,\n\t\tworkflowExecutionService?: WorkflowExecutionService,\n\t) {\n\t\tthis.modelModule = modelModule;\n\t\tthis.memoryModule = memoryModule;\n\t\tthis.toolCallingService = toolCallingService;\n\t\tthis.onIntentFallback = onIntentFallback;\n\t\tthis.aggregateService = new AggregateService(modelModule, memoryModule);\n\t\tthis.piiService = piiService;\n\t\tthis.workflowExecutionService = workflowExecutionService;\n\t}\n\n\t/**\n\t * Fulfills the detected intent by generating a streaming response.\n\t *\n\t * Manages the complete inference loop including:\n\t * - Loading prompts and conversation history\n\t * - Collecting available tools from modules\n\t * - Executing model inference with tool support\n\t * - Processing tool calls iteratively until completion\n\t * - Streaming results as Server-Sent Events\n\t *\n\t * @param query - The user's input query\n\t * @param threadId - Thread identifier for context\n\t * @param thread - Previous conversation history\n\t * @param intent - Optional detected intent with custom prompt\n\t * @returns AsyncGenerator yielding StreamEvent objects\n\t */\n\tprivate async *intentFulfilling(\n\t\tquery: string,\n\t\tthread: ThreadObject,\n\t\tintent?: Intent,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst prompt = await fulfillPrompt(this.memoryModule, intent);\n\n\t\tconst modelInstance = this.modelModule.getModel();\n\t\tconst messages = modelInstance.generateMessages({\n\t\t\tquery,\n\t\t\tthread,\n\t\t\tsystemPrompt: prompt.trim(),\n\t\t});\n\n\t\tloggers.intent.debug(\"Intent fulfillment start\", {\n\t\t\tthreadId: thread.threadId,\n\t\t\tmessages,\n\t\t});\n\n\t\tconst toolPrompt = await toolSelectPrompt(this.memoryModule);\n\t\tconst tools = await this.toolCallingService.getTools({\n\t\t\ttoolPrompt,\n\t\t\tmode: \"all\",\n\t\t});\n\t\tconst stream = this.toolCallingService.run({\n\t\t\tmessages,\n\t\t\ttools,\n\t\t\tquery,\n\t\t\tthread,\n\t\t\ttoolChoice: intent?.toolChoice === \"required\" ? \"required\" : \"auto\",\n\t\t});\n\n\t\tlet result = await stream.next();\n\t\twhile (!result.done) {\n\t\t\tyield result.value;\n\t\t\tresult = await stream.next();\n\t\t}\n\n\t\tloggers.intent.debug(\"Intent fulfillment completed\", {\n\t\t\tthreadId: thread.threadId,\n\t\t\ttoolCallsExecuted: result.value.toolCallsExecuted,\n\t\t\tintentName: intent?.name,\n\t\t});\n\t}\n\n\t/**\n\t * Returns the appropriate stream for a triggered intent.\n\t * Uses fallback handler if no intent matched and fallback is configured.\n\t */\n\tprivate getIntentStream(\n\t\ttriggeredIntent: TriggeredIntent,\n\t\tthread: ThreadObject,\n\t): AsyncGenerator<StreamEvent> | undefined {\n\t\tconst { subquery = \"\", intent } = triggeredIntent;\n\n\t\tif (!intent && this.onIntentFallback) {\n\t\t\tloggers.intent.info(\"No intent matched, calling fallback handler\");\n\t\t\tconst fallbackStream = this.onIntentFallback({\n\t\t\t\ttriggeredIntent,\n\t\t\t\tthread,\n\t\t\t});\n\t\t\tif (fallbackStream !== undefined) {\n\t\t\t\treturn fallbackStream;\n\t\t\t}\n\t\t}\n\n\t\tif (intent?.workflowId && this.workflowExecutionService) {\n\t\t\treturn this.intentWorkflowFulfilling(triggeredIntent, thread);\n\t\t}\n\n\t\treturn this.intentFulfilling(subquery, thread, intent);\n\t}\n\n\t/**\n\t * Fulfills an intent by running its mapped workflow. If the workflow\n\t * fails before producing any event (unknown id, missing definition),\n\t * falls back to the prompt-based fulfillment path so a stale mapping\n\t * never disables the intent. Failures after streaming started are\n\t * rethrown to avoid duplicating a partial response.\n\t */\n\tprivate async *intentWorkflowFulfilling(\n\t\ttriggeredIntent: TriggeredIntent,\n\t\tthread: ThreadObject,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst { subquery = \"\", intent } = triggeredIntent;\n\t\tconst workflowId = intent?.workflowId;\n\t\tif (!workflowId || !this.workflowExecutionService) {\n\t\t\tyield* this.intentFulfilling(subquery, thread, intent);\n\t\t\treturn;\n\t\t}\n\n\t\tloggers.intent.info(\"Fulfilling intent via mapped workflow\", {\n\t\t\tintentName: intent?.name,\n\t\t\tworkflowId,\n\t\t});\n\n\t\tlet yielded = false;\n\t\ttry {\n\t\t\tconst stream = this.workflowExecutionService.executeIntentWorkflowStream(\n\t\t\t\tworkflowId,\n\t\t\t\tthread,\n\t\t\t\tsubquery,\n\t\t\t);\n\t\t\tfor await (const event of stream) {\n\t\t\t\tyielded = true;\n\t\t\t\tyield event;\n\t\t\t}\n\t\t\treturn;\n\t\t} catch (error) {\n\t\t\tif (yielded) {\n\t\t\t\tthrow error;\n\t\t\t}\n\t\t\tloggers.intent.warn(\n\t\t\t\t\"Intent workflow unavailable, falling back to prompt fulfillment\",\n\t\t\t\t{ intentName: intent?.name, workflowId, error },\n\t\t\t);\n\t\t}\n\n\t\tyield* this.intentFulfilling(subquery, thread, intent);\n\t}\n\n\t/**\n\t * Processes all triggered intents and generates a unified response.\n\t *\n\t * Workflow:\n\t * 1. Process each intent sequentially, collecting results\n\t * 2. Yield thinking_process events for progress visibility\n\t * 3. Use AggregateService to unify results if needsAggregation is true\n\t * 4. Stream the final (possibly aggregated) response\n\t *\n\t * @param intents - Array of triggered intents to process\n\t * @param thread - The thread history\n\t * @param originalQuery - The user's original query (for aggregate context)\n\t * @param needsAggregation - Whether the results need to be aggregated\n\t * @returns AsyncGenerator yielding StreamEvent objects\n\t */\n\tpublic async *intentFulfill(\n\t\tintents: Array<TriggeredIntent>,\n\t\tthread: ThreadObject,\n\t\toriginalQuery: string,\n\t\tneedsAggregation: boolean,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst streamStartTime = Date.now();\n\t\tloggers.intentStream.info(\"Stream session started\", {\n\t\t\tthreadId: thread.threadId,\n\t\t\tintentCount: intents.length,\n\t\t\tneedsAggregation,\n\t\t\tstartTime: new Date(streamStartTime).toISOString(),\n\t\t});\n\n\t\tlet finalResponseText = \"\";\n\t\tlet collectionName: string | undefined;\n\n\t\tif (intents.length <= 1) {\n\t\t\t// Single intent: stream response directly\n\t\t\tconst triggeredIntent = intents[0];\n\t\t\tif (!triggeredIntent) {\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tconst { subquery = \"\", intent, actionPlan } = triggeredIntent;\n\t\t\tloggers.intent.info(\n\t\t\t\t`Process single intent: ${subquery}, ${intent?.name}`,\n\t\t\t);\n\t\t\tloggers.intent.info(`Action plan: ${actionPlan}`);\n\n\t\t\tconst intentThinking = this.buildIntentThinkingData(triggeredIntent);\n\t\t\tif (intentThinking) {\n\t\t\t\tyield { event: \"thinking_process\", data: intentThinking };\n\t\t\t}\n\n\t\t\t// Get the stream for this intent\n\t\t\tconst stream = this.getIntentStream(triggeredIntent, thread);\n\t\t\tif (!stream) {\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\t// Stream response directly\n\t\t\tfor await (const event of stream) {\n\t\t\t\tif (event.event === \"text_chunk\" && event.data.delta) {\n\t\t\t\t\tfinalResponseText += event.data.delta;\n\t\t\t\t} else if (event.event === \"collection_name\") {\n\t\t\t\t\tcollectionName = event.data.name;\n\t\t\t\t}\n\t\t\t\tyield event;\n\t\t\t}\n\t\t} else if (!needsAggregation) {\n\t\t\t// Multiple intents but no aggregation needed: collect intermediate results, stream only last\n\t\t\tfor (let i = 0; i < intents.length; i++) {\n\t\t\t\tconst triggeredIntent = intents[i];\n\t\t\t\tconst { subquery = \"\", intent, actionPlan } = triggeredIntent;\n\t\t\t\tloggers.intent.info(`Process query: ${subquery}, ${intent?.name}`);\n\t\t\t\tloggers.intent.info(`Action plan: ${actionPlan}`);\n\n\t\t\t\tconst isLastIntent = i === intents.length - 1;\n\n\t\t\t\tconst intentThinking = this.buildIntentThinkingData(triggeredIntent);\n\t\t\t\tif (intentThinking) {\n\t\t\t\t\tyield { event: \"thinking_process\", data: intentThinking };\n\t\t\t\t}\n\n\t\t\t\t// Get the stream for this intent\n\t\t\t\tconst stream = this.getIntentStream(triggeredIntent, thread);\n\t\t\t\tif (!stream) {\n\t\t\t\t\tcontinue;\n\t\t\t\t}\n\n\t\t\t\tif (isLastIntent) {\n\t\t\t\t\t// Stream last intent response directly\n\t\t\t\t\tfor await (const event of stream) {\n\t\t\t\t\t\tif (event.event === \"text_chunk\" && event.data.delta) {\n\t\t\t\t\t\t\tfinalResponseText += event.data.delta;\n\t\t\t\t\t\t} else if (event.event === \"collection_name\") {\n\t\t\t\t\t\t\tcollectionName = event.data.name;\n\t\t\t\t\t\t}\n\t\t\t\t\t\tyield event;\n\t\t\t\t\t}\n\t\t\t\t} else {\n\t\t\t\t\t// Collect intermediate results without streaming text_chunk\n\t\t\t\t\tlet responseText = \"\";\n\t\t\t\t\tfor await (const event of stream) {\n\t\t\t\t\t\tif (event.event === \"text_chunk\" && event.data.delta) {\n\t\t\t\t\t\t\tresponseText += event.data.delta;\n\t\t\t\t\t\t} else if (event.event === \"collection_name\") {\n\t\t\t\t\t\t\tcollectionName = event.data.name;\n\t\t\t\t\t\t} else if (event.event === \"thinking_process\") {\n\t\t\t\t\t\t\t// Tool execution thinking_process events are yielded immediately\n\t\t\t\t\t\t\tyield event;\n\t\t\t\t\t\t}\n\t\t\t\t\t}\n\t\t\t\t\t// Add intermediate result to thread context for next intent\n\t\t\t\t\tthread.messages.push({\n\t\t\t\t\t\tmessageId: randomUUID(),\n\t\t\t\t\t\trole: MessageRole.MODEL,\n\t\t\t\t\t\ttimestamp: Date.now(),\n\t\t\t\t\t\tcontent: { type: \"text\", parts: [responseText] },\n\t\t\t\t\t\tmetadata: { isThinking: true },\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t}\n\t\t} else {\n\t\t\t// Multi-intent mode with aggregation: collect all results then aggregate\n\t\t\tconst fulfillmentResults: FulfillmentResult[] = [];\n\n\t\t\tfor (let i = 0; i < intents.length; i++) {\n\t\t\t\tconst triggeredIntent = intents[i];\n\t\t\t\tconst { subquery = \"\", intent, actionPlan } = triggeredIntent;\n\t\t\t\tloggers.intent.info(`Process query: ${subquery}, ${intent?.name}`);\n\t\t\t\tloggers.intent.info(`Action plan: ${actionPlan}`);\n\n\t\t\t\t// Add previous result to thread context for inference (not stored in memory)\n\t\t\t\tif (fulfillmentResults.length > 0) {\n\t\t\t\t\tconst lastResult = fulfillmentResults[fulfillmentResults.length - 1];\n\t\t\t\t\tthread.messages.push({\n\t\t\t\t\t\tmessageId: randomUUID(),\n\t\t\t\t\t\trole: MessageRole.MODEL,\n\t\t\t\t\t\ttimestamp: Date.now(),\n\t\t\t\t\t\tcontent: { type: \"text\", parts: [lastResult.response] },\n\t\t\t\t\t\tmetadata: { isThinking: true },\n\t\t\t\t\t});\n\t\t\t\t}\n\n\t\t\t\tconst intentThinking = this.buildIntentThinkingData(triggeredIntent);\n\t\t\t\tif (intentThinking) {\n\t\t\t\t\tyield { event: \"thinking_process\", data: intentThinking };\n\t\t\t\t}\n\n\t\t\t\t// Get the stream for this intent\n\t\t\t\tconst stream = this.getIntentStream(triggeredIntent, thread);\n\t\t\t\tif (!stream) {\n\t\t\t\t\tcontinue;\n\t\t\t\t}\n\n\t\t\t\t// Collect response text (don't yield text_chunk yet)\n\t\t\t\tlet responseText = \"\";\n\t\t\t\tfor await (const event of stream) {\n\t\t\t\t\tif (event.event === \"text_chunk\" && event.data.delta) {\n\t\t\t\t\t\tresponseText += event.data.delta;\n\t\t\t\t\t} else if (event.event === \"collection_name\") {\n\t\t\t\t\t\tcollectionName = event.data.name;\n\t\t\t\t\t} else if (event.event === \"thinking_process\") {\n\t\t\t\t\t\t// Tool execution thinking_process events are yielded immediately\n\t\t\t\t\t\tyield event;\n\t\t\t\t\t}\n\t\t\t\t}\n\n\t\t\t\tfulfillmentResults.push({\n\t\t\t\t\tsubquery,\n\t\t\t\t\tintent,\n\t\t\t\t\tactionPlan,\n\t\t\t\t\tresponse: responseText,\n\t\t\t\t});\n\t\t\t}\n\n\t\t\t// Aggregate step: generate unified response\n\t\t\tconst aggregateStream = this.aggregateService.aggregate(\n\t\t\t\toriginalQuery,\n\t\t\t\tfulfillmentResults,\n\t\t\t);\n\n\t\t\tfor await (const event of aggregateStream) {\n\t\t\t\tif (event.event === \"text_chunk\" && event.data.delta) {\n\t\t\t\t\tfinalResponseText += event.data.delta;\n\t\t\t\t}\n\t\t\t\tyield event;\n\t\t\t}\n\t\t}\n\n\t\t// PII filtering on output before saving to memory (mask mode only)\n\t\tif (this.piiService?.getMode() === PIIFilterMode.MASK) {\n\t\t\tfinalResponseText = await this.piiService.filterText(finalResponseText);\n\t\t}\n\n\t\t// Save final response to memory\n\t\ttry {\n\t\t\tawait appendTextMessageToThread(\n\t\t\t\tthis.memoryModule,\n\t\t\t\tthread,\n\t\t\t\tMessageRole.MODEL,\n\t\t\t\tfinalResponseText,\n\t\t\t\tcollectionName ? { collectionName } : undefined,\n\t\t\t);\n\t\t} catch (error) {\n\t\t\tloggers.intentStream.error(\"Error adding message to thread\", error);\n\t\t}\n\n\t\tconst streamEndTime = Date.now();\n\t\tconst streamDuration = streamEndTime - streamStartTime;\n\n\t\tloggers.intentStream.info(\"Stream session completed\", {\n\t\t\tthreadId: thread.threadId,\n\t\t\tduration: `${streamDuration}ms`,\n\t\t\tendTime: new Date(streamEndTime).toISOString(),\n\t\t});\n\t}\n\n\tprivate buildIntentThinkingData(\n\t\ttriggeredIntent: TriggeredIntent,\n\t): Extract<StreamEvent, { event: \"thinking_process\" }>[\"data\"] | null {\n\t\tconst { intent, actionPlan } = triggeredIntent;\n\t\tif (!intent && !actionPlan) {\n\t\t\treturn null;\n\t\t}\n\t\treturn sanitizeThinkingData({\n\t\t\ttitle: `[${getManifest().name}] ${intent?.name || \"intent\"}`,\n\t\t\tdescription: actionPlan || \"\",\n\t\t});\n\t}\n}\n"]}