{"version":3,"sources":["/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-UKYQ3XWI.cjs","../../src/services/job-runner.service.ts"],"names":[],"mappings":"AAAA;AACE;AACF,wDAA6B;AAC7B;AACE;AACF,wDAA6B;AAC7B;AACA;ACyBA,IAAM,wBAAA,EAA0B,CAAC,GAAA,EAAQ,IAAA,EAAS,IAAO,CAAA;AAEzD,SAAS,KAAA,CAAM,EAAA,EAA2B;AACzC,EAAA,OAAO,IAAI,OAAA,CAAQ,CAAC,OAAA,EAAA,GAAY,UAAA,CAAW,OAAA,EAAS,EAAE,CAAC,CAAA;AACxD;AAEO,IAAM,iBAAA,YAAN,MAAM,kBAAiB;AAAA,EACrB;AAAA,EACA;AAAA,EACA;AAAA,iBAEA,OAAA,EAAS,EAAA;AAAA,kBACT,UAAA,EAA+B,CAAC,EAAA;AAAA,kBAChC,YAAA,kBAAc,IAAI,GAAA,CAAY,EAAA;AAAA,kBAC9B,cAAA,EAAgB,EAAA;AAAA,kBAChB,SAAA,kBAAW,IAAI,GAAA,CAAyB,EAAA;AAAA,EAEhD,WAAA,CAAY,OAAA,EAA4B;AACvC,IAAA,IAAA,CAAK,cAAA,EAAgB,iBAAA,CAAiB,oBAAA;AAAA,sBACrC,OAAA,2BAAS;AAAA,IACV,CAAA;AACA,IAAA,IAAA,CAAK,cAAA,mCAAgB,OAAA,6BAAS,eAAA,UAAiB,yBAAA;AAC/C,IAAA,IAAA,CAAK,WAAA,mCAAa,OAAA,6BAAS,YAAA,UAAc,KAAA;AAAA,EAC1C;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAOA,OAAe,oBAAA,CAAqB,QAAA,EAA2B;AAC9D,IAAA,GAAA,CAAI,SAAA,IAAa,KAAA,CAAA,EAAW;AAC3B,MAAA,OAAO,MAAA,CAAO,QAAA,CAAS,QAAQ,EAAA,GAAK,SAAA,GAAY,EAAA,EAAI,SAAA,EAAW,CAAA;AAAA,IAChE;AACA,IAAA,MAAM,IAAA,EAAM,OAAA,CAAQ,GAAA,CAAI,wBAAA;AACxB,IAAA,GAAA,CAAI,IAAA,IAAQ,KAAA,CAAA,EAAW;AACtB,MAAA,OAAO,CAAA;AAAA,IACR;AACA,IAAA,MAAM,OAAA,EAAS,MAAA,CAAO,QAAA,CAAS,GAAA,EAAK,EAAE,CAAA;AACtC,IAAA,GAAA,CAAI,MAAA,CAAO,QAAA,CAAS,MAAM,EAAA,GAAK,OAAA,GAAU,CAAA,EAAG;AAC3C,MAAA,OAAO,MAAA;AAAA,IACR;AACA,IAAA,yBAAA,CAAQ,KAAA,CAAM,IAAA;AAAA,MACb,CAAA,kCAAA,EAAqC,GAAG,CAAA,kCAAA;AAAA,IACzC,CAAA;AACA,IAAA,OAAO,CAAA;AAAA,EACR;AAAA,EAEA,MAAM,MAAA,CAAO,GAAA,EAA+B;AAC3C,IAAA,GAAA,CAAI,IAAA,CAAK,WAAA,CAAY,GAAA,CAAI,GAAA,CAAI,MAAM,CAAA,EAAG;AACrC,MAAA,yBAAA,CAAQ,KAAA,CAAM,IAAA,CAAK,CAAA,uBAAA,EAA0B,GAAA,CAAI,MAAM,CAAA,CAAA;AACP,MAAA;AACjD,IAAA;AAC+B,IAAA;AACP,IAAA;AACH,IAAA;AACjB,IAAA;AACU,MAAA;AACZ,IAAA;AACuB,MAAA;AACU,MAAA;AACnC,IAAA;AACD,EAAA;AAAA;AAG+C,EAAA;AAC1C,IAAA;AAC2C,IAAA;AACD,MAAA;AAC7C,IAAA;AACG,IAAA;AACqD,MAAA;AACvD,IAAA;AACyB,MAAA;AAC3B,IAAA;AACD,EAAA;AAEiD,EAAA;AAC7B,IAAA;AACf,IAAA;AAC6C,MAAA;AAChC,MAAA;AAC8B,MAAA;AAClB,QAAA;AACvB,QAAA;AACe,UAAA;AAC4B,UAAA;AAC/B,QAAA;AAC8B,UAAA;AACF,UAAA;AAC7B,UAAA;AACoC,YAAA;AACH,YAAA;AAC/C,UAAA;AACgC,UAAA;AACL,YAAA;AACpB,cAAA;AACwC,cAAA;AAC9C,YAAA;AACD,UAAA;AAC6C,UAAA;AACE,YAAA;AAC/C,UAAA;AACM,UAAA;AAC6C,6BAAA;AACnD,UAAA;AACD,QAAA;AACD,MAAA;AACkD,MAAA;AACjD,IAAA;AACY,MAAA;AACd,IAAA;AACD,EAAA;AAE+C,EAAA;AACN,IAAA;AACI,MAAA;AAC5C,IAAA;AACD,EAAA;AAEiC,EAAA;AACM,IAAA;AAChC,MAAA;AACkB,MAAA;AACxB,IAAA;AACgC,IAAA;AACL,MAAA;AACpB,QAAA;AACG,QAAA;AACR,MAAA;AACD,IAAA;AACF,EAAA;AAEwB,EAAA;AAClB,IAAA;AAC6B,IAAA;AACnB,IAAA;AAChB,EAAA;AACD;ADlC8D;AACA;AACA;AACA","file":"/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-UKYQ3XWI.cjs","sourcesContent":[null,"import { classifyJobError } from \"@/utils/job-error.js\";\nimport { loggers } from \"@/utils/logger.js\";\n\n/**\n * Reliability layer for scheduled jobs. Serializes LLM pressure through a\n * semaphore, skips overlapping runs of the same jobKey, retries retryable\n * errors with backoff, and pauses ALL dispatch during a rate-limit cooldown.\n *\n * Deliberately storage-free: run history is recorded by the caller\n * (SchedulerService) around submit().\n */\n\nexport interface Job {\n\t/** Overlap key. WORKFLOW: workflowId, slot refresh: `${documentId}:${slotId}`. */\n\tjobKey: string;\n\texecute: () => Promise<void>;\n}\n\nexport type JobOutcome =\n\t| { status: \"success\"; attempts: number }\n\t| { status: \"failed\"; attempts: number; error: string }\n\t| { status: \"skipped_overlap\"; attempts: 0 };\n\nexport interface JobRunnerOptions {\n\t/** Max jobs executing at once. Default: env SCHEDULER_MAX_CONCURRENT or 2. */\n\tmaxConcurrent?: number;\n\t/** Waits between attempts; attempts = length + 1. */\n\tretryDelaysMs?: number[];\n\t/** Global dispatch pause after a rate limit without Retry-After. */\n\tcooldownMs?: number;\n}\n\nconst DEFAULT_RETRY_DELAYS_MS = [30_000, 120_000, 480_000];\n\nfunction sleep(ms: number): Promise<void> {\n\treturn new Promise((resolve) => setTimeout(resolve, ms));\n}\n\nexport class JobRunnerService {\n\tprivate maxConcurrent: number;\n\tprivate retryDelaysMs: number[];\n\tprivate cooldownMs: number;\n\n\tprivate active = 0;\n\tprivate waitQueue: Array<() => void> = [];\n\tprivate runningKeys = new Set<string>();\n\tprivate cooldownUntil = 0;\n\tprivate inFlight = new Set<Promise<JobOutcome>>();\n\n\tconstructor(options?: JobRunnerOptions) {\n\t\tthis.maxConcurrent = JobRunnerService.resolveMaxConcurrent(\n\t\t\toptions?.maxConcurrent,\n\t\t);\n\t\tthis.retryDelaysMs = options?.retryDelaysMs ?? DEFAULT_RETRY_DELAYS_MS;\n\t\tthis.cooldownMs = options?.cooldownMs ?? 60_000;\n\t}\n\n\t/**\n\t * A non-finite or <1 maxConcurrent would make `active < maxConcurrent`\n\t * permanently false, deadlocking every submitted job with no error.\n\t * Guard both the explicit option and the env-var fallback against that.\n\t */\n\tprivate static resolveMaxConcurrent(explicit?: number): number {\n\t\tif (explicit !== undefined) {\n\t\t\treturn Number.isFinite(explicit) && explicit >= 1 ? explicit : 2;\n\t\t}\n\t\tconst raw = process.env.SCHEDULER_MAX_CONCURRENT;\n\t\tif (raw === undefined) {\n\t\t\treturn 2;\n\t\t}\n\t\tconst parsed = Number.parseInt(raw, 10);\n\t\tif (Number.isFinite(parsed) && parsed >= 1) {\n\t\t\treturn parsed;\n\t\t}\n\t\tloggers.agent.warn(\n\t\t\t`Invalid SCHEDULER_MAX_CONCURRENT=\"${raw}\"; falling back to maxConcurrent=2`,\n\t\t);\n\t\treturn 2;\n\t}\n\n\tasync submit(job: Job): Promise<JobOutcome> {\n\t\tif (this.runningKeys.has(job.jobKey)) {\n\t\t\tloggers.agent.warn(`Job overlap, skipping: ${job.jobKey}`);\n\t\t\treturn { status: \"skipped_overlap\", attempts: 0 };\n\t\t}\n\t\tthis.runningKeys.add(job.jobKey);\n\t\tconst run = this.run(job);\n\t\tthis.inFlight.add(run);\n\t\ttry {\n\t\t\treturn await run;\n\t\t} finally {\n\t\t\tthis.inFlight.delete(run);\n\t\t\tthis.runningKeys.delete(job.jobKey);\n\t\t}\n\t}\n\n\t/** Waits for in-flight jobs to settle (graceful shutdown). */\n\tasync drain(timeoutMs = 30_000): Promise<void> {\n\t\tlet timeoutHandle: NodeJS.Timeout | undefined;\n\t\tconst timeout = new Promise<void>((resolve) => {\n\t\t\ttimeoutHandle = setTimeout(resolve, timeoutMs);\n\t\t});\n\t\ttry {\n\t\t\tawait Promise.race([Promise.allSettled([...this.inFlight]), timeout]);\n\t\t} finally {\n\t\t\tclearTimeout(timeoutHandle);\n\t\t}\n\t}\n\n\tprivate async run(job: Job): Promise<JobOutcome> {\n\t\tawait this.acquire();\n\t\ttry {\n\t\t\tconst maxAttempts = this.retryDelaysMs.length + 1;\n\t\t\tlet lastError = \"\";\n\t\t\tfor (let attempt = 1; attempt <= maxAttempts; attempt++) {\n\t\t\t\tawait this.waitForCooldown();\n\t\t\t\ttry {\n\t\t\t\t\tawait job.execute();\n\t\t\t\t\treturn { status: \"success\", attempts: attempt };\n\t\t\t\t} catch (error) {\n\t\t\t\t\tconst classification = classifyJobError(error);\n\t\t\t\t\tlastError = error instanceof Error ? error.message : String(error);\n\t\t\t\t\tloggers.agent.error(\n\t\t\t\t\t\t`Job attempt ${attempt}/${maxAttempts} failed: ${job.jobKey}`,\n\t\t\t\t\t\t{ error: lastError, retryable: classification.retryable },\n\t\t\t\t\t);\n\t\t\t\t\tif (classification.rateLimited) {\n\t\t\t\t\t\tthis.cooldownUntil = Math.max(\n\t\t\t\t\t\t\tthis.cooldownUntil,\n\t\t\t\t\t\t\tDate.now() + (classification.retryAfterMs ?? this.cooldownMs),\n\t\t\t\t\t\t);\n\t\t\t\t\t}\n\t\t\t\t\tif (!classification.retryable || attempt === maxAttempts) {\n\t\t\t\t\t\treturn { status: \"failed\", attempts: attempt, error: lastError };\n\t\t\t\t\t}\n\t\t\t\t\tawait sleep(\n\t\t\t\t\t\tclassification.retryAfterMs ?? this.retryDelaysMs[attempt - 1],\n\t\t\t\t\t);\n\t\t\t\t}\n\t\t\t}\n\t\t\treturn { status: \"failed\", attempts: maxAttempts, error: lastError };\n\t\t} finally {\n\t\t\tthis.release();\n\t\t}\n\t}\n\n\tprivate async waitForCooldown(): Promise<void> {\n\t\twhile (Date.now() < this.cooldownUntil) {\n\t\t\tawait sleep(this.cooldownUntil - Date.now());\n\t\t}\n\t}\n\n\tprivate acquire(): Promise<void> {\n\t\tif (this.active < this.maxConcurrent) {\n\t\t\tthis.active++;\n\t\t\treturn Promise.resolve();\n\t\t}\n\t\treturn new Promise((resolve) => {\n\t\t\tthis.waitQueue.push(() => {\n\t\t\t\tthis.active++;\n\t\t\t\tresolve();\n\t\t\t});\n\t\t});\n\t}\n\n\tprivate release(): void {\n\t\tthis.active--;\n\t\tconst next = this.waitQueue.shift();\n\t\tif (next) next();\n\t}\n}\n"]}