Skip to content
File

Blob: src/worker/durable/run-do/logs.ts

typescript103 lines
1import { asc, desc, eq } from "drizzle-orm";
2 
3import { LogStream, OpaqueId, RunId, UnixTimestampMs } from "@/contracts";
4import { type AppendRunLogsInput, expectTrusted, PositiveInteger, type RunLogRecord } from "@/worker/contracts";
5import * as runSchema from "@/worker/db/durable/schema/run-do";
6 
7import { chunkRows, toRunLogRecord, type RunDb, type RunTx } from "./repo";
8 
9const MAX_HOT_LOG_BYTES = 2 * 1024 * 1024;
10const MAX_SQL_BOUND_PARAMETERS = 100;
11const RUN_LOG_INSERT_COLUMN_COUNT = 6;
12const MAX_RUN_LOG_ROWS_PER_INSERT = Math.floor(MAX_SQL_BOUND_PARAMETERS / RUN_LOG_INSERT_COLUMN_COUNT);
13const textEncoder = new TextEncoder();
14 
15const toAppendedRunLogRecord = (
16 row: Pick<typeof runSchema.runLogs.$inferInsert, "id" | "runId" | "seq" | "stream" | "chunk" | "createdAt">,
17): RunLogRecord => ({
18 id: expectTrusted(OpaqueId, row.id, "OpaqueId"),
19 runId: expectTrusted(RunId, row.runId, "RunId"),
20 seq: expectTrusted(PositiveInteger, row.seq, "PositiveInteger"),
21 stream: expectTrusted(LogStream, row.stream, "LogStream"),
22 chunk: row.chunk,
23 createdAt: expectTrusted(UnixTimestampMs, row.createdAt, "UnixTimestampMs"),
24});
25 
26export const listRunLogs = async (db: RunDb, runId: RunId): Promise<RunLogRecord[]> => {
27 const rows = await db
28 .select()
29 .from(runSchema.runLogs)
30 .where(eq(runSchema.runLogs.runId, runId))
31 .orderBy(asc(runSchema.runLogs.seq));
32 
33 return rows.map(toRunLogRecord);
34};
35 
36export const pruneLogs = async (db: RunDb, runId: RunId): Promise<void> => {
37 const rows = await db
38 .select({
39 id: runSchema.runLogs.id,
40 chunk: runSchema.runLogs.chunk,
41 })
42 .from(runSchema.runLogs)
43 .where(eq(runSchema.runLogs.runId, runId))
44 .orderBy(desc(runSchema.runLogs.seq));
45 
46 let retainedBytes = 0;
47 const deleteIds: string[] = [];
48 
49 for (const row of rows) {
50 retainedBytes += textEncoder.encode(row.chunk).length;
51 if (retainedBytes > MAX_HOT_LOG_BYTES) {
52 deleteIds.push(row.id);
53 }
54 }
55 
56 for (const id of deleteIds) {
57 await db.delete(runSchema.runLogs).where(eq(runSchema.runLogs.id, id));
58 }
59};
60 
61export const appendLogs = async (db: RunDb, input: AppendRunLogsInput): Promise<RunLogRecord[]> => {
62 const payload = input;
63 if (payload.events.length === 0) {
64 return [];
65 }
66 
67 let appendedLogs: RunLogRecord[] = [];
68 
69 db.transaction((tx: RunTx) => {
70 const latestRow = tx
71 .select({ seq: runSchema.runLogs.seq })
72 .from(runSchema.runLogs)
73 .where(eq(runSchema.runLogs.runId, payload.runId))
74 .orderBy(desc(runSchema.runLogs.seq))
75 .limit(1)
76 .get();
77 
78 let nextSeq = (latestRow?.seq ?? 0) + 1;
79 const logRows = payload.events.map((event) => {
80 const currentSeq = nextSeq;
81 nextSeq += 1;
82 
83 return {
84 id: `${payload.runId}:log:${currentSeq}`,
85 runId: payload.runId,
86 seq: currentSeq,
87 stream: event.stream,
88 chunk: event.chunk,
89 createdAt: event.createdAt,
90 };
91 });
92 
93 appendedLogs = logRows.map(toAppendedRunLogRecord);
94 
95 for (const batch of chunkRows(logRows, MAX_RUN_LOG_ROWS_PER_INSERT)) {
96 tx.insert(runSchema.runLogs).values(batch).run();
97 }
98 });
99 
100 await pruneLogs(db, payload.runId);
101 return appendedLogs;
102};