{"version":3,"file":"jsonl-writer.d.ts","sourceRoot":"","sources":["../../../src/shared/jsonl-writer.ts"],"names":[],"mappings":"AAEA,MAAM,WAAW,eAAe;IAC/B,KAAK,IAAI,IAAI,CAAC;IACd,MAAM,IAAI,IAAI,CAAC;CACf;AAED,MAAM,WAAW,gBAAgB;IAChC,KAAK,CAAC,KAAK,EAAE,MAAM,GAAG,OAAO,CAAC;IAC9B,IAAI,CAAC,KAAK,EAAE,OAAO,EAAE,QAAQ,EAAE,MAAM,IAAI,GAAG,gBAAgB,CAAC;IAC7D,GAAG,CAAC,QAAQ,CAAC,EAAE,MAAM,IAAI,GAAG,IAAI,CAAC;CACjC;AAID,UAAU,eAAe;IACxB,iBAAiB,CAAC,EAAE,CAAC,QAAQ,EAAE,MAAM,KAAK,gBAAgB,CAAC;IAC3D,QAAQ,CAAC,EAAE,MAAM,CAAC;CAClB;AAED,UAAU,WAAW;IACpB,SAAS,CAAC,IAAI,EAAE,MAAM,GAAG,IAAI,CAAC;IAC9B,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC,CAAC;CACvB;AAED,wBAAgB,iBAAiB,CAChC,QAAQ,EAAE,MAAM,GAAG,SAAS,EAC5B,MAAM,EAAE,eAAe,EACvB,IAAI,GAAE,eAAoB,GACxB,WAAW,CAoDb","sourcesContent":["import * as fs from \"node:fs\";\n\nexport interface DrainableSource {\n\tpause(): void;\n\tresume(): void;\n}\n\nexport interface JsonlWriteStream {\n\twrite(chunk: string): boolean;\n\tonce(event: \"drain\", listener: () => void): JsonlWriteStream;\n\tend(callback?: () => void): void;\n}\n\nconst DEFAULT_MAX_JSONL_BYTES = 50 * 1024 * 1024;\n\ninterface JsonlWriterDeps {\n\tcreateWriteStream?: (filePath: string) => JsonlWriteStream;\n\tmaxBytes?: number;\n}\n\ninterface JsonlWriter {\n\twriteLine(line: string): void;\n\tclose(): Promise<void>;\n}\n\nexport function createJsonlWriter(\n\tfilePath: string | undefined,\n\tsource: DrainableSource,\n\tdeps: JsonlWriterDeps = {},\n): JsonlWriter {\n\tif (!filePath) {\n\t\treturn {\n\t\t\twriteLine() {},\n\t\t\tasync close() {},\n\t\t};\n\t}\n\n\tconst createWriteStream =\n\t\tdeps.createWriteStream ?? ((targetPath: string) => fs.createWriteStream(targetPath, { flags: \"a\" }));\n\tlet stream: JsonlWriteStream | undefined;\n\ttry {\n\t\tstream = createWriteStream(filePath);\n\t} catch {\n\t\treturn {\n\t\t\twriteLine() {},\n\t\t\tasync close() {},\n\t\t};\n\t}\n\n\tlet backpressured = false;\n\tlet closed = false;\n\tlet bytesWritten = 0;\n\tconst maxBytes = deps.maxBytes ?? DEFAULT_MAX_JSONL_BYTES;\n\n\treturn {\n\t\twriteLine(line: string) {\n\t\t\tif (!stream || closed || !line.trim()) return;\n\t\t\tconst chunk = `${line}\\n`;\n\t\t\tconst chunkBytes = Buffer.byteLength(chunk, \"utf-8\");\n\t\t\tif (bytesWritten + chunkBytes > maxBytes) return;\n\t\t\ttry {\n\t\t\t\tconst ok = stream.write(chunk);\n\t\t\t\tbytesWritten += chunkBytes;\n\t\t\t\tif (!ok && !backpressured) {\n\t\t\t\t\tbackpressured = true;\n\t\t\t\t\tsource.pause();\n\t\t\t\t\tstream.once(\"drain\", () => {\n\t\t\t\t\t\tbackpressured = false;\n\t\t\t\t\t\tif (!closed) source.resume();\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t} catch {}\n\t\t},\n\t\tasync close() {\n\t\t\tif (!stream || closed) return;\n\t\t\tclosed = true;\n\t\t\tconst current = stream;\n\t\t\tstream = undefined;\n\t\t\tawait new Promise<void>((resolve) => current.end(() => resolve()));\n\t\t},\n\t};\n}\n"]}