{"version":3,"sources":["/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-I4NHSSD7.cjs","../../src/services/scheduler.service.ts"],"names":[],"mappings":"AAAA;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACA;ACPA,gCAA2B;AAC3B,yFAAyC;AAuBlC,IAAM,iBAAA,YAAN,MAAM,kBAAiB;AAAA,EACrB;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,iBACA,MAAA,kBAAoC,IAAI,GAAA,CAAI,EAAA;AAAA,kBAC5C,mBAAA,kBAA0C,IAAI,GAAA,CAAI,EAAA;AAAA,EAClD;AAAA,EACR,4BAAwB,iBAAA,EAAmB,IAAA;AAAA,EAC3C,6BAAwB,yBAAA,EAA2B,EAAA;AAAA;AAAA,kBAE3C,oBAAA,kBAA2C,IAAI,GAAA,CAAI,EAAA;AAAA,EAE3D,WAAA,CACC,mBAAA,EACA,wBAAA,EACA,SAAA,EACA,YAAA,EACC;AACD,IAAA,IAAA,CAAK,oBAAA,EAAsB,mBAAA;AAC3B,IAAA,IAAA,CAAK,yBAAA,EAA2B,wBAAA;AAChC,IAAA,IAAA,CAAK,UAAA,EAAY,SAAA;AACjB,IAAA,IAAA,CAAK,aAAA,EAAe,YAAA;AAAA,EACrB;AAAA,EAEA,MAAM,KAAA,CAAA,EAAuB;AAC5B,IAAA,MAAM,kBAAA,EAAoB,IAAA,CAAK,YAAA,CAAa,oBAAA,CAAqB,CAAA;AACjE,IAAA,GAAA,CAAI,iBAAA,EAAmB;AACtB,MAAA,MAAM,YAAA,EAAc,MAAM,iBAAA,CAAkB,mBAAA,CAAoB,CAAA;AAChE,MAAA,GAAA,CAAI,YAAA,EAAc,CAAA,EAAG;AACpB,QAAA,yBAAA,CAAQ,KAAA,CAAM,IAAA;AAAA,UACb,CAAA,OAAA,EAAU,WAAW,CAAA,sCAAA;AAAA,QACtB,CAAA;AAAA,MACD;AAAA,IACD;AAEA,IAAA,MAAM,gBAAA,EACL,MAAM,IAAA,CAAK,mBAAA,CAAoB,4BAAA,CAA6B,CAAA;AAC7D,IAAA,yBAAA,CAAQ,KAAA,CAAM,IAAA;AAAA,MACb,CAAA,wBAAA,EAA2B,eAAA,CAAgB,MAAM,CAAA,mBAAA;AAAA,IAClD,CAAA;AACA,IAAA,IAAA,CAAA,MAAW,SAAA,GAAY,eAAA,EAAiB;AAEvC,MAAA,MAAM,QAAA,EACL,QAAA,CAAS,UAAA,IAAc,KAAA,EAAA,GAAa,QAAA,CAAS,UAAA,GAAa,IAAA,CAAK,GAAA,CAAI,CAAA;AACpE,MAAA,MAAM,IAAA,CAAK,gBAAA,CAAiB,QAAQ,CAAA;AACpC,MAAA,GAAA,CAAI,OAAA,EAAS;AACZ,QAAA,KAAK,IAAA,CAAK,cAAA;AAAA,UACT,QAAA,CAAS,UAAA;AAAA,UACT,SAAA;AAAA,2BACA,QAAA,CAAS,SAAA,UAAa,IAAA,CAAK,GAAA,CAAI;AAAA,QAChC,CAAA;AAAA,MACD;AAAA,IACD;AAEA,IAAA,MAAM,IAAA,CAAK,wBAAA,CAAyB,CAAA;AACpC,IAAA,IAAA,CAAK,SAAA,CAAU,CAAA;AACf,IAAA,IAAA,CAAK,IAAA,CAAK,CAAA;AAAA,EACX;AAAA,EAEA,MAAM,IAAA,CAAA,EAAsB;AAC3B,IAAA,yBAAA,CAAQ,KAAA,CAAM,IAAA;AAAA,MACb,CAAA,6BAAA,EAAgC,IAAA,CAAK,KAAA,CAAM,IAAI,CAAA,QAAA;AAAA,IAChD,CAAA;AACA,IAAA,IAAA,CAAA,MAAW,CAAC,UAAA,EAAY,IAAI,EAAA,GAAK,IAAA,CAAK,KAAA,EAAO;AAC5C,MAAA,MAAM,IAAA,CAAK,IAAA,CAAK,CAAA;AAChB,MAAA,yBAAA,CAAQ,KAAA,CAAM,KAAA,CAAM,CAAA,wBAAA,EAA2B,UAAU,CAAA,CAAA;AAC1D,IAAA;AACiB,IAAA;AACG,IAAA;AACS,MAAA;AACX,MAAA;AAClB,IAAA;AAC2B,IAAA;AAC5B,EAAA;AAE8D,EAAA;AACrC,IAAA;AACvB,MAAA;AACD,IAAA;AACyC,IAAA;AACS,MAAA;AAClD,IAAA;AACuC,IAAA;AACxB,MAAA;AACoC,QAAA;AAClD,MAAA;AACA,MAAA;AACD,IAAA;AAEkB,IAAA;AACR,MAAA;AACW,MAAA;AACL,QAAA;AACkC,UAAA;AAChD,QAAA;AACuD,QAAA;AACxD,MAAA;AACA,MAAA;AACoB,QAAA;AACJ,QAAA;AAChB,MAAA;AACD,IAAA;AACwC,IAAA;AAER,IAAA;AACuB,IAAA;AACrC,MAAA;AACwB,MAAA;AACzC,IAAA;AACa,IAAA;AACsC,MAAA;AACpD,IAAA;AACD,EAAA;AAE4D,EAAA;AACrB,IAAA;AAC5B,IAAA;AACO,MAAA;AACY,MAAA;AAC6B,MAAA;AAC1D,IAAA;AAC0C,IAAA;AAC3C,EAAA;AAEgE,EAAA;AACd,IAAA;AACP,IAAA;AACL,MAAA;AACrC,IAAA;AACD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAeiB,EAAA;AAGS,IAAA;AAClB,IAAA;AAAyC,MAAA;AACC,MAAA;AACjD,IAAA;AACD,EAAA;AAMC,EAAA;AAEI,IAAA;AACyC,MAAA;AACjB,MAAA;AAG2B,MAAA;AACrD,QAAA;AACA,QAAA;AACiD,QAAA;AAC5B,QAAA;AACrB,MAAA;AAC0C,MAAA;AAC1C,QAAA;AACS,QAAA;AACD,QAAA;AACR,QAAA;AACA,QAAA;AACA,QAAA;AACQ,QAAA;AACE,QAAA;AACV,MAAA;AAE+C,MAAA;AACjC,MAAA;AAE0B,QAAA;AACU,QAAA;AACzC,UAAA;AACa,UAAA;AACX,UAAA;AACH,UAAA;AACP,QAAA;AACD,QAAA;AACD,MAAA;AAE4C,MAAA;AACnC,QAAA;AACa,QAAA;AACgC,UAAA;AACrD,QAAA;AACA,MAAA;AAQgC,MAAA;AACe,QAAA;AACE,QAAA;AAChB,QAAA;AAClB,UAAA;AACoC,YAAA;AAClD,UAAA;AACwC,UAAA;AACzC,QAAA;AACwC,MAAA;AACE,QAAA;AAC3C,MAAA;AAEkD,MAAA;AACjC,QAAA;AACK,QAAA;AACH,QAAA;AACmC,QAAA;AACrD,MAAA;AAEsD,MAAA;AACG,MAAA;AACxC,QAAA;AACN,QAAA;AAC8B,QAAA;AACzC,MAAA;AACc,IAAA;AACyC,MAAA;AACvD,QAAA;AACA,QAAA;AACA,MAAA;AACF,IAAA;AACD,EAAA;AAAA;AAGoD,EAAA;AACtB,IAAA;AACwB,IAAA;AACH,MAAA;AAC3C,IAAA;AAC4C,MAAA;AACnD,IAAA;AACD,EAAA;AAAA;AAGoD,EAAA;AACV,IAAA;AAC1C,EAAA;AAAA;AAGyB,EAAA;AACT,IAAA;AAChB,EAAA;AAE0B,EAAA;AACL,IAAA;AACH,IAAA;AACA,MAAA;AACC,MAAA;AAClB,IAAA;AACuB,oBAAA;AACxB,EAAA;AAEwD,EAAA;AACI,IAAA;AACL,IAAA;AACrD,MAAA;AACD,IAAA;AACuC,IAAA;AACL,IAAA;AACM,MAAA;AACxC,IAAA;AACc,IAAA;AACmC,MAAA;AACjD,IAAA;AACD,EAAA;AAEqB,EAAA;AACC,IAAA;AACsC,IAAA;AACxC,MAAA;AACwB,QAAA;AAER,QAAA;AACqB,QAAA;AACvD,MAAA;AACD,IAAA;AACD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAiBiB,EAAA;AAES,IAAA;AAClB,IAAA;AAAyC,MAAA;AACI,MAAA;AACpD,IAAA;AACD,EAAA;AAMC,EAAA;AAEI,IAAA;AACsC,MAAA;AACG,MAAA;AACvB,MAAA;AACM,MAAA;AAIqB,MAAA;AAC/C,QAAA;AACA,QAAA;AACiD,QAAA;AAC5B,QAAA;AACrB,MAAA;AAC0C,MAAA;AAC1C,QAAA;AACS,QAAA;AACD,QAAA;AACR,QAAA;AACA,QAAA;AACA,QAAA;AACQ,QAAA;AACE,QAAA;AACV,MAAA;AAEiD,MAAA;AACnC,MAAA;AACoC,QAAA;AACzC,UAAA;AACa,UAAA;AACX,UAAA;AACH,UAAA;AACP,QAAA;AACD,QAAA;AACD,MAAA;AAC6B,MAAA;AACwB,MAAA;AACF,QAAA;AACzC,UAAA;AACa,UAAA;AACX,UAAA;AAEP,UAAA;AAEH,QAAA;AACD,QAAA;AACD,MAAA;AAKkD,MAAA;AACtB,MAAA;AACwB,MAAA;AACjB,MAAA;AACf,QAAA;AAClB,UAAA;AACe,UAAA;AACK,UAAA;AACpB,QAAA;AACF,MAAA;AAE4B,MAAA;AAC4B,QAAA;AACL,QAAA;AACzC,UAAA;AACa,UAAA;AACX,UAAA;AACE,UAAA;AACZ,QAAA;AACD,QAAA;AACD,MAAA;AAE8C,MAAA;AACtB,MAAA;AACV,MAAA;AACmB,QAAA;AACa,UAAA;AACZ,YAAA;AACV,YAAA;AACgB,cAAA;AACnC,gBAAA;AACA,gBAAA;AACA,gBAAA;AACyC,kBAAA;AACzC,gBAAA;AACD,cAAA;AACD,YAAA;AACA,UAAA;AACiC,UAAA;AAM7B,YAAA;AAS4C,cAAA;AACD,cAAA;AACR,cAAA;AACpB,gBAAA;AAChB,kBAAA;AACQ,kBAAA;AACU,kBAAA;AACS,kBAAA;AAC3B,gBAAA;AACD,gBAAA;AACD,cAAA;AACqB,cAAA;AACpB,gBAAA;AACA,gBAAA;AACD,cAAA;AACe,YAAA;AACK,cAAA;AACA,cAAA;AACnB,gBAAA;AACA,gBAAA;AACA,gBAAA;AACA,cAAA;AACF,YAAA;AACuC,UAAA;AAQnC,YAAA;AACqC,cAAA;AAC/B,gBAAA;AACO,gBAAA;AACf,cAAA;AACc,YAAA;AACK,cAAA;AACnB,gBAAA;AACA,gBAAA;AACA,gBAAA;AACA,cAAA;AACF,YAAA;AACD,UAAA;AACiB,UAAA;AAChB,YAAA;AACgB,YAAA;AACE,YAAA;AAC2B,YAAA;AAC7C,UAAA;AACD,QAAA;AACF,MAAA;AAEsD,MAAA;AACP,MAAA;AACS,QAAA;AACxD,MAAA;AACkD,MAAA;AACP,QAAA;AACrB,QAAA;AACX,QAAA;AAGc,QAAA;AAExB,QAAA;AACA,MAAA;AACc,IAAA;AACK,MAAA;AACnB,QAAA;AACA,QAAA;AACA,MAAA;AACF,IAAA;AACD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAWY,EAAA;AAED,IAAA;AAEoB,IAAA;AAC/B,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAawB,EAAA;AACU,IAAA;AACsB,IAAA;AACrC,IAAA;AACuC,MAAA;AACzD,IAAA;AACoD,IAAA;AACb,IAAA;AACW,IAAA;AAEjB,IAAA;AACS,IAAA;AACT,IAAA;AACV,MAAA;AAC2B,QAAA;AAC1C,MAAA;AACmB,QAAA;AAC1B,MAAA;AACD,IAAA;AAEsC,IAAA;AACP,MAAA;AAChB,MAAA;AACb,QAAA;AACgD,QAAA;AAChD,MAAA;AACF,IAAA;AAEO,IAAA;AACN,MAAA;AAAA;AAAA;AAAA;AAI0D,MAAA;AAC1D,MAAA;AACA,MAAA;AACD,IAAA;AACD,EAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAYiB,EAAA;AACZ,IAAA;AACsC,MAAA;AACpB,MAAA;AAE6B,MAAA;AACpB,MAAA;AACuB,MAAA;AACpD,QAAA;AACD,MAAA;AAEiD,MAAA;AAClB,MAAA;AAMuB,MAAA;AACrB,MAAA;AAClB,QAAA;AACb,UAAA;AAC2C,UAAA;AAC5C,QAAA;AACA,QAAA;AACD,MAAA;AAE+C,MAAA;AAET,MAAA;AACG,MAAA;AACe,QAAA;AACd,QAAA;AAC1C,MAAA;AACe,IAAA;AACK,MAAA;AACnB,QAAA;AACA,QAAA;AACA,QAAA;AACA,MAAA;AACF,IAAA;AACD,EAAA;AACD;ADzIgE;AACA;AACA;AACA","file":"/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-I4NHSSD7.cjs","sourcesContent":[null,"import { randomUUID } from \"node:crypto\";\nimport cron, { type ScheduledTask } from \"node-cron\";\nimport type { MemoryModule } from \"@/modules/memory/memory.module.js\";\nimport type { Document, DocumentAutoRefresh } from \"@/types/document.js\";\nimport type { UserWorkflow } from \"@/types/memory.js\";\nimport type {\n\tScheduleRunExclusion,\n\tScheduleRunSlotResult,\n\tScheduleRunTargeting,\n\tScheduleTrigger,\n} from \"@/types/schedule.js\";\nimport { loggers } from \"@/utils/logger.js\";\nimport { runWithRequestContext } from \"@/utils/request-context.js\";\nimport type { JobRunnerService } from \"./job-runner.service.js\";\nimport type { UserWorkflowService } from \"./user-workflow.service.js\";\nimport type { WorkflowExecutionService } from \"./workflow-execution.service.js\";\n\n/**\n * Cron-based scheduler for user workflows plus one-shot document auto\n * refreshes. Triggering (node-cron / minute tick) is separated from\n * execution: every run goes through the JobRunner, which owns concurrency,\n * retries and the rate-limit cooldown. This service owns run history and\n * schedule state (nextRunAt, autoRefresh bookkeeping).\n */\nexport class SchedulerService {\n\tprivate userWorkflowService: UserWorkflowService;\n\tprivate workflowExecutionService: WorkflowExecutionService;\n\tprivate jobRunner: JobRunnerService;\n\tprivate memoryModule: MemoryModule;\n\tprivate tasks: Map<string, ScheduledTask> = new Map();\n\tprivate pendingAutoRefresh: Map<string, number> = new Map();\n\tprivate tickTimer?: ReturnType<typeof setInterval>;\n\tprivate static readonly TICK_INTERVAL_MS = 60_000;\n\tprivate static readonly MAX_CONSECUTIVE_FAILURES = 3;\n\t/** Consecutive failed cron runs per workflowId; reset on success. */\n\tprivate consecutiveFailures: Map<string, number> = new Map();\n\n\tconstructor(\n\t\tuserWorkflowService: UserWorkflowService,\n\t\tworkflowExecutionService: WorkflowExecutionService,\n\t\tjobRunner: JobRunnerService,\n\t\tmemoryModule: MemoryModule,\n\t) {\n\t\tthis.userWorkflowService = userWorkflowService;\n\t\tthis.workflowExecutionService = workflowExecutionService;\n\t\tthis.jobRunner = jobRunner;\n\t\tthis.memoryModule = memoryModule;\n\t}\n\n\tasync start(): Promise<void> {\n\t\tconst scheduleRunMemory = this.memoryModule.getScheduleRunMemory();\n\t\tif (scheduleRunMemory) {\n\t\t\tconst interrupted = await scheduleRunMemory.failInterruptedRuns();\n\t\t\tif (interrupted > 0) {\n\t\t\t\tloggers.agent.warn(\n\t\t\t\t\t`Marked ${interrupted} interrupted schedule run(s) as failed`,\n\t\t\t\t);\n\t\t\t}\n\t\t}\n\n\t\tconst activeWorkflows =\n\t\t\tawait this.userWorkflowService.listActiveScheduledWorkflows();\n\t\tloggers.agent.info(\n\t\t\t`Scheduler starting with ${activeWorkflows.length} active workflow(s)`,\n\t\t);\n\t\tfor (const workflow of activeWorkflows) {\n\t\t\t// Catch-up BEFORE scheduleWorkflow refreshes nextRunAt.\n\t\t\tconst overdue =\n\t\t\t\tworkflow.nextRunAt !== undefined && workflow.nextRunAt <= Date.now();\n\t\t\tawait this.scheduleWorkflow(workflow);\n\t\t\tif (overdue) {\n\t\t\t\tvoid this.runWorkflowJob(\n\t\t\t\t\tworkflow.workflowId,\n\t\t\t\t\t\"catchup\",\n\t\t\t\t\tworkflow.nextRunAt ?? Date.now(),\n\t\t\t\t);\n\t\t\t}\n\t\t}\n\n\t\tawait this.loadAutoRefreshDocuments();\n\t\tthis.startTick();\n\t\tthis.tick(); // 부팅 즉시 1회 — runAt이 이미 지난 문서의 catch-up\n\t}\n\n\tasync stop(): Promise<void> {\n\t\tloggers.agent.info(\n\t\t\t`Scheduler stopping, clearing ${this.tasks.size} task(s)`,\n\t\t);\n\t\tfor (const [workflowId, task] of this.tasks) {\n\t\t\tawait task.stop();\n\t\t\tloggers.agent.debug(`Stopped scheduled task: ${workflowId}`);\n\t\t}\n\t\tthis.tasks.clear();\n\t\tif (this.tickTimer) {\n\t\t\tclearInterval(this.tickTimer);\n\t\t\tthis.tickTimer = undefined;\n\t\t}\n\t\tawait this.jobRunner.drain();\n\t}\n\n\tasync scheduleWorkflow(workflow: UserWorkflow): Promise<void> {\n\t\tif (!workflow.schedule) {\n\t\t\treturn;\n\t\t}\n\t\tif (this.tasks.has(workflow.workflowId)) {\n\t\t\tawait this.unscheduleWorkflow(workflow.workflowId);\n\t\t}\n\t\tif (!cron.validate(workflow.schedule)) {\n\t\t\tloggers.agent.error(\n\t\t\t\t`Invalid cron expression for workflow ${workflow.workflowId}: ${workflow.schedule}`,\n\t\t\t);\n\t\t\treturn;\n\t\t}\n\n\t\tconst task = cron.schedule(\n\t\t\tworkflow.schedule,\n\t\t\tasync (_context) => {\n\t\t\t\tloggers.agent.info(\n\t\t\t\t\t`Cron triggered workflow: ${workflow.title} (${workflow.workflowId})`,\n\t\t\t\t);\n\t\t\t\tawait this.runWorkflowJob(workflow.workflowId, \"cron\", Date.now());\n\t\t\t},\n\t\t\t{\n\t\t\t\ttimezone: workflow.timezone,\n\t\t\t\tname: workflow.workflowId,\n\t\t\t},\n\t\t);\n\t\tthis.tasks.set(workflow.workflowId, task);\n\n\t\tconst nextRun = task.getNextRun();\n\t\tawait this.userWorkflowService.updateWorkflow(workflow.workflowId, {\n\t\t\tuserId: workflow.userId,\n\t\t\tnextRunAt: nextRun ? nextRun.getTime() : undefined,\n\t\t});\n\t\tloggers.agent.info(\n\t\t\t`Scheduled workflow: ${workflow.title} (${workflow.workflowId}) with cron \"${workflow.schedule}\"${workflow.timezone ? ` [${workflow.timezone}]` : \"\"}`,\n\t\t);\n\t}\n\n\tasync unscheduleWorkflow(workflowId: string): Promise<void> {\n\t\tconst task = this.tasks.get(workflowId);\n\t\tif (task) {\n\t\t\tawait task.stop();\n\t\t\tthis.tasks.delete(workflowId);\n\t\t\tloggers.agent.debug(`Unscheduled workflow: ${workflowId}`);\n\t\t}\n\t\tthis.consecutiveFailures.delete(workflowId);\n\t}\n\n\tasync rescheduleWorkflow(workflow: UserWorkflow): Promise<void> {\n\t\tawait this.unscheduleWorkflow(workflow.workflowId);\n\t\tif (workflow.active && workflow.schedule) {\n\t\t\tawait this.scheduleWorkflow(workflow);\n\t\t}\n\t}\n\n\t/**\n\t * Executes one scheduled workflow run through the JobRunner and records\n\t * it in schedule_runs. Public for tests and manual triggering.\n\t *\n\t * Never rejects: execution errors are absorbed by the JobRunner, and\n\t * bookkeeping (memory) errors are caught and logged here so that\n\t * fire-and-forget callers (boot catch-up) cannot crash the process\n\t * with an unhandled rejection.\n\t */\n\tasync runWorkflowJob(\n\t\tworkflowId: string,\n\t\ttrigger: ScheduleTrigger,\n\t\tscheduledFor: number,\n\t): Promise<void> {\n\t\t// The runId doubles as the correlation id so every log line of this\n\t\t// run carries it, the same way requestId does for HTTP requests.\n\t\tconst runId = randomUUID();\n\t\treturn runWithRequestContext({ requestId: runId }, () =>\n\t\t\tthis.runWorkflowJobWithRunId(runId, workflowId, trigger, scheduledFor),\n\t\t);\n\t}\n\n\tprivate async runWorkflowJobWithRunId(\n\t\trunId: string,\n\t\tworkflowId: string,\n\t\ttrigger: ScheduleTrigger,\n\t\tscheduledFor: number,\n\t): Promise<void> {\n\t\ttry {\n\t\t\tconst scheduleRunMemory = this.memoryModule.getScheduleRunMemory();\n\t\t\tconst startedAt = Date.now();\n\t\t\t// Catch-up runs fire at boot, far from the cron time they stand in\n\t\t\t// for — record the planned time so the gap is explainable from logs.\n\t\t\tloggers.agent.info(\"Scheduled workflow run starting\", {\n\t\t\t\tworkflowId,\n\t\t\t\ttrigger,\n\t\t\t\tscheduledFor: new Date(scheduledFor).toISOString(),\n\t\t\t\tdelayMs: startedAt - scheduledFor,\n\t\t\t});\n\t\t\tawait scheduleRunMemory?.createScheduleRun({\n\t\t\t\trunId,\n\t\t\t\tjobType: \"WORKFLOW\",\n\t\t\t\tjobKey: workflowId,\n\t\t\t\ttrigger,\n\t\t\t\tscheduledFor,\n\t\t\t\tstartedAt,\n\t\t\t\tstatus: \"running\",\n\t\t\t\tattempts: 0,\n\t\t\t});\n\n\t\t\tconst workflow = await this.userWorkflowService.getWorkflow(workflowId);\n\t\t\tif (!workflow) {\n\t\t\t\t// Deleted since scheduling: stop repeating a doomed job.\n\t\t\t\tawait this.unscheduleWorkflow(workflowId);\n\t\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, {\n\t\t\t\t\tstatus: \"failed\",\n\t\t\t\t\tfinishedAt: Date.now(),\n\t\t\t\t\tattempts: 1,\n\t\t\t\t\terror: \"Workflow not found; unscheduled\",\n\t\t\t\t});\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tconst outcome = await this.jobRunner.submit({\n\t\t\t\tjobKey: workflowId,\n\t\t\t\texecute: async () => {\n\t\t\t\t\tawait this.workflowExecutionService.executeWorkflow(workflowId);\n\t\t\t\t},\n\t\t\t});\n\n\t\t\t// Spec §8: a workflow that fails deterministically (broken definition,\n\t\t\t// etc.) would otherwise fail every cron period forever. Tracked\n\t\t\t// in-memory only (deliberate, non-destructive): a process restart\n\t\t\t// re-arms the schedule for one more attempt cycle rather than\n\t\t\t// permanently stranding a workflow whose definition was fixed but\n\t\t\t// whose persisted counter was never cleared.\n\t\t\tif (outcome.status === \"failed\") {\n\t\t\t\tconst failures = (this.consecutiveFailures.get(workflowId) ?? 0) + 1;\n\t\t\t\tthis.consecutiveFailures.set(workflowId, failures);\n\t\t\t\tif (failures >= SchedulerService.MAX_CONSECUTIVE_FAILURES) {\n\t\t\t\t\tloggers.agent.warn(\n\t\t\t\t\t\t`Auto-unscheduled workflow ${workflowId} after ${failures} consecutive failures`,\n\t\t\t\t\t);\n\t\t\t\t\tawait this.unscheduleWorkflow(workflowId);\n\t\t\t\t}\n\t\t\t} else if (outcome.status === \"success\") {\n\t\t\t\tthis.consecutiveFailures.delete(workflowId);\n\t\t\t}\n\n\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, {\n\t\t\t\tstatus: outcome.status,\n\t\t\t\tfinishedAt: Date.now(),\n\t\t\t\tattempts: outcome.attempts,\n\t\t\t\terror: outcome.status === \"failed\" ? outcome.error : undefined,\n\t\t\t});\n\n\t\t\tconst nextRun = this.tasks.get(workflowId)?.getNextRun();\n\t\t\tawait this.userWorkflowService.updateWorkflow(workflowId, {\n\t\t\t\tuserId: workflow.userId,\n\t\t\t\tlastRunAt: startedAt,\n\t\t\t\tnextRunAt: nextRun ? nextRun.getTime() : undefined,\n\t\t\t});\n\t\t} catch (error) {\n\t\t\tloggers.agent.error(\"Scheduled run bookkeeping failed\", {\n\t\t\t\tworkflowId,\n\t\t\t\terror,\n\t\t\t});\n\t\t}\n\t}\n\n\t/** Reflects a created/updated document in the pending auto-refresh list. */\n\tnotifyDocumentAutoRefresh(document: Document): void {\n\t\tconst autoRefresh = document.autoRefresh;\n\t\tif (autoRefresh?.active && !autoRefresh.completedAt) {\n\t\t\tthis.pendingAutoRefresh.set(document.documentId, autoRefresh.runAt);\n\t\t} else {\n\t\t\tthis.pendingAutoRefresh.delete(document.documentId);\n\t\t}\n\t}\n\n\t/** Drops a (deleted) document from the pending list. */\n\tremoveDocumentAutoRefresh(documentId: string): void {\n\t\tthis.pendingAutoRefresh.delete(documentId);\n\t}\n\n\t/** Exposed for tests: starts the minute tick without full start(). */\n\tstartTickForTest(): void {\n\t\tthis.startTick();\n\t}\n\n\tprivate startTick(): void {\n\t\tif (this.tickTimer) return;\n\t\tthis.tickTimer = setInterval(\n\t\t\t() => this.tick(),\n\t\t\tSchedulerService.TICK_INTERVAL_MS,\n\t\t);\n\t\tthis.tickTimer.unref?.();\n\t}\n\n\tprivate async loadAutoRefreshDocuments(): Promise<void> {\n\t\tconst documentMemory = this.memoryModule.getDocumentMemory();\n\t\tif (!documentMemory?.listAutoRefreshPendingDocuments) {\n\t\t\treturn;\n\t\t}\n\t\tconst documents = await documentMemory.listAutoRefreshPendingDocuments();\n\t\tfor (const document of documents) {\n\t\t\tthis.notifyDocumentAutoRefresh(document);\n\t\t}\n\t\tloggers.agent.info(\n\t\t\t`Scheduler loaded ${this.pendingAutoRefresh.size} pending auto-refresh document(s)`,\n\t\t);\n\t}\n\n\tprivate tick(): void {\n\t\tconst now = Date.now();\n\t\tfor (const [documentId, runAt] of this.pendingAutoRefresh) {\n\t\t\tif (runAt <= now) {\n\t\t\t\tthis.pendingAutoRefresh.delete(documentId);\n\t\t\t\tconst trigger: ScheduleTrigger =\n\t\t\t\t\trunAt <= now - SchedulerService.TICK_INTERVAL_MS ? \"catchup\" : \"once\";\n\t\t\t\tvoid this.runAutoRefreshJob(documentId, trigger, runAt);\n\t\t\t}\n\t\t}\n\t}\n\n\t/**\n\t * Expands a document auto refresh into per-slot jobs (the JobRunner\n\t * throttles them), accumulates doneSlotIds, and completes the refresh\n\t * only when every target slot succeeded. Failed slots stay pending so\n\t * the next boot catch-up retries ONLY them.\n\t *\n\t * Never rejects: slot execution failures are absorbed by the JobRunner\n\t * (submit() never rejects), and bookkeeping (memory) errors are caught\n\t * and logged here so that fire-and-forget callers (tick) cannot crash\n\t * the process with an unhandled rejection.\n\t */\n\tasync runAutoRefreshJob(\n\t\tdocumentId: string,\n\t\ttrigger: ScheduleTrigger,\n\t\tscheduledFor: number,\n\t): Promise<void> {\n\t\t// See runWorkflowJob: runId doubles as the correlation id.\n\t\tconst runId = randomUUID();\n\t\treturn runWithRequestContext({ requestId: runId }, () =>\n\t\t\tthis.runAutoRefreshJobWithRunId(runId, documentId, trigger, scheduledFor),\n\t\t);\n\t}\n\n\tprivate async runAutoRefreshJobWithRunId(\n\t\trunId: string,\n\t\tdocumentId: string,\n\t\ttrigger: ScheduleTrigger,\n\t\tscheduledFor: number,\n\t): Promise<void> {\n\t\ttry {\n\t\t\tconst documentMemory = this.memoryModule.getDocumentMemory();\n\t\t\tconst scheduleRunMemory = this.memoryModule.getScheduleRunMemory();\n\t\t\tif (!documentMemory) return;\n\t\t\tconst startedAt = Date.now();\n\t\t\t// \"once\" fires within a tick of runAt; \"catchup\" is a boot/tick\n\t\t\t// recovery of an already-passed runAt, so it starts at an arbitrary\n\t\t\t// time — log the planned time and the gap to make that visible.\n\t\t\tloggers.agent.info(\"Auto-refresh run starting\", {\n\t\t\t\tdocumentId,\n\t\t\t\ttrigger,\n\t\t\t\tscheduledFor: new Date(scheduledFor).toISOString(),\n\t\t\t\tdelayMs: startedAt - scheduledFor,\n\t\t\t});\n\t\t\tawait scheduleRunMemory?.createScheduleRun({\n\t\t\t\trunId,\n\t\t\t\tjobType: \"SLOT_REFRESH\",\n\t\t\t\tjobKey: documentId,\n\t\t\t\ttrigger,\n\t\t\t\tscheduledFor,\n\t\t\t\tstartedAt,\n\t\t\t\tstatus: \"running\",\n\t\t\t\tattempts: 0,\n\t\t\t});\n\n\t\t\tconst document = await documentMemory.getDocument(documentId);\n\t\t\tif (!document) {\n\t\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, {\n\t\t\t\t\tstatus: \"failed\",\n\t\t\t\t\tfinishedAt: Date.now(),\n\t\t\t\t\tattempts: 1,\n\t\t\t\t\terror: \"Document not found\",\n\t\t\t\t});\n\t\t\t\treturn;\n\t\t\t}\n\t\t\tconst autoRefresh = document.autoRefresh;\n\t\t\tif (!autoRefresh?.active || autoRefresh.completedAt) {\n\t\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, {\n\t\t\t\t\tstatus: \"success\",\n\t\t\t\t\tfinishedAt: Date.now(),\n\t\t\t\t\tattempts: 0,\n\t\t\t\t\tnoopReason: autoRefresh?.completedAt\n\t\t\t\t\t\t? \"auto_refresh_completed\"\n\t\t\t\t\t\t: \"auto_refresh_inactive\",\n\t\t\t\t});\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\t// Persisted BEFORE the first slot job is submitted: a run killed by a\n\t\t\t// restart mid-way still shows what it set out to cover, and a run\n\t\t\t// covering fewer slots than the document has stops being silent.\n\t\t\tconst targeting = this.deriveAutoRefreshTargeting(document, autoRefresh);\n\t\t\tconst targetIds = targeting.targetSlotIds;\n\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, { targeting });\n\t\t\tif (targeting.excluded.length > 0) {\n\t\t\t\tloggers.agent.info(\"Auto-refresh skipped some document slots\", {\n\t\t\t\t\tdocumentId,\n\t\t\t\t\ttargetSlotIds: targetIds,\n\t\t\t\t\texcluded: targeting.excluded,\n\t\t\t\t});\n\t\t\t}\n\n\t\t\tif (targetIds.length === 0) {\n\t\t\t\tawait documentMemory.completeAutoRefresh?.(documentId, Date.now());\n\t\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, {\n\t\t\t\t\tstatus: \"success\",\n\t\t\t\t\tfinishedAt: Date.now(),\n\t\t\t\t\tattempts: 0,\n\t\t\t\t\tnoopReason: \"no_pending_slots\",\n\t\t\t\t});\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tconst slotResults: ScheduleRunSlotResult[] = [];\n\t\t\tlet doneMarkingFailed = false;\n\t\t\tawait Promise.all(\n\t\t\t\ttargetIds.map(async (slotId) => {\n\t\t\t\t\tconst outcome = await this.jobRunner.submit({\n\t\t\t\t\t\tjobKey: `${documentId}:${slotId}`,\n\t\t\t\t\t\texecute: async () => {\n\t\t\t\t\t\t\tawait this.workflowExecutionService.fillDocumentSlot(\n\t\t\t\t\t\t\t\tdocumentId,\n\t\t\t\t\t\t\t\tslotId,\n\t\t\t\t\t\t\t\t{\n\t\t\t\t\t\t\t\t\tinitiator: { type: \"schedule\", trigger, scheduledFor, runId },\n\t\t\t\t\t\t\t\t},\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\tif (outcome.status === \"success\") {\n\t\t\t\t\t\t// A rejection here must not fail-fast the Promise.all and\n\t\t\t\t\t\t// skip the aggregate run bookkeeping below. The fill itself\n\t\t\t\t\t\t// succeeded, so the slot still counts as success; only\n\t\t\t\t\t\t// completion is withheld so the next catch-up redoes the\n\t\t\t\t\t\t// (idempotent) done-marking.\n\t\t\t\t\t\ttry {\n\t\t\t\t\t\t\t// The fill resolving is not proof the slot resolved: a\n\t\t\t\t\t\t\t// \"no content produced\" fill persists status:\"failed\" on\n\t\t\t\t\t\t\t// the slot and returns without throwing. Re-read the\n\t\t\t\t\t\t\t// document and require the fresh persisted status to be\n\t\t\t\t\t\t\t// \"resolved\" before ledgering the slot as done —\n\t\t\t\t\t\t\t// otherwise record it as a failed slot so the run\n\t\t\t\t\t\t\t// aggregates to failed, completion is withheld, and the\n\t\t\t\t\t\t\t// boot catch-up retries it.\n\t\t\t\t\t\t\tconst fresh = await documentMemory.getDocument(documentId);\n\t\t\t\t\t\t\tconst freshSlot = fresh?.slots?.find((s) => s.slotId === slotId);\n\t\t\t\t\t\t\tif (freshSlot?.status !== \"resolved\") {\n\t\t\t\t\t\t\t\tslotResults.push({\n\t\t\t\t\t\t\t\t\tslotId,\n\t\t\t\t\t\t\t\t\tstatus: \"failed\",\n\t\t\t\t\t\t\t\t\tattempts: outcome.attempts,\n\t\t\t\t\t\t\t\t\terror: freshSlot?.error ?? \"No content produced\",\n\t\t\t\t\t\t\t\t});\n\t\t\t\t\t\t\t\treturn;\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t\tawait documentMemory.markAutoRefreshSlotDone?.(\n\t\t\t\t\t\t\t\tdocumentId,\n\t\t\t\t\t\t\t\tslotId,\n\t\t\t\t\t\t\t);\n\t\t\t\t\t\t} catch (error) {\n\t\t\t\t\t\t\tdoneMarkingFailed = true;\n\t\t\t\t\t\t\tloggers.agent.error(\"Auto-refresh slot bookkeeping failed\", {\n\t\t\t\t\t\t\t\tdocumentId,\n\t\t\t\t\t\t\t\tslotId,\n\t\t\t\t\t\t\t\terror,\n\t\t\t\t\t\t\t});\n\t\t\t\t\t\t}\n\t\t\t\t\t} else if (outcome.status === \"failed\") {\n\t\t\t\t\t\t// fillDocumentSlotStream throws for several pre-execution\n\t\t\t\t\t\t// cases (document/slot/binding/workflow/definition gone)\n\t\t\t\t\t\t// BEFORE ever writing status:\"failed\" to the slot itself —\n\t\t\t\t\t\t// only post-start execution errors do that. Without this,\n\t\t\t\t\t\t// the frontend badge (derived solely from slot.status)\n\t\t\t\t\t\t// stays stuck \"in progress\" forever. Idempotent with the\n\t\t\t\t\t\t// execution-failure path, which already sets this status.\n\t\t\t\t\t\ttry {\n\t\t\t\t\t\t\tawait documentMemory.updateDocumentSlot(documentId, slotId, {\n\t\t\t\t\t\t\t\tstatus: \"failed\",\n\t\t\t\t\t\t\t\terror: outcome.error,\n\t\t\t\t\t\t\t});\n\t\t\t\t\t\t} catch (error) {\n\t\t\t\t\t\t\tloggers.agent.error(\"Auto-refresh slot bookkeeping failed\", {\n\t\t\t\t\t\t\t\tdocumentId,\n\t\t\t\t\t\t\t\tslotId,\n\t\t\t\t\t\t\t\terror,\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\tslotResults.push({\n\t\t\t\t\t\tslotId,\n\t\t\t\t\t\tstatus: outcome.status,\n\t\t\t\t\t\tattempts: outcome.attempts,\n\t\t\t\t\t\terror: outcome.status === \"failed\" ? outcome.error : undefined,\n\t\t\t\t\t});\n\t\t\t\t}),\n\t\t\t);\n\n\t\t\tconst failed = slotResults.filter((r) => r.status !== \"success\");\n\t\t\tif (failed.length === 0 && !doneMarkingFailed) {\n\t\t\t\tawait documentMemory.completeAutoRefresh?.(documentId, Date.now());\n\t\t\t}\n\t\t\tawait scheduleRunMemory?.updateScheduleRun(runId, {\n\t\t\t\tstatus: failed.length === 0 ? \"success\" : \"failed\",\n\t\t\t\tfinishedAt: Date.now(),\n\t\t\t\tattempts: 1,\n\t\t\t\terror:\n\t\t\t\t\tfailed.length > 0\n\t\t\t\t\t\t? `${failed.length}/${slotResults.length} slot(s) failed`\n\t\t\t\t\t\t: undefined,\n\t\t\t\tslotResults,\n\t\t\t});\n\t\t} catch (error) {\n\t\t\tloggers.agent.error(\"Auto-refresh run bookkeeping failed\", {\n\t\t\t\tdocumentId,\n\t\t\t\terror,\n\t\t\t});\n\t\t}\n\t}\n\n\t/**\n\t * Candidate slot ids for a document's auto-refresh: an explicit allowlist,\n\t * or (default) every slot with a binding. Shared by\n\t * {@link deriveAutoRefreshTargeting} and {@link reconcileManualSlotFill} so\n\t * the two stay in sync.\n\t */\n\tprivate getAutoRefreshTargetSlotIds(\n\t\tdocument: Document,\n\t\tautoRefresh: DocumentAutoRefresh,\n\t): string[] {\n\t\tconst boundSlotIds = (document.slots ?? [])\n\t\t\t.filter((slot) => slot.binding)\n\t\t\t.map((slot) => slot.slotId);\n\t\treturn autoRefresh.slotIds ?? boundSlotIds;\n\t}\n\n\t/**\n\t * Expands {@link getAutoRefreshTargetSlotIds} into the derivation recorded\n\t * on the run: what the document offered, what this run will actually\n\t * submit, and every slot dropped along the way with its reason. Three\n\t * different filters can shrink the target set (no binding, an explicit\n\t * allowlist, the done ledger) and none of them left a trace before — a\n\t * 3-of-6 run and a genuine 3-slot run looked identical in storage.\n\t */\n\tprivate deriveAutoRefreshTargeting(\n\t\tdocument: Document,\n\t\tautoRefresh: DocumentAutoRefresh,\n\t): ScheduleRunTargeting {\n\t\tconst slots = document.slots ?? [];\n\t\tconst documentSlotIds = slots.map((slot) => slot.slotId);\n\t\tconst bound = new Set(\n\t\t\tslots.filter((slot) => slot.binding).map((slot) => slot.slotId),\n\t\t);\n\t\tconst candidates = this.getAutoRefreshTargetSlotIds(document, autoRefresh);\n\t\tconst candidateSet = new Set(candidates);\n\t\tconst done = new Set(autoRefresh.doneSlotIds ?? []);\n\n\t\tconst targetSlotIds: string[] = [];\n\t\tconst excluded: ScheduleRunExclusion[] = [];\n\t\tfor (const slotId of candidates) {\n\t\t\tif (done.has(slotId)) {\n\t\t\t\texcluded.push({ slotId, reason: \"already_done\" });\n\t\t\t} else {\n\t\t\t\ttargetSlotIds.push(slotId);\n\t\t\t}\n\t\t}\n\t\t// Slots the document carries that never became candidates at all.\n\t\tfor (const slotId of documentSlotIds) {\n\t\t\tif (candidateSet.has(slotId)) continue;\n\t\t\texcluded.push({\n\t\t\t\tslotId,\n\t\t\t\treason: bound.has(slotId) ? \"not_in_slot_ids\" : \"no_binding\",\n\t\t\t});\n\t\t}\n\n\t\treturn {\n\t\t\tdocumentSlotIds,\n\t\t\t// Omitted rather than set to undefined: mongoose stores an explicit\n\t\t\t// undefined as null, which does not match the optional string[] type\n\t\t\t// that readers of the stored run see.\n\t\t\t...(autoRefresh.slotIds ? { requestedSlotIds: autoRefresh.slotIds } : {}),\n\t\t\ttargetSlotIds,\n\t\t\texcluded,\n\t\t};\n\t}\n\n\t/**\n\t * Reconciles a successful MANUAL slot fill into the auto-refresh ledger.\n\t * If the slot is a target of a pending (active, incomplete) autoRefresh,\n\t * mark it done; if that completes every target, stamp completedAt and drop\n\t * the pending entry. Idempotent; never throws (bookkeeping must not fail\n\t * the fill request).\n\t */\n\tasync reconcileManualSlotFill(\n\t\tdocumentId: string,\n\t\tslotId: string,\n\t): Promise<void> {\n\t\ttry {\n\t\t\tconst documentMemory = this.memoryModule.getDocumentMemory();\n\t\t\tif (!documentMemory) return;\n\n\t\t\tconst document = await documentMemory.getDocument(documentId);\n\t\t\tconst autoRefresh = document?.autoRefresh;\n\t\t\tif (!document || !autoRefresh?.active || autoRefresh.completedAt) {\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tconst targets = this.getAutoRefreshTargetSlotIds(document, autoRefresh);\n\t\t\tif (!targets.includes(slotId)) return;\n\n\t\t\t// The fill call resolving is not proof the slot resolved: a \"no\n\t\t\t// content produced\" fill persists status:\"failed\" on the slot and\n\t\t\t// returns without throwing. Only a slot whose persisted status is\n\t\t\t// \"resolved\" may be ledgered as done.\n\t\t\tconst slot = document.slots?.find((s) => s.slotId === slotId);\n\t\t\tif (slot?.status !== \"resolved\") {\n\t\t\t\tloggers.agent.debug(\n\t\t\t\t\t\"Manual fill reconciliation skipped: slot not resolved\",\n\t\t\t\t\t{ documentId, slotId, status: slot?.status },\n\t\t\t\t);\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tawait documentMemory.markAutoRefreshSlotDone?.(documentId, slotId);\n\n\t\t\tconst done = new Set([...(autoRefresh.doneSlotIds ?? []), slotId]);\n\t\t\tif (targets.every((id) => done.has(id))) {\n\t\t\t\tawait documentMemory.completeAutoRefresh?.(documentId, Date.now());\n\t\t\t\tthis.removeDocumentAutoRefresh(documentId);\n\t\t\t}\n\t\t} catch (error) {\n\t\t\tloggers.agent.error(\"Manual fill auto-refresh reconciliation failed\", {\n\t\t\t\tdocumentId,\n\t\t\t\tslotId,\n\t\t\t\terror,\n\t\t\t});\n\t\t}\n\t}\n}\n"]}