File
Blob: src/worker/dispatch/shared/run-steps/logging.ts
| 1 | import type { LogStream } from "@/contracts"; |
| 2 | import type { LogAppendEvent } from "@/worker/contracts"; |
| 3 | |
| 4 | import { now, type RunExecutionContext } from "@/worker/dispatch/shared/run-execution-context"; |
| 5 | |
| 6 | type RunLoggingContext = Pick<RunExecutionContext, "runStore" | "scope">; |
| 7 | |
| 8 | export interface LogBatcher { |
| 9 | push(stream: Exclude<LogStream, "system">, chunk: string): void; |
| 10 | flush(): Promise<void>; |
| 11 | } |
| 12 | |
| 13 | export const createLogBatcher = (context: RunLoggingContext): LogBatcher => { |
| 14 | const buffer: LogAppendEvent[] = []; |
| 15 | let bufferedBytes = 0; |
| 16 | let inFlightFlush: Promise<void> | null = null; |
| 17 | |
| 18 | const startFlush = (): Promise<void> => { |
| 19 | if (inFlightFlush) { |
| 20 | return inFlightFlush; |
| 21 | } |
| 22 | |
| 23 | inFlightFlush = (async () => { |
| 24 | try { |
| 25 | while (buffer.length > 0) { |
| 26 | const events = buffer.splice(0, buffer.length); |
| 27 | bufferedBytes = 0; |
| 28 | await context.runStore.appendLogs(events); |
| 29 | } |
| 30 | } finally { |
| 31 | inFlightFlush = null; |
| 32 | if (buffer.length > 0) { |
| 33 | void startFlush(); |
| 34 | } |
| 35 | } |
| 36 | })(); |
| 37 | |
| 38 | return inFlightFlush; |
| 39 | }; |
| 40 | |
| 41 | return { |
| 42 | push(stream: Exclude<LogStream, "system">, chunk: string) { |
| 43 | if (chunk.length === 0) { |
| 44 | return; |
| 45 | } |
| 46 | |
| 47 | buffer.push({ |
| 48 | stream, |
| 49 | chunk, |
| 50 | createdAt: now(), |
| 51 | }); |
| 52 | bufferedBytes += new TextEncoder().encode(chunk).length; |
| 53 | |
| 54 | if (bufferedBytes >= 4096) { |
| 55 | void startFlush(); |
| 56 | } |
| 57 | }, |
| 58 | async flush() { |
| 59 | if (buffer.length > 0) { |
| 60 | await startFlush(); |
| 61 | } |
| 62 | |
| 63 | while (inFlightFlush) { |
| 64 | await inFlightFlush; |
| 65 | } |
| 66 | }, |
| 67 | }; |
| 68 | }; |