Skip to content
File

Blob: src/worker/dispatch/shared/run-steps/logging.ts

typescript69 lines
1import type { LogStream } from "@/contracts";
2import type { LogAppendEvent } from "@/worker/contracts";
3 
4import { now, type RunExecutionContext } from "@/worker/dispatch/shared/run-execution-context";
5 
6type RunLoggingContext = Pick<RunExecutionContext, "runStore" | "scope">;
7 
8export interface LogBatcher {
9 push(stream: Exclude<LogStream, "system">, chunk: string): void;
10 flush(): Promise<void>;
11}
12 
13export 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};