import { createHash } from "node:crypto"; import { promises as fs } from "node:fs"; import { join } from "node:path"; import type { PFrameInternal } from "@milaboratories/pl-model-middle-layer"; import type { PObjectId, DataQuery } from "@milaboratories/pl-model-common"; export function hashColumnId(columnId: PObjectId): PObjectId { return createHash("sha256").update(columnId).digest("hex") as PObjectId; } export function hashAddColumnEntry( entry: PFrameInternal.AddColumnEntryV2, ): PFrameInternal.AddColumnEntryV2 { return { ...entry, id: hashColumnId(entry.id), }; } export function hashUniqueValuesRequestColumnId( request: PFrameInternal.UniqueValuesRequestV2, ): PFrameInternal.UniqueValuesRequestV2 { // V2 filters are index-based (no column ids), so only `columnId` is rewritten. return { ...request, columnId: hashColumnId(request.columnId), }; } export function hashDataQuery(query: DataQuery): DataQuery { switch (query.type) { case "column": return { ...query, column: hashColumnId(query.column), }; case "inlineColumn": return { ...query, spec: { ...query.spec, id: hashColumnId(query.spec.id), }, }; case "sparseToDenseColumn": return { ...query, column: hashColumnId(query.column), ...(query.specOverride ? { specOverride: { ...query.specOverride, id: hashColumnId(query.specOverride.id), }, } : {}), }; case "innerJoin": case "fullJoin": return { ...query, entries: query.entries.map((e) => ({ ...e, entry: hashDataQuery(e.entry), })), }; case "outerJoin": return { ...query, primary: { ...query.primary, entry: hashDataQuery(query.primary.entry), }, secondary: query.secondary.map((e) => ({ ...e, entry: hashDataQuery(e.entry), })), }; case "linkerJoin": return { ...query, linker: { ...query.linker, column: hashColumnId(query.linker.column), }, secondary: query.secondary.map((e) => ({ ...e, entry: hashDataQuery(e.entry), })), }; case "sliceAxes": return { ...query, input: hashDataQuery(query.input), }; case "sort": return { ...query, input: hashDataQuery(query.input), }; case "filter": return { ...query, input: hashDataQuery(query.input), }; case "transformColumns": return { ...query, input: hashDataQuery(query.input), columns: query.columns.map((entry) => entry.specOverride ? { ...entry, specOverride: { ...entry.specOverride, id: hashColumnId(entry.specOverride.id), }, } : entry, ) as typeof query.columns, }; } } export async function dump( relativePath: string[], data: { [key: string]: string | number | boolean | object } | Uint8Array, logger?: PFrameInternal.Logger, ): Promise { if (!process.env.MI_DUMP_PFRAMES_RS) return; try { const relativeUri = relativePath.map((part) => encodeURIComponent(part)); const fileDir = join(process.env.MI_DUMP_PFRAMES_RS, ...relativeUri.slice(0, -1)); await fs.mkdir(fileDir, { recursive: true }); const filePath = join(process.env.MI_DUMP_PFRAMES_RS, ...relativeUri); const fileData = ArrayBuffer.isView(data) ? data : JSON.stringify(data, null, 2); await fs.writeFile(filePath, fileData, { flag: "wx" }); } catch (error: unknown) { logger?.("info", `error while dumping PFrames data: ${error}`); } }