File
Blob: src/worker/durable/run-do/logs.ts
| 1 | import { asc, desc, eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { LogStream, OpaqueId, RunId, UnixTimestampMs } from "@/contracts"; |
| 4 | import { type AppendRunLogsInput, expectTrusted, PositiveInteger, type RunLogRecord } from "@/worker/contracts"; |
| 5 | import * as runSchema from "@/worker/db/durable/schema/run-do"; |
| 6 | |
| 7 | import { chunkRows, toRunLogRecord, type RunDb, type RunTx } from "./repo"; |
| 8 | |
| 9 | const MAX_HOT_LOG_BYTES = 2 * 1024 * 1024; |
| 10 | const MAX_SQL_BOUND_PARAMETERS = 100; |
| 11 | const RUN_LOG_INSERT_COLUMN_COUNT = 6; |
| 12 | const MAX_RUN_LOG_ROWS_PER_INSERT = Math.floor(MAX_SQL_BOUND_PARAMETERS / RUN_LOG_INSERT_COLUMN_COUNT); |
| 13 | const textEncoder = new TextEncoder(); |
| 14 | |
| 15 | const 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 | |
| 26 | export 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 | |
| 36 | export 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 | |
| 61 | export 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 | }; |