{"version":3,"sources":["/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-DSBE2JX2.cjs","../../src/utils/sse-stream.ts"],"names":[],"mappings":"AAAA;AACE;AACF,wDAA6B;AAC7B;AACA;ACiBA,MAAA,SAAsB,iBAAA,CACrB,GAAA,EACA,GAAA,EACA,OAAA,EACgB;AAChB,EAAA,GAAA,CAAI,SAAA,CAAU,GAAA,EAAK;AAAA,IAClB,cAAA,EAAgB,mBAAA;AAAA,IAChB,eAAA,EAAiB,UAAA;AAAA,IACjB,UAAA,EAAY,YAAA;AAAA,IACZ,mBAAA,EAAqB;AAAA;AAAA,EACtB,CAAC,CAAA;AACD,EAAA,GAAA,CAAI,YAAA,CAAa,CAAA;AACjB,EAAA,GAAA,CAAI,KAAA,CAAM,SAAS,CAAA;AAEnB,EAAA,MAAM,kBAAA,EAAoB,WAAA,CAAY,CAAA,EAAA,GAAM;AAC3C,IAAA,GAAA,CAAI,KAAA,CAAM,gBAAgB,CAAA;AAAA,EAC3B,CAAA,EAAG,GAAK,CAAA;AAER,EAAA,MAAM,gBAAA,EAAkB,IAAI,eAAA,CAAgB,CAAA;AAC5C,EAAA,IAAI,eAAA;AACJ,EAAA,GAAA,CAAI,EAAA,CAAG,OAAA,EAAS,CAAA,EAAA,GAAM;AACrB,IAAA,eAAA,CAAgB,KAAA,CAAM,CAAA;AACtB,IAAA,yBAAA,CAAQ,YAAA,CAAa,IAAA,CAAK,CAAA,EAAA;AACf,MAAA;AACM,MAAA;AACL,MAAA;AACX,IAAA;AACD,EAAA;AAEG,EAAA;AACkB,IAAA;AACK,IAAA;AACL,MAAA;AAEA,MAAA;AACD,QAAA;AACG,wBAAA;AAEf,MAAA;AAIQ,QAAA;AACf,MAAA;AAEI,MAAA;AACkB,QAAA;AAAgC,MAAA;AAAK;AAAA;AAC3D,MAAA;AACD,IAAA;AACqB,IAAA;AACN,MAAA;AACf,IAAA;AACwB,EAAA;AAEL,IAAA;AAI4B,IAAA;AAGpB,IAAA;AACV,MAAA;AACN,MAAA;AACV,MAAA;AACO,MAAA;AACI,MAAA;AACX,IAAA;AACG,IAAA;AACH,MAAA;AAAwC,MAAA;AAA0B;AAAA;AACnE,IAAA;AACC,EAAA;AACa,IAAA;AACN,IAAA;AACT,EAAA;AACD;AD3B+B;AACA;AACA;AACA","file":"/Users/shyun/comcom/ain-enterprise/ain-adk/dist/cjs/chunk-DSBE2JX2.cjs","sourcesContent":[null,"import type { Request, Response } from \"express\";\nimport type { StreamEvent } from \"@/types/stream\";\nimport { loggers } from \"./logger\";\n\nexport type SSEStreamOptions = {\n\tlogLabel: string;\n\tuserId: string;\n\tlogContext?: Record<string, unknown>;\n\tonThreadId?: (threadId: string) => void;\n\tonThinkingProcess?: (\n\t\tthreadId: string,\n\t\tdata: Extract<StreamEvent, { event: \"thinking_process\" }>[\"data\"],\n\t) => Promise<void> | void;\n\t/**\n\t * Called once the stream has been fully consumed with no error and no\n\t * client-initiated abort — i.e. on successful completion only.\n\t */\n\tonComplete?: () => Promise<void> | void;\n\tsetup: (signal: AbortSignal) => Promise<AsyncIterable<StreamEvent>>;\n};\n\nexport async function streamEventsToSSE(\n\treq: Request,\n\tres: Response,\n\toptions: SSEStreamOptions,\n): Promise<void> {\n\tres.writeHead(200, {\n\t\t\"Content-Type\": \"text/event-stream\",\n\t\t\"Cache-Control\": \"no-cache\",\n\t\tConnection: \"keep-alive\",\n\t\t\"X-Accel-Buffering\": \"no\", // nginx 버퍼링 비활성화\n\t});\n\tres.flushHeaders();\n\tres.write(\":ok\\n\\n\");\n\n\tconst keepaliveInterval = setInterval(() => {\n\t\tres.write(\":keepalive\\n\\n\");\n\t}, 10000); // 10초마다 keepalive 전송\n\n\tconst abortController = new AbortController();\n\tlet currentThreadId: string | undefined;\n\treq.on(\"close\", () => {\n\t\tabortController.abort();\n\t\tloggers.intentStream.info(`${options.logLabel} client connection closed`, {\n\t\t\tthreadId: currentThreadId,\n\t\t\tuserId: options.userId,\n\t\t\t...options.logContext,\n\t\t});\n\t});\n\n\ttry {\n\t\tconst stream = await options.setup(abortController.signal);\n\t\tfor await (const event of stream) {\n\t\t\tif (abortController.signal.aborted) break;\n\n\t\t\tif (event.event === \"thread_id\") {\n\t\t\t\tcurrentThreadId = event.data.threadId;\n\t\t\t\toptions.onThreadId?.(event.data.threadId);\n\t\t\t} else if (\n\t\t\t\tevent.event === \"thinking_process\" &&\n\t\t\t\tcurrentThreadId &&\n\t\t\t\toptions.onThinkingProcess\n\t\t\t) {\n\t\t\t\tawait options.onThinkingProcess(currentThreadId, event.data);\n\t\t\t}\n\n\t\t\tres.write(\n\t\t\t\t`event: ${event.event}\\ndata: ${JSON.stringify(event.data)}\\n\\n`,\n\t\t\t);\n\t\t}\n\t\tif (!abortController.signal.aborted) {\n\t\t\tawait options.onComplete?.();\n\t\t}\n\t} catch (error: unknown) {\n\t\tconst errMsg =\n\t\t\t(error as Error)?.message || `Failed to handle ${options.logLabel}`;\n\t\t// Carry the HTTP status (when the error is an AinHttpError) so the client\n\t\t// can distinguish e.g. 403 (no permission) from a generic failure — the\n\t\t// SSE headers are already sent, so it can't come back as a real status.\n\t\tconst status = (error as { status?: number })?.status;\n\t\t// The error only reaches the client as an SSE event; without this log a\n\t\t// failed stream leaves no server-side trace at all.\n\t\tloggers.intentStream.error(`${options.logLabel} failed`, {\n\t\t\tuserId: options.userId,\n\t\t\tthreadId: currentThreadId,\n\t\t\tstatus,\n\t\t\terror: errMsg,\n\t\t\t...options.logContext,\n\t\t});\n\t\tres.write(\n\t\t\t`event: error\\ndata: ${JSON.stringify({ message: errMsg, status })}\\n\\n`,\n\t\t);\n\t} finally {\n\t\tclearInterval(keepaliveInterval);\n\t\tres.end();\n\t}\n}\n"]}