{"version":3,"sources":["/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-FGJNT6OA.cjs","../../src/services/workflow-execution.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;AACA;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACA;AC1BA,gCAA2B;AA0CpB,IAAM,yBAAA,EAAN,MAA+B;AAAA,EAC7B;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EAER,WAAA,CACC,mBAAA,EACA,wBAAA,EACA,WAAA,EACA,YAAA,EACA,kBAAA,EACA,SAAA,EACC;AACD,IAAA,IAAA,CAAK,oBAAA,EAAsB,mBAAA;AAC3B,IAAA,IAAA,CAAK,yBAAA,EAA2B,wBAAA;AAChC,IAAA,IAAA,CAAK,aAAA,EAAe,YAAA;AACpB,IAAA,IAAA,CAAK,mBAAA,EAAqB,IAAI,yCAAA;AAAA,MAC7B,WAAA;AAAA,MACA,kBAAA;AAAA,MACA;AAAA,IACD,CAAA;AACA,IAAA,IAAA,CAAK,yBAAA,EAA2B,IAAI,+CAAA;AAAA,MACnC,WAAA;AAAA,MACA,IAAI,2CAAA,CAAqB;AAAA,IAC1B,CAAA;AACA,IAAA,IAAA,CAAK,2BAAA,EAA6B,IAAI,wDAAA;AAAA,MACrC;AAAA,IACD,CAAA;AAAA,EACD;AAAA,EAEA,MAAM,eAAA,CACL,UAAA,EACA,kBAAA,EACA,MAAA,EACmC;AACnC,IAAA,IAAI,QAAA,EAAU,EAAA;AACd,IAAA,IAAI,QAAA;AACJ,IAAA,MAAM,OAAA,EAAS,IAAA,CAAK,qBAAA;AAAA,MACnB,UAAA;AAAA,MACA,kBAAA;AAAA,MACA;AAAA,IACD,CAAA;AAEA,IAAA,IAAA,MAAA,CAAA,MAAiB,MAAA,GAAS,MAAA,EAAQ;AACjC,MAAA,GAAA,CAAI,KAAA,CAAM,MAAA,IAAU,WAAA,EAAa;AAChC,QAAA,SAAA,EAAW,KAAA,CAAM,IAAA,CAAK,QAAA;AAAA,MACvB,EAAA,KAAA,GAAA,CAAW,KAAA,CAAM,MAAA,IAAU,YAAA,EAAc;AACxC,QAAA,QAAA,GAAW,KAAA,CAAM,IAAA,CAAK,KAAA;AAAA,MACvB;AAAA,IACD;AAEA,IAAA,OAAO,EAAE,OAAA,EAAS,SAAS,CAAA;AAAA,EAC5B;AAAA,EAEA,MAAA,CAAO,qBAAA,CACN,UAAA,EACA,kBAAA,EACA,MAAA,EAC8B;AAC9B,IAAA,MAAM,SAAA,EAAW,MAAM,IAAA,CAAK,mBAAA,CAAoB,WAAA,CAAY,UAAU,CAAA;AACtE,IAAA,GAAA,CAAI,CAAC,QAAA,EAAU;AACd,MAAA,MAAM,IAAI,KAAA,CAAM,CAAA,yBAAA,EAA4B,UAAU,CAAA,CAAA;AACvD,IAAA;AAGM,IAAA;AACJ,MAAA;AACA,MAAA;AACD,IAAA;AAEgB,IAAA;AACN,MAAA;AACa,QAAA;AACvB,MAAA;AACD,IAAA;AAEc,IAAA;AACwC,MAAA;AACrD,MAAA;AACC,QAAA;AAC4B,QAAA;AAC7B,MAAA;AACD,IAAA;AAE0B,IAAA;AACzB,MAAA;AACyB,MAAA;AACzB,MAAA;AACD,IAAA;AACM,IAAA;AACE,MAAA;AACD,MAAA;AACQ,QAAA;AACE,QAAA;AACE,QAAA;AACH,QAAA;AACK,QAAA;AACpB,MAAA;AACD,IAAA;AAGC,IAAA;AACC,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AAGiB,IAAA;AACd,IAAA;AACsC,MAAA;AACc,MAAA;AAGxB,QAAA;AACK,QAAA;AACC,QAAA;AACnC,UAAA;AACiB,UAAA;AACH,UAAA;AACd,UAAA;AACS,UAAA;AACD,UAAA;AACR,UAAA;AACA,UAAA;AACiB,UAAA;AACR,UAAA;AACE,UAAA;AACA,UAAA;AACX,QAAA;AACK,QAAA;AACA,UAAA;AACL,UAAA;AAAA,UAAA;AAE+C,UAAA;AAC/C,UAAA;AACC,YAAA;AACa,YAAA;AACb,YAAA;AACD,UAAA;AACD,QAAA;AACM,MAAA;AAGA,QAAA;AACA,UAAA;AACL,UAAA;AAAA,UAAA;AAEA,UAAA;AACA,UAAA;AACC,YAAA;AACa,YAAA;AACG,YAAA;AAC6B,YAAA;AAC9C,UAAA;AACD,QAAA;AACD,MAAA;AACmB,IAAA;AACC,MAAA;AACnB,QAAA;AACiB,QAAA;AACV,QAAA;AACP,MAAA;AACF,IAAA;AAEI,IAAA;AAC2C,MAAA;AAC5B,QAAA;AACG,QAAA;AACC,QAAA;AACrB,MAAA;AACoB,IAAA;AACD,MAAA;AACnB,QAAA;AACO,QAAA;AACP,MAAA;AACF,IAAA;AAEoB,IAAA;AACb,MAAA;AACP,IAAA;AACD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAgBC,EAAA;AAGgD,IAAA;AACjC,IAAA;AACE,MAAA;AACjB,IAAA;AAEgD,IAAA;AACtC,MAAA;AACV,IAAA;AACI,IAAA;AACgD,IAAA;AAEvC,MAAA;AACV,QAAA;AACA,QAAA;AAC6C,QAAA;AAC9C,MAAA;AACF,IAAA;AAMqD,IAAA;AACpD,MAAA;AACA,MAAA;AACD,IAAA;AACiB,IAAA;AACN,MAAA;AACa,QAAA;AACvB,MAAA;AACD,IAAA;AAEuD,IAAA;AACtD,MAAA;AACwB,MAAA;AACI,MAAA;AACH,MAAA;AACS,MAAA;AAClC,IAAA;AAI4B,IAAA;AAC5B,MAAA;AACmB,MAAA;AACE,MAAA;AACL,MAAA;AAChB,MAAA;AACW,MAAA;AACZ,IAAA;AAEuC,IAAA;AACtC,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AAEoB,IAAA;AACb,MAAA;AACP,IAAA;AACD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAWC,EAAA;AAWyD,IAAA;AACR,IAAA;AAC9B,IAAA;AACf,IAAA;AACA,IAAA;AAEA,IAAA;AACG,MAAA;AACE,QAAA;AACD,QAAA;AACE,UAAA;AACM,UAAA;AACH,UAAA;AACF,YAAA;AACP,YAAA;AACD,UAAA;AACD,QAAA;AACD,MAAA;AAEkD,MAAA;AAC5B,QAAA;AACJ,UAAA;AACjB,QAAA;AAC+B,QAAA;AAER,QAAA;AACK,UAAA;AACb,YAAA;AACgB,YAAA;AACjB,YAAA;AACJ,YAAA;AACC,YAAA;AACgC,YAAA;AACrB,YAAA;AACE,YAAA;AACvB,UAAA;AACM,UAAA;AACE,YAAA;AACD,YAAA;AACE,cAAA;AACe,cAAA;AACZ,cAAA;AACF,gBAAA;AACM,gBAAA;AACL,gBAAA;AACM,gBAAA;AACf,cAAA;AACD,YAAA;AACD,UAAA;AACA,UAAA;AACD,QAAA;AAEc,QAAA;AACkC,UAAA;AAC/C,UAAA;AACC,YAAA;AACiB,YAAA;AACJ,YAAA;AACuB,YAAA;AACxB,YAAA;AAC4B,YAAA;AACzC,UAAA;AACD,QAAA;AACuC,QAAA;AACtC,UAAA;AACA,UAAA;AACA,UAAA;AACD,QAAA;AAC+B,QAAA;AACV,QAAA;AACqB,UAAA;AAC1B,YAAA;AACb,cAAA;AACA,cAAA;AACC,gBAAA;AACiB,gBAAA;AACJ,gBAAA;AACoB,gBAAA;AACW,gBAAA;AAC7C,cAAA;AACD,YAAA;AACM,UAAA;AACO,YAAA;AACd,UAAA;AAC2B,UAAA;AAC5B,QAAA;AACkC,QAAA;AACpB,QAAA;AACkC,UAAA;AAC/C,UAAA;AACC,YAAA;AACiB,YAAA;AACJ,YAAA;AACQ,YAAA;AACyB,YAAA;AACC,YAAA;AACA,YAAA;AAC3B,YAAA;AACrB,UAAA;AACD,QAAA;AAEsC,QAAA;AACZ,UAAA;AAC1B,QAAA;AACD,MAAA;AAEuB,MAAA;AACsB,QAAA;AAClC,QAAA;AACuC,UAAA;AACjD,QAAA;AACD,MAAA;AAEM,MAAA;AACE,QAAA;AACD,QAAA;AACE,UAAA;AACM,UAAA;AACH,UAAA;AACF,YAAA;AACP,YAAA;AACuC,YAAA;AACxC,UAAA;AACD,QAAA;AACD,MAAA;AAEuD,MAAA;AACjC,QAAA;AACJ,UAAA;AACjB,QAAA;AAC0C,QAAA;AAC5B,QAAA;AACkC,UAAA;AAC/C,UAAA;AACC,YAAA;AACiB,YAAA;AACF,YAAA;AACE,YAAA;AAEW,YAAA;AAC7B,UAAA;AACD,QAAA;AAC6C,QAAA;AAC5C,UAAA;AACA,UAAA;AACA,UAAA;AACD,QAAA;AAC+B,QAAA;AACV,QAAA;AACqB,UAAA;AACN,YAAA;AACnC,UAAA;AACa,UAAA;AACc,UAAA;AAC5B,QAAA;AACgC,QAAA;AAClB,QAAA;AACiC,UAAA;AAC9C,UAAA;AACC,YAAA;AACiB,YAAA;AACK,YAAA;AACE,YAAA;AACuB,YAAA;AACA,YAAA;AAChD,UAAA;AACD,QAAA;AACD,MAAA;AAEmB,MAAA;AAClB,QAAA;AACiB,QAAA;AACjB,MAAA;AACc,IAAA;AAEuB,MAAA;AAClB,MAAA;AACnB,QAAA;AACiB,QAAA;AACK,QAAA;AACtB,MAAA;AACF,IAAA;AAEsD,IAAA;AACvD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAcC,EAAA;AAEyC,IAAA;AACpB,IAAA;AACgC,MAAA;AACrD,IAAA;AACkD,IAAA;AACnC,IAAA;AACqC,MAAA;AACpD,IAAA;AAEwD,IAAA;AACzC,IAAA;AACJ,MAAA;AACuC,QAAA;AACjD,MAAA;AACD,IAAA;AAO+C,IAAA;AACM,IAAA;AACpD,MAAA;AACA,MAAA;AACa,QAAA;AACD,QAAA;AAAA;AAAA;AAGD,QAAA;AACX,MAAA;AACD,IAAA;AACiB,IAAA;AACN,MAAA;AACqB,QAAA;AAC/B,MAAA;AACD,IAAA;AAE2B,IAAA;AACR,IAAA;AAClB,MAAA;AACoB,MAAA;AACI,MAAA;AACI,MAAA;AACG,MAAA;AAC/B,IAAA;AAI4B,IAAA;AAC5B,MAAA;AACiB,MAAA;AACI,MAAA;AACL,MAAA;AACI,MAAA;AACT,MAAA;AACZ,IAAA;AAGa,IAAA;AACX,MAAA;AACA,MAAA;AACQ,MAAA;AACR,MAAA;AACD,IAAA;AAEmB,IAAA;AAGoC,MAAA;AACtD,QAAA;AACoB,QAAA;AACK,QAAA;AACH,QAAA;AACtB,MAAA;AACK,MAAA;AACP,IAAA;AAC0B,IAAA;AACN,MAAA;AAClB,QAAA;AACoB,QAAA;AACK,QAAA;AACzB,MAAA;AACD,MAAA;AACD,IAAA;AAEI,IAAA;AAG6C,MAAA;AACvC,QAAA;AACE,UAAA;AACe,UAAA;AACzB,QAAA;AACA,MAAA;AACkB,IAAA;AACoC,MAAA;AACtD,QAAA;AACO,QAAA;AACP,MAAA;AACF,IAAA;AAEmB,IAAA;AAClB,MAAA;AACoB,MAAA;AACO,MAAA;AACF,MAAA;AACzB,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AA2B+B,EAAA;AAE5B,IAAA;AAEa,IAAA;AACkC,MAAA;AACjD,IAAA;AAEqD,IAAA;AACpD,MAAA;AAC+B,uBAAA;AAChC,IAAA;AACiB,IAAA;AACN,MAAA;AACsB,QAAA;AAChC,MAAA;AACD,IAAA;AAE2B,IAAA;AAC4B,IAAA;AACtD,MAAA;AACwB,MAAA;AACI,MAAA;AAC5B,IAAA;AAI4B,IAAA;AAC5B,MAAA;AACgB,MAAA;AACK,MAAA;AACL,MAAA;AACJ,MAAA;AACD,MAAA;AACZ,IAAA;AAGa,IAAA;AACX,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AAEmB,IAAA;AAGC,MAAA;AACnB,QAAA;AACyB,QAAA;AACH,QAAA;AACtB,MAAA;AACK,MAAA;AACP,IAAA;AAEmB,IAAA;AAClB,MAAA;AAC4B,MAAA;AACH,MAAA;AACzB,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAaC,EAAA;AAEc,IAAA;AACiB,IAAA;AAC9B,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACE,IAAA;AACgC,MAAA;AACX,QAAA;AACvB,MAAA;AACD,IAAA;AACqC,IAAA;AACtC,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAiBC,EAAA;AAEyC,IAAA;AACpB,IAAA;AACgC,MAAA;AACrD,IAAA;AAEkD,IAAA;AACnC,IAAA;AACqC,MAAA;AACpD,IAAA;AAEsD,IAAA;AAC3C,IAAA;AAC4C,MAAA;AACvD,IAAA;AAKgB,IAAA;AACC,IAAA;AACN,MAAA;AACwC,QAAA;AAClD,MAAA;AACD,IAAA;AAEU,IAAA;AAKJ,IAAA;AACE,MAAA;AACoB,MAAA;AAC5B,IAAA;AAGgD,IAAA;AACjC,IAAA;AACE,MAAA;AACjB,IAAA;AAEqD,IAAA;AACpD,MAAA;AACA,MAAA;AACD,IAAA;AACiB,IAAA;AACN,MAAA;AACa,QAAA;AACvB,MAAA;AACD,IAAA;AAK2B,IAAA;AACsB,IAAA;AAChD,MAAA;AACA,MAAA;AACA,MAAA;AACwB,MAAA;AACQ,MAAA;AAE7B,MAAA;AAC2B,QAAA;AACF,QAAA;AACsB,QAAA;AACP,QAAA;AAEvC,MAAA;AAED,MAAA;AAEH,IAAA;AAEuD,IAAA;AAC/C,MAAA;AACD,MAAA;AACP,IAAA;AAI4B,IAAA;AAC5B,MAAA;AACiB,MAAA;AACI,MAAA;AACL,MAAA;AAChB,MAAA;AACW,MAAA;AACZ,IAAA;AAGC,IAAA;AACC,MAAA;AACA,MAAA;AACA,MAAA;AACA,MAAA;AACD,IAAA;AAEoC,IAAA;AACY,MAAA;AACvC,QAAA;AAC0B,QAAA;AAClC,MAAA;AACmB,MAAA;AACb,QAAA;AACP,MAAA;AACA,MAAA;AACD,IAAA;AAEmC,IAAA;AACzB,MAAA;AACD,MAAA;AAC+B,MAAA;AACJ,MAAA;AACpC,IAAA;AACwD,IAAA;AAC/C,MAAA;AACR,MAAA;AACO,MAAA;AACP,IAAA;AACF,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAciB,EAAA;AACiC,IAAA;AAClD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AASwD,EAAA;AACH,IAAA;AAClC,IAAA;AACV,MAAA;AACR,IAAA;AAGE,IAAA;AACH,EAAA;AAKC,EAAA;AAEuD,IAAA;AAC3B,IAAA;AACW,IAAA;AACe,IAAA;AAAA,MAAA;AAE5C,MAAA;AACT,MAAA;AACA,MAAA;AACS,MAAA;AACJ,IAAA;AACL,MAAA;AACiB,MAAA;AACjB,MAAA;AACA,MAAA;AACqB,MAAA;AACtB,IAAA;AAEyD,IAAA;AACnD,IAAA;AACA,MAAA;AACL,MAAA;AAAA,MAAA;AAEA,MAAA;AACA,MAAA;AACsB,QAAA;AACR,QAAA;AACN,QAAA;AACR,MAAA;AACD,IAAA;AAEO,IAAA;AACR,EAAA;AACD;ADnN6D;AACA;AACA;AACA","file":"/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-FGJNT6OA.cjs","sourcesContent":[null,"import { randomUUID } from \"node:crypto\";\nimport type { A2AModule, MemoryModule, ModelModule } from \"@/modules\";\nimport {\n\ttype Document,\n\tDocumentFormat,\n\ttype DocumentFragment,\n\ttype DocumentSlot,\n\tDocumentSource,\n} from \"@/types/document.js\";\nimport {\n\tMessageRole,\n\ttype ThreadMetadata,\n\ttype ThreadObject,\n\tThreadType,\n\ttype UserWorkflow,\n\ttype WorkflowDefinition,\n\ttype WorkflowRenderedBlock,\n\ttype WorkflowTaskResult,\n\ttype WorkflowTemplate,\n} from \"@/types/memory.js\";\nimport type { SlotFillInitiator } from \"@/types/schedule.js\";\nimport type { StreamEvent } from \"@/types/stream.js\";\nimport { renderDocument } from \"@/utils/document-render.js\";\nimport { loggers } from \"@/utils/logger.js\";\nimport {\n\tappendRichMessageToThread,\n\tappendTextMessageToThread,\n} from \"@/utils/thread-messages.js\";\nimport { workflowTaskLabel } from \"@/utils/workflow-task-results.js\";\nimport type { ToolCallingService } from \"./tool-calling.service.js\";\nimport type { UserWorkflowService } from \"./user-workflow.service.js\";\nimport { WorkflowResponseComposer } from \"./workflow-response-composer.service.js\";\nimport { WorkflowTableService } from \"./workflow-table.service.js\";\nimport { WorkflowTaskRunner } from \"./workflow-task-runner.service.js\";\nimport { WorkflowVariableExtractionService } from \"./workflow-variable-extraction.service.js\";\nimport type { WorkflowVariableResolver } from \"./workflow-variable-resolver.service.js\";\n\ntype WorkflowExecutionResult = {\n\tcontent: string;\n\tthreadId?: string;\n};\n\nexport class WorkflowExecutionService {\n\tprivate userWorkflowService: UserWorkflowService;\n\tprivate workflowVariableResolver: WorkflowVariableResolver;\n\tprivate memoryModule: MemoryModule;\n\tprivate workflowTaskRunner: WorkflowTaskRunner;\n\tprivate workflowResponseComposer: WorkflowResponseComposer;\n\tprivate workflowVariableExtraction: WorkflowVariableExtractionService;\n\n\tconstructor(\n\t\tuserWorkflowService: UserWorkflowService,\n\t\tworkflowVariableResolver: WorkflowVariableResolver,\n\t\tmodelModule: ModelModule,\n\t\tmemoryModule: MemoryModule,\n\t\ttoolCallingService: ToolCallingService,\n\t\ta2aModule?: A2AModule,\n\t) {\n\t\tthis.userWorkflowService = userWorkflowService;\n\t\tthis.workflowVariableResolver = workflowVariableResolver;\n\t\tthis.memoryModule = memoryModule;\n\t\tthis.workflowTaskRunner = new WorkflowTaskRunner(\n\t\t\tmodelModule,\n\t\t\ttoolCallingService,\n\t\t\ta2aModule,\n\t\t);\n\t\tthis.workflowResponseComposer = new WorkflowResponseComposer(\n\t\t\tmodelModule,\n\t\t\tnew WorkflowTableService(),\n\t\t);\n\t\tthis.workflowVariableExtraction = new WorkflowVariableExtractionService(\n\t\t\tmodelModule,\n\t\t);\n\t}\n\n\tasync executeWorkflow(\n\t\tworkflowId: string,\n\t\texecutionVariables?: Record<string, string>,\n\t\tsignal?: AbortSignal,\n\t): Promise<WorkflowExecutionResult> {\n\t\tlet content = \"\";\n\t\tlet threadId: string | undefined;\n\t\tconst stream = this.executeWorkflowStream(\n\t\t\tworkflowId,\n\t\t\texecutionVariables,\n\t\t\tsignal,\n\t\t);\n\n\t\tfor await (const event of stream) {\n\t\t\tif (event.event === \"thread_id\") {\n\t\t\t\tthreadId = event.data.threadId;\n\t\t\t} else if (event.event === \"text_chunk\") {\n\t\t\t\tcontent += event.data.delta;\n\t\t\t}\n\t\t}\n\n\t\treturn { content, threadId };\n\t}\n\n\tasync *executeWorkflowStream(\n\t\tworkflowId: string,\n\t\texecutionVariables?: Record<string, string>,\n\t\tsignal?: AbortSignal,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst workflow = await this.userWorkflowService.getWorkflow(workflowId);\n\t\tif (!workflow) {\n\t\t\tthrow new Error(`User workflow not found: ${workflowId}`);\n\t\t}\n\n\t\tconst { query, displayQuery, definition } =\n\t\t\tthis.workflowVariableResolver.resolveForExecution(\n\t\t\t\tworkflow,\n\t\t\t\texecutionVariables,\n\t\t\t);\n\n\t\tif (!definition) {\n\t\t\tthrow new Error(\n\t\t\t\t`Workflow ${workflowId} has no valid structured definition; legacy content execution is no longer supported`,\n\t\t\t);\n\t\t}\n\n\t\tloggers.agent.info(\n\t\t\t`Executing structured user workflow: ${workflow.title}`,\n\t\t\t{\n\t\t\t\tworkflowId,\n\t\t\t\ttaskCount: definition.tasks.length,\n\t\t\t},\n\t\t);\n\n\t\tconst thread = await this.createWorkflowThread(\n\t\t\tworkflow,\n\t\t\tdisplayQuery || workflow.title,\n\t\t\tquery,\n\t\t);\n\t\tyield {\n\t\t\tevent: \"thread_id\",\n\t\t\tdata: {\n\t\t\t\ttype: thread.type,\n\t\t\t\tuserId: thread.userId,\n\t\t\t\tthreadId: thread.threadId,\n\t\t\t\ttitle: thread.title,\n\t\t\t\tworkflowId: thread.workflowId,\n\t\t\t},\n\t\t};\n\n\t\tconst { finalContent, renderedBlocks, executionError } =\n\t\t\tyield* this.renderStructuredDefinition(\n\t\t\t\tdefinition,\n\t\t\t\tthread,\n\t\t\t\tworkflowId,\n\t\t\t\tsignal,\n\t\t\t);\n\n\t\tconst responseContent =\n\t\t\tfinalContent || (executionError ? `오류: ${executionError.message}` : \"\");\n\t\ttry {\n\t\t\tconst documentMemory = this.memoryModule.getDocumentMemory();\n\t\t\tif (documentMemory && !executionError && finalContent) {\n\t\t\t\t// Promote the workflow result to a first-class document and\n\t\t\t\t// reference it from the thread (body is resolved on demand).\n\t\t\t\tconst documentId = randomUUID();\n\t\t\t\tconst now = new Date().toISOString();\n\t\t\t\tawait documentMemory.createDocument({\n\t\t\t\t\tdocumentId,\n\t\t\t\t\tuserId: workflow.userId,\n\t\t\t\t\ttitle: thread.title,\n\t\t\t\t\tformat: DocumentFormat.MARKDOWN,\n\t\t\t\t\tcontent: finalContent,\n\t\t\t\t\tblocks: renderedBlocks,\n\t\t\t\t\tsource: DocumentSource.WORKFLOW,\n\t\t\t\t\tworkflowId,\n\t\t\t\t\tthreadId: thread.threadId,\n\t\t\t\t\tversion: 1,\n\t\t\t\t\tcreatedAt: now,\n\t\t\t\t\tupdatedAt: now,\n\t\t\t\t});\n\t\t\t\tawait appendRichMessageToThread(\n\t\t\t\t\tthis.memoryModule,\n\t\t\t\t\tthread,\n\t\t\t\t\tMessageRole.MODEL,\n\t\t\t\t\t[{ type: \"document\", documentId, title: thread.title }],\n\t\t\t\t\t{\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tworkflowRun: true,\n\t\t\t\t\t\tdocumentId,\n\t\t\t\t\t},\n\t\t\t\t);\n\t\t\t} else {\n\t\t\t\t// No document memory (or execution failed): keep the legacy\n\t\t\t\t// inline text message with structured blocks in metadata.\n\t\t\t\tawait appendTextMessageToThread(\n\t\t\t\t\tthis.memoryModule,\n\t\t\t\t\tthread,\n\t\t\t\t\tMessageRole.MODEL,\n\t\t\t\t\tresponseContent,\n\t\t\t\t\t{\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tworkflowRun: true,\n\t\t\t\t\t\tresponseBlocks: renderedBlocks,\n\t\t\t\t\t\t...(executionError ? { error: executionError.message } : {}),\n\t\t\t\t\t},\n\t\t\t\t);\n\t\t\t}\n\t\t} catch (saveError) {\n\t\t\tloggers.agent.error(\"Failed to save workflow response message\", {\n\t\t\t\tworkflowId,\n\t\t\t\tthreadId: thread.threadId,\n\t\t\t\terror: saveError,\n\t\t\t});\n\t\t}\n\n\t\ttry {\n\t\t\tawait this.userWorkflowService.updateWorkflow(workflowId, {\n\t\t\t\tuserId: workflow.userId,\n\t\t\t\tlastRunAt: Date.now(),\n\t\t\t\tlastThreadId: thread.threadId,\n\t\t\t});\n\t\t} catch (updateError) {\n\t\t\tloggers.agent.error(\"Failed to update workflow lastRunAt\", {\n\t\t\t\tworkflowId,\n\t\t\t\terror: updateError,\n\t\t\t});\n\t\t}\n\n\t\tif (executionError) {\n\t\t\tthrow executionError;\n\t\t}\n\t}\n\n\t/**\n\t * Runs an intent-mapped workflow and streams its progress into the chat\n\t * stream. Unlike {@link executeWorkflowStream}, this does NOT create or\n\t * persist a workflow thread, document, or lastRunAt bookkeeping — the\n\t * chat thread (persisted by the intent fulfillment path) is the only\n\t * artifact. Accepts a user workflow id or a template id, like slot\n\t * bindings. Variable values are extracted from the subquery via one LLM\n\t * call and merged over the workflow's stored variableValues. Setup\n\t * failures (unknown id, no definition) throw before the first yield so\n\t * callers can fall back; task failures throw after streaming.\n\t */\n\tasync *executeIntentWorkflowStream(\n\t\tworkflowId: string,\n\t\tchatThread: ThreadObject,\n\t\tsubquery: string,\n\t\tsignal?: AbortSignal,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst workflow = await this.getFillableWorkflow(workflowId);\n\t\tif (!workflow) {\n\t\t\tthrow new Error(`User workflow or template not found: ${workflowId}`);\n\t\t}\n\n\t\tconst variables = this.workflowVariableResolver.normalizeVariables(\n\t\t\tworkflow.variables,\n\t\t);\n\t\tlet extractedVariables: Record<string, string> | undefined;\n\t\tif (variables && Object.keys(variables).length > 0) {\n\t\t\textractedVariables =\n\t\t\t\tawait this.workflowVariableExtraction.extractFromQuery(\n\t\t\t\t\tvariables,\n\t\t\t\t\tsubquery,\n\t\t\t\t\t\"timezone\" in workflow ? workflow.timezone : undefined,\n\t\t\t\t);\n\t\t}\n\n\t\t// Document-fill semantics: an intent-mapped run has no creation step,\n\t\t// so every provided value applies regardless of its declared\n\t\t// resolveAt; values the model couldn't extract fall back to the\n\t\t// workflow's stored variableValues.\n\t\tconst { definition } = this.workflowVariableResolver.resolveForDocumentFill(\n\t\t\tworkflow,\n\t\t\textractedVariables,\n\t\t);\n\t\tif (!definition) {\n\t\t\tthrow new Error(\n\t\t\t\t`Workflow ${workflowId} has no valid structured definition; cannot fulfill intent`,\n\t\t\t);\n\t\t}\n\n\t\tloggers.agent.info(\"Executing intent-mapped workflow\", {\n\t\t\tworkflowId,\n\t\t\tworkflowTitle: workflow.title,\n\t\t\ttaskCount: definition.tasks.length,\n\t\t\tchatThreadId: chatThread.threadId,\n\t\t\textractedVariableIds: Object.keys(extractedVariables ?? {}),\n\t\t});\n\n\t\t// Ephemeral, non-persisted thread: carries threadId for A2A correlation\n\t\t// and task context, but is never written to the thread store.\n\t\tconst thread: ThreadObject = {\n\t\t\ttype: ThreadType.WORKFLOW,\n\t\t\tuserId: chatThread.userId,\n\t\t\tthreadId: randomUUID(),\n\t\t\ttitle: workflow.title,\n\t\t\tworkflowId,\n\t\t\tmessages: [],\n\t\t};\n\n\t\tconst { executionError } = yield* this.renderStructuredDefinition(\n\t\t\tdefinition,\n\t\t\tthread,\n\t\t\tworkflowId,\n\t\t\tsignal,\n\t\t);\n\n\t\tif (executionError) {\n\t\t\tthrow executionError;\n\t\t}\n\t}\n\n\t/**\n\t * Runs a structured workflow definition (tasks → response blocks) against a\n\t * thread context, streaming progress events. Never throws — any failure is\n\t * captured and returned as `executionError` so callers can decide how to\n\t * persist the (partial) result.\n\t */\n\tprivate async *renderStructuredDefinition(\n\t\tdefinition: WorkflowDefinition,\n\t\tthread: ThreadObject,\n\t\tworkflowId: string,\n\t\tsignal?: AbortSignal,\n\t): AsyncGenerator<\n\t\tStreamEvent,\n\t\t{\n\t\t\tfinalContent: string;\n\t\t\trenderedBlocks: WorkflowRenderedBlock[];\n\t\t\texecutionError?: Error;\n\t\t},\n\t\tunknown\n\t> {\n\t\tconst taskResults: Record<string, WorkflowTaskResult> = {};\n\t\tconst renderedBlocks: WorkflowRenderedBlock[] = [];\n\t\tlet finalContent = \"\";\n\t\tlet executionError: Error | undefined;\n\t\tlet firstFailedTaskId: string | undefined;\n\n\t\ttry {\n\t\t\tyield {\n\t\t\t\tevent: \"thinking_process\",\n\t\t\t\tdata: {\n\t\t\t\t\ttitle: \"[워크플로우] 실행\",\n\t\t\t\t\tdescription: \"워크플로우 실행을 시작합니다.\",\n\t\t\t\t\tmetadata: {\n\t\t\t\t\t\tphase: \"workflow_start\",\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t},\n\t\t\t\t},\n\t\t\t};\n\n\t\t\tfor (let i = 0; i < definition.tasks.length; i++) {\n\t\t\t\tif (signal?.aborted) {\n\t\t\t\t\tthrow new Error(\"Workflow execution aborted by client\");\n\t\t\t\t}\n\t\t\t\tconst task = definition.tasks[i];\n\n\t\t\t\tif (firstFailedTaskId) {\n\t\t\t\t\ttaskResults[task.taskId] = {\n\t\t\t\t\t\ttaskId: task.taskId,\n\t\t\t\t\t\ttitle: workflowTaskLabel(task),\n\t\t\t\t\t\tagent: task.agent,\n\t\t\t\t\t\tstatus: \"skipped\",\n\t\t\t\t\t\tcontent: \"\",\n\t\t\t\t\t\terror: `Skipped due to failure of task ${firstFailedTaskId}`,\n\t\t\t\t\t\tstartedAt: Date.now(),\n\t\t\t\t\t\tcompletedAt: Date.now(),\n\t\t\t\t\t};\n\t\t\t\t\tyield {\n\t\t\t\t\t\tevent: \"thinking_process\",\n\t\t\t\t\t\tdata: {\n\t\t\t\t\t\t\ttitle: `[워크플로우] 작업 건너뜀: ${workflowTaskLabel(task)}`,\n\t\t\t\t\t\t\tdescription: `이전 작업(${firstFailedTaskId}) 실패로 인해 건너뜁니다.`,\n\t\t\t\t\t\t\tmetadata: {\n\t\t\t\t\t\t\t\tphase: \"task_skipped\",\n\t\t\t\t\t\t\t\ttaskId: task.taskId,\n\t\t\t\t\t\t\t\treason: \"previous_task_failed\",\n\t\t\t\t\t\t\t\tfailedTaskId: firstFailedTaskId,\n\t\t\t\t\t\t\t},\n\t\t\t\t\t\t},\n\t\t\t\t\t};\n\t\t\t\t\tcontinue;\n\t\t\t\t}\n\n\t\t\t\tloggers.agent.debug(\n\t\t\t\t\t`Workflow task starting (${i + 1}/${definition.tasks.length}): ${workflowTaskLabel(task)}`,\n\t\t\t\t\t{\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tthreadId: thread.threadId,\n\t\t\t\t\t\ttaskId: task.taskId,\n\t\t\t\t\t\texecutionType: task.agent ? \"a2a\" : \"local\",\n\t\t\t\t\t\tagent: task.agent,\n\t\t\t\t\t\tpromptPreview: task.prompt?.slice(0, 200),\n\t\t\t\t\t},\n\t\t\t\t);\n\t\t\t\tconst stream = this.workflowTaskRunner.executeTask(\n\t\t\t\t\ttask,\n\t\t\t\t\tthread,\n\t\t\t\t\ttaskResults,\n\t\t\t\t);\n\t\t\t\tlet result = await stream.next();\n\t\t\t\twhile (!result.done) {\n\t\t\t\t\tif (result.value.event === \"text_chunk\") {\n\t\t\t\t\t\tloggers.agent.warn(\n\t\t\t\t\t\t\t\"Suppressed unexpected workflow task text_chunk before response phase\",\n\t\t\t\t\t\t\t{\n\t\t\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\t\t\tthreadId: thread.threadId,\n\t\t\t\t\t\t\t\ttaskId: task.taskId,\n\t\t\t\t\t\t\t\ttaskTitle: workflowTaskLabel(task),\n\t\t\t\t\t\t\t\tdeltaPreview: result.value.data.delta.slice(0, 200),\n\t\t\t\t\t\t\t},\n\t\t\t\t\t\t);\n\t\t\t\t\t} else {\n\t\t\t\t\t\tyield result.value;\n\t\t\t\t\t}\n\t\t\t\t\tresult = await stream.next();\n\t\t\t\t}\n\t\t\t\ttaskResults[task.taskId] = result.value;\n\t\t\t\tloggers.agent.debug(\n\t\t\t\t\t`Workflow task finished (${i + 1}/${definition.tasks.length}): ${workflowTaskLabel(task)}`,\n\t\t\t\t\t{\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tthreadId: thread.threadId,\n\t\t\t\t\t\ttaskId: task.taskId,\n\t\t\t\t\t\tstatus: result.value.status,\n\t\t\t\t\t\tdurationMs: result.value.completedAt - result.value.startedAt,\n\t\t\t\t\t\tcontentLength: result.value.content?.length ?? 0,\n\t\t\t\t\t\tcontentPreview: result.value.content?.slice(0, 500),\n\t\t\t\t\t\terror: result.value.error,\n\t\t\t\t\t},\n\t\t\t\t);\n\n\t\t\t\tif (result.value.status === \"failed\") {\n\t\t\t\t\tfirstFailedTaskId = task.taskId;\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tif (firstFailedTaskId) {\n\t\t\t\tconst failed = taskResults[firstFailedTaskId];\n\t\t\t\tthrow new Error(\n\t\t\t\t\t`Workflow task failed: ${firstFailedTaskId} - ${failed?.error ?? \"unknown error\"}`,\n\t\t\t\t);\n\t\t\t}\n\n\t\t\tyield {\n\t\t\t\tevent: \"thinking_process\",\n\t\t\t\tdata: {\n\t\t\t\t\ttitle: \"[워크플로우] 응답 구성 중\",\n\t\t\t\t\tdescription: \"작업 결과를 바탕으로 워크플로우 응답을 구성합니다.\",\n\t\t\t\t\tmetadata: {\n\t\t\t\t\t\tphase: \"response_start\",\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tblockCount: definition.response.blocks.length,\n\t\t\t\t\t},\n\t\t\t\t},\n\t\t\t};\n\n\t\t\tfor (let i = 0; i < definition.response.blocks.length; i++) {\n\t\t\t\tif (signal?.aborted) {\n\t\t\t\t\tthrow new Error(\"Workflow execution aborted by client\");\n\t\t\t\t}\n\t\t\t\tconst block = definition.response.blocks[i];\n\t\t\t\tloggers.agent.debug(\n\t\t\t\t\t`Workflow response block rendering (${i + 1}/${definition.response.blocks.length})`,\n\t\t\t\t\t{\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tthreadId: thread.threadId,\n\t\t\t\t\t\tblockId: block.blockId,\n\t\t\t\t\t\tblockType: block.type,\n\t\t\t\t\t\tsourceTaskIds:\n\t\t\t\t\t\t\tblock.type === \"heading\" ? undefined : block.sourceTaskIds,\n\t\t\t\t\t},\n\t\t\t\t);\n\t\t\t\tconst stream = this.workflowResponseComposer.renderResponseBlock(\n\t\t\t\t\tblock,\n\t\t\t\t\ttaskResults,\n\t\t\t\t\trenderedBlocks,\n\t\t\t\t);\n\t\t\t\tlet result = await stream.next();\n\t\t\t\twhile (!result.done) {\n\t\t\t\t\tif (result.value.event === \"text_chunk\") {\n\t\t\t\t\t\tfinalContent += result.value.data.delta;\n\t\t\t\t\t}\n\t\t\t\t\tyield result.value;\n\t\t\t\t\tresult = await stream.next();\n\t\t\t\t}\n\t\t\t\trenderedBlocks.push(result.value);\n\t\t\t\tloggers.agent.debug(\n\t\t\t\t\t`Workflow response block rendered (${i + 1}/${definition.response.blocks.length})`,\n\t\t\t\t\t{\n\t\t\t\t\t\tworkflowId,\n\t\t\t\t\t\tthreadId: thread.threadId,\n\t\t\t\t\t\tblockId: result.value.blockId,\n\t\t\t\t\t\tblockType: result.value.type,\n\t\t\t\t\t\tcontentLength: result.value.content?.length ?? 0,\n\t\t\t\t\t\tcontentPreview: result.value.content?.slice(0, 500),\n\t\t\t\t\t},\n\t\t\t\t);\n\t\t\t}\n\n\t\t\tloggers.agent.info(\"Structured workflow definition completed\", {\n\t\t\t\tworkflowId,\n\t\t\t\tthreadId: thread.threadId,\n\t\t\t});\n\t\t} catch (error) {\n\t\t\texecutionError =\n\t\t\t\terror instanceof Error ? error : new Error(String(error));\n\t\t\tloggers.agent.error(\"Structured workflow definition failed\", {\n\t\t\t\tworkflowId,\n\t\t\t\tthreadId: thread.threadId,\n\t\t\t\terror: executionError.message,\n\t\t\t});\n\t\t}\n\n\t\treturn { finalContent, renderedBlocks, executionError };\n\t}\n\n\t/**\n\t * Generates AI advice for a document by running the bound advice workflow\n\t * over the document's rendered content, then caches the result on\n\t * `document.advice`. Mirrors {@link fillDocumentSlotStream}: ephemeral\n\t * non-persisted thread — the advice field is the only artifact.\n\t */\n\tasync *generateDocumentAdviceStream(\n\t\tdocumentId: string,\n\t\toptions: {\n\t\t\tworkflowId: string;\n\t\t\texecutionVariables?: Record<string, string>;\n\t\t},\n\t\tsignal?: AbortSignal,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst documentMemory = this.memoryModule.getDocumentMemory();\n\t\tif (!documentMemory) {\n\t\t\tthrow new Error(\"Document memory is not initialized\");\n\t\t}\n\t\tconst document = await documentMemory.getDocument(documentId);\n\t\tif (!document) {\n\t\t\tthrow new Error(`Document not found: ${documentId}`);\n\t\t}\n\n\t\tconst workflow = await this.getFillableWorkflow(options.workflowId);\n\t\tif (!workflow) {\n\t\t\tthrow new Error(\n\t\t\t\t`User workflow or template not found: ${options.workflowId}`,\n\t\t\t);\n\t\t}\n\n\t\t// The document IS the variable source: its labels become variables of\n\t\t// the same name (so a new document kind needs no code change), the\n\t\t// caller's executionVariables override them by name, and the rendered\n\t\t// body arrives as {{document}}. Variable substitution is deep, so all\n\t\t// three reach task AND response-block prompts alike.\n\t\tconst renderedContent = renderDocument(document);\n\t\tconst { definition } = this.workflowVariableResolver.resolveForDocumentFill(\n\t\t\tworkflow,\n\t\t\t{\n\t\t\t\t...document.labels,\n\t\t\t\t...options.executionVariables,\n\t\t\t\t// Last on purpose: replacements apply in insertion order, so a\n\t\t\t\t// later one would otherwise rewrite text inside the body.\n\t\t\t\tdocument: renderedContent,\n\t\t\t},\n\t\t);\n\t\tif (!definition) {\n\t\t\tthrow new Error(\n\t\t\t\t`Workflow ${options.workflowId} has no valid structured definition; cannot generate advice`,\n\t\t\t);\n\t\t}\n\n\t\tconst startedAt = Date.now();\n\t\tloggers.agent.info(\"Generating document advice via workflow\", {\n\t\t\tdocumentId,\n\t\t\tworkflowId: options.workflowId,\n\t\t\tworkflowTitle: workflow.title,\n\t\t\ttaskCount: definition.tasks.length,\n\t\t\tcontentLength: renderedContent.length,\n\t\t});\n\n\t\t// Ephemeral, non-persisted thread: carries threadId for A2A correlation\n\t\t// and task context, but is never written to the thread store.\n\t\tconst thread: ThreadObject = {\n\t\t\ttype: ThreadType.WORKFLOW,\n\t\t\tuserId: document.userId,\n\t\t\tthreadId: randomUUID(),\n\t\t\ttitle: workflow.title,\n\t\t\tworkflowId: options.workflowId,\n\t\t\tmessages: [],\n\t\t};\n\n\t\tconst { finalContent, executionError } =\n\t\t\tyield* this.renderStructuredDefinition(\n\t\t\t\tdefinition,\n\t\t\t\tthread,\n\t\t\t\toptions.workflowId,\n\t\t\t\tsignal,\n\t\t\t);\n\n\t\tif (executionError) {\n\t\t\t// renderStructuredDefinition already logged the task-level failure;\n\t\t\t// this ties it to the advice request before the SSE layer reports it.\n\t\t\tloggers.agent.error(\"Document advice workflow failed\", {\n\t\t\t\tdocumentId,\n\t\t\t\tworkflowId: options.workflowId,\n\t\t\t\tdurationMs: Date.now() - startedAt,\n\t\t\t\terror: executionError.message,\n\t\t\t});\n\t\t\tthrow executionError;\n\t\t}\n\t\tif (!finalContent.trim()) {\n\t\t\tloggers.agent.warn(\"Document advice workflow produced no content\", {\n\t\t\t\tdocumentId,\n\t\t\t\tworkflowId: options.workflowId,\n\t\t\t\tdurationMs: Date.now() - startedAt,\n\t\t\t});\n\t\t\treturn;\n\t\t}\n\n\t\ttry {\n\t\t\t// Persist only the advice (metadata); do NOT bump version off a\n\t\t\t// pre-stream read, which could clobber a concurrent edit (lost update).\n\t\t\tawait documentMemory.updateDocument(documentId, {\n\t\t\t\tadvice: {\n\t\t\t\t\tcontent: finalContent,\n\t\t\t\t\tgeneratedAt: new Date().toISOString(),\n\t\t\t\t},\n\t\t\t});\n\t\t} catch (saveError) {\n\t\t\tloggers.agent.error(\"Failed to cache document advice\", {\n\t\t\t\tdocumentId,\n\t\t\t\terror: saveError,\n\t\t\t});\n\t\t}\n\n\t\tloggers.agent.info(\"Document advice generated via workflow\", {\n\t\t\tdocumentId,\n\t\t\tworkflowId: options.workflowId,\n\t\t\tadviceLength: finalContent.length,\n\t\t\tdurationMs: Date.now() - startedAt,\n\t\t});\n\t}\n\n\t/**\n\t * Runs a workflow **template** with caller-supplied variables and streams the\n\t * rendered result. Nothing is persisted — no thread, no document — so the\n\t * caller's screen is the only artifact.\n\t *\n\t * This is {@link generateDocumentAdviceStream} without the document. Dashboards\n\t * show a workflow's output against a selection (a period, a group of outlets)\n\t * that no document represents, and {@link executeWorkflowStream} cannot serve\n\t * them: it resolves user workflows only and it persists a chat thread per run.\n\t *\n\t * Templates only, deliberately: the screens that need this have no per-user\n\t * workflow copy — everyone shares one catalog entry. A user workflow id is not\n\t * accepted here, so there is no ownership question to answer.\n\t *\n\t * Variables resolve with document-fill semantics — there is no creation step, so\n\t * every supplied value applies regardless of its declared `resolveAt`, and values\n\t * the caller omits fall back to the template's stored ones.\n\t */\n\tasync *runTemplateStream(\n\t\ttemplateId: string,\n\t\toptions: {\n\t\t\tuserId: string;\n\t\t\texecutionVariables?: Record<string, string>;\n\t\t},\n\t\tsignal?: AbortSignal,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst template = await this.memoryModule\n\t\t\t.getWorkflowTemplateMemory()\n\t\t\t.getTemplate(templateId);\n\t\tif (!template) {\n\t\t\tthrow new Error(`Workflow template not found: ${templateId}`);\n\t\t}\n\n\t\tconst { definition } = this.workflowVariableResolver.resolveForDocumentFill(\n\t\t\ttemplate,\n\t\t\toptions.executionVariables ?? {},\n\t\t);\n\t\tif (!definition) {\n\t\t\tthrow new Error(\n\t\t\t\t`Workflow template ${templateId} has no valid structured definition; cannot run`,\n\t\t\t);\n\t\t}\n\n\t\tconst startedAt = Date.now();\n\t\tloggers.agent.info(\"Running workflow template ad hoc\", {\n\t\t\ttemplateId,\n\t\t\ttemplateTitle: template.title,\n\t\t\ttaskCount: definition.tasks.length,\n\t\t});\n\n\t\t// Ephemeral, non-persisted thread: carries threadId for A2A correlation\n\t\t// and task context, but is never written to the thread store.\n\t\tconst thread: ThreadObject = {\n\t\t\ttype: ThreadType.WORKFLOW,\n\t\t\tuserId: options.userId,\n\t\t\tthreadId: randomUUID(),\n\t\t\ttitle: template.title,\n\t\t\tworkflowId: templateId,\n\t\t\tmessages: [],\n\t\t};\n\n\t\tconst { finalContent, executionError } =\n\t\t\tyield* this.renderStructuredDefinition(\n\t\t\t\tdefinition,\n\t\t\t\tthread,\n\t\t\t\ttemplateId,\n\t\t\t\tsignal,\n\t\t\t);\n\n\t\tif (executionError) {\n\t\t\t// renderStructuredDefinition already logged the task-level failure;\n\t\t\t// this ties it to the run request before the SSE layer reports it.\n\t\t\tloggers.agent.error(\"Ad hoc workflow template run failed\", {\n\t\t\t\ttemplateId,\n\t\t\t\tdurationMs: Date.now() - startedAt,\n\t\t\t\terror: executionError.message,\n\t\t\t});\n\t\t\tthrow executionError;\n\t\t}\n\n\t\tloggers.agent.info(\"Ad hoc workflow template run finished\", {\n\t\t\ttemplateId,\n\t\t\tcontentLength: finalContent.length,\n\t\t\tdurationMs: Date.now() - startedAt,\n\t\t});\n\t}\n\n\t/**\n\t * Non-streaming variant of {@link fillDocumentSlotStream}.\n\t */\n\tasync fillDocumentSlot(\n\t\tdocumentId: string,\n\t\tslotId: string,\n\t\toptions?: {\n\t\t\tworkflowId?: string;\n\t\t\texecutionVariables?: Record<string, string>;\n\t\t\tinitiator?: SlotFillInitiator;\n\t\t},\n\t\tsignal?: AbortSignal,\n\t): Promise<{ documentId: string; slotId: string; content: string }> {\n\t\tlet content = \"\";\n\t\tfor await (const event of this.fillDocumentSlotStream(\n\t\t\tdocumentId,\n\t\t\tslotId,\n\t\t\toptions,\n\t\t\tsignal,\n\t\t)) {\n\t\t\tif (event.event === \"text_chunk\") {\n\t\t\t\tcontent += event.data.delta;\n\t\t\t}\n\t\t}\n\t\treturn { documentId, slotId, content };\n\t}\n\n\t/**\n\t * Fills a single document slot by running its bound workflow (or an\n\t * explicitly provided one). Unlike {@link executeWorkflowStream}, this does\n\t * NOT create or persist a thread — the document slot is the only artifact.\n\t * Progress is streamed live but not persisted anywhere.\n\t */\n\tasync *fillDocumentSlotStream(\n\t\tdocumentId: string,\n\t\tslotId: string,\n\t\toptions?: {\n\t\t\tworkflowId?: string;\n\t\t\texecutionVariables?: Record<string, string>;\n\t\t\t/** Who initiated this fill; logged so unexpected fills are traceable. */\n\t\t\tinitiator?: SlotFillInitiator;\n\t\t},\n\t\tsignal?: AbortSignal,\n\t): AsyncGenerator<StreamEvent> {\n\t\tconst documentMemory = this.memoryModule.getDocumentMemory();\n\t\tif (!documentMemory) {\n\t\t\tthrow new Error(\"Document memory is not initialized\");\n\t\t}\n\n\t\tconst document = await documentMemory.getDocument(documentId);\n\t\tif (!document) {\n\t\t\tthrow new Error(`Document not found: ${documentId}`);\n\t\t}\n\n\t\tconst slot = document.slots?.find((s) => s.slotId === slotId);\n\t\tif (!slot) {\n\t\t\tthrow new Error(`Document slot not found: ${documentId}/${slotId}`);\n\t\t}\n\n\t\t// Resolve which workflow fills this slot (explicit override > binding).\n\t\tconst workflowId =\n\t\t\toptions?.workflowId ??\n\t\t\t(slot.binding?.type === \"WORKFLOW\" ? slot.binding.workflowId : undefined);\n\t\tif (!workflowId) {\n\t\t\tthrow new Error(\n\t\t\t\t`No workflow bound to slot ${documentId}/${slotId}; provide workflowId`,\n\t\t\t);\n\t\t}\n\t\tconst executionVariables =\n\t\t\toptions?.executionVariables ??\n\t\t\t(slot.binding?.type === \"WORKFLOW\"\n\t\t\t\t? slot.binding.executionVariables\n\t\t\t\t: undefined);\n\n\t\tyield {\n\t\t\tevent: \"document_id\",\n\t\t\tdata: { documentId, slotId },\n\t\t};\n\n\t\t// A slot may bind to either a user workflow or a workflow template.\n\t\tconst workflow = await this.getFillableWorkflow(workflowId);\n\t\tif (!workflow) {\n\t\t\tthrow new Error(`User workflow or template not found: ${workflowId}`);\n\t\t}\n\n\t\tconst { definition } = this.workflowVariableResolver.resolveForDocumentFill(\n\t\t\tworkflow,\n\t\t\texecutionVariables,\n\t\t);\n\t\tif (!definition) {\n\t\t\tthrow new Error(\n\t\t\t\t`Workflow ${workflowId} has no structured definition; cannot fill slot`,\n\t\t\t);\n\t\t}\n\n\t\t// Fills can fire long after their cause (boot catch-up, delayed tick), so\n\t\t// name the initiator here — \"the log stream suddenly shows a fill\" must be\n\t\t// answerable from this line alone.\n\t\tconst initiator = options?.initiator;\n\t\tloggers.agent.info(\"Document slot fill started\", {\n\t\t\tdocumentId,\n\t\t\tslotId,\n\t\t\tworkflowId,\n\t\t\tworkflowTitle: workflow.title,\n\t\t\tinitiatedBy: initiator?.type ?? \"unknown\",\n\t\t\t...(initiator?.type === \"schedule\"\n\t\t\t\t? {\n\t\t\t\t\t\tscheduleTrigger: initiator.trigger,\n\t\t\t\t\t\tscheduleRunId: initiator.runId,\n\t\t\t\t\t\tscheduledFor: new Date(initiator.scheduledFor).toISOString(),\n\t\t\t\t\t\tscheduleDelayMs: Date.now() - initiator.scheduledFor,\n\t\t\t\t\t}\n\t\t\t\t: {}),\n\t\t\t...(initiator?.type === \"manual\" && initiator.userId\n\t\t\t\t? { requestedBy: initiator.userId }\n\t\t\t\t: {}),\n\t\t});\n\n\t\tawait this.updateSlot(documentMemory, document, slotId, {\n\t\t\tstatus: \"running\",\n\t\t\terror: undefined,\n\t\t});\n\n\t\t// Ephemeral, non-persisted thread: carries threadId for A2A correlation\n\t\t// and task context, but is never written to the thread store.\n\t\tconst thread: ThreadObject = {\n\t\t\ttype: ThreadType.WORKFLOW,\n\t\t\tuserId: document.userId,\n\t\t\tthreadId: randomUUID(),\n\t\t\ttitle: workflow.title,\n\t\t\tworkflowId,\n\t\t\tmessages: [],\n\t\t};\n\n\t\tconst { finalContent, renderedBlocks, executionError } =\n\t\t\tyield* this.renderStructuredDefinition(\n\t\t\t\tdefinition,\n\t\t\t\tthread,\n\t\t\t\tworkflowId,\n\t\t\t\tsignal,\n\t\t\t);\n\n\t\tif (executionError || !finalContent) {\n\t\t\tawait this.updateSlot(documentMemory, document, slotId, {\n\t\t\t\tstatus: \"failed\",\n\t\t\t\terror: executionError?.message ?? \"No content produced\",\n\t\t\t});\n\t\t\tif (executionError) {\n\t\t\t\tthrow executionError;\n\t\t\t}\n\t\t\treturn;\n\t\t}\n\n\t\tconst fragment: DocumentFragment = {\n\t\t\tcontent: finalContent,\n\t\t\tblocks: renderedBlocks,\n\t\t\tsource: { type: \"WORKFLOW\", workflowId },\n\t\t\tresolvedAt: new Date().toISOString(),\n\t\t};\n\t\tawait this.updateSlot(documentMemory, document, slotId, {\n\t\t\tstatus: \"resolved\",\n\t\t\tfragment,\n\t\t\terror: undefined,\n\t\t});\n\t}\n\n\t/**\n\t * Atomically patches a single slot via the memory layer. Must NOT rebuild\n\t * the whole slots array from `document` (a snapshot taken when the fill\n\t * started): concurrent fills of other slots would clobber each other's\n\t * results. `updateDocumentSlot` targets only the matched slot and bumps\n\t * `version`/`updatedAt` in the same write.\n\t */\n\tprivate async updateSlot(\n\t\tdocumentMemory: NonNullable<ReturnType<MemoryModule[\"getDocumentMemory\"]>>,\n\t\tdocument: Document,\n\t\tslotId: string,\n\t\tpatch: Partial<DocumentSlot>,\n\t): Promise<void> {\n\t\tawait documentMemory.updateDocumentSlot(document.documentId, slotId, patch);\n\t}\n\n\t/**\n\t * Resolves a slot binding's `workflowId` to a runnable workflow, accepting\n\t * either a user workflow or a workflow template. User workflows take\n\t * precedence; falls back to a template with the same id.\n\t */\n\tprivate async getFillableWorkflow(\n\t\tworkflowId: string,\n\t): Promise<UserWorkflow | WorkflowTemplate | undefined> {\n\t\tconst userWorkflow = await this.userWorkflowService.getWorkflow(workflowId);\n\t\tif (userWorkflow) {\n\t\t\treturn userWorkflow;\n\t\t}\n\t\treturn this.memoryModule\n\t\t\t.getWorkflowTemplateMemory()\n\t\t\t.getTemplate(workflowId);\n\t}\n\n\tprivate async createWorkflowThread(\n\t\tworkflow: UserWorkflow,\n\t\tdisplayQuery: string,\n\t\tresolvedQuery: string,\n\t): Promise<ThreadObject> {\n\t\tconst threadMemory = this.memoryModule.getThreadMemory();\n\t\tconst threadId = randomUUID();\n\t\tconst title = displayQuery || workflow.title;\n\t\tconst metadata: ThreadMetadata = (await threadMemory?.createThread(\n\t\t\tThreadType.WORKFLOW,\n\t\t\tworkflow.userId,\n\t\t\tthreadId,\n\t\t\ttitle,\n\t\t\tworkflow.workflowId,\n\t\t)) || {\n\t\t\ttype: ThreadType.WORKFLOW,\n\t\t\tuserId: workflow.userId,\n\t\t\tthreadId,\n\t\t\ttitle,\n\t\t\tworkflowId: workflow.workflowId,\n\t\t};\n\n\t\tconst thread: ThreadObject = { ...metadata, messages: [] };\n\t\tawait appendTextMessageToThread(\n\t\t\tthis.memoryModule,\n\t\t\tthread,\n\t\t\tMessageRole.USER,\n\t\t\ttitle,\n\t\t\t{\n\t\t\t\tworkflowId: workflow.workflowId,\n\t\t\t\tworkflowRun: true,\n\t\t\t\tquery: resolvedQuery,\n\t\t\t},\n\t\t);\n\n\t\treturn thread;\n\t}\n}\n"]}