import { ulid } from "ulid"; import { AbortError, PFrameError, type PTableShape, type PTableVector, type TableRange, type UniqueValuesResponse, type DataQuery, ensureError, } from "@milaboratories/pl-model-common"; import type { PFrameInternal } from "@milaboratories/pl-model-middle-layer"; import { PerfTimer } from "@milaboratories/helpers"; import type { CancellationTokenSymbol, NodeFrameSymbol, NodeTableSymbol } from "./addon-def"; import { AddonSymbol } from "./addon"; import { dump, hashColumnId, hashAddColumnEntry, hashUniqueValuesRequestColumnId, hashDataQuery, } from "./dump"; /** * Run `callback` with a native cancellation token bridged to a JS `AbortSignal`. * * The abort listener is registered with `{ once: true }` and removed in * `finally`, so it is always cleaned up — including when `throwIfAborted` * short-circuits an already-aborted signal before the addon is ever called. * Detecting whether the operation was actually cancelled is left to the caller * via `signal.aborted`. */ async function callAbortable( callback: (token: CancellationTokenSymbol | undefined) => Promise, signal: AbortSignal | undefined, ): Promise { if (!signal) { return await callback(undefined); } const token = AddonSymbol.cancellationTokenCreate(); const abort = (): void => AddonSymbol.cancellationTokenCancel(token); try { signal.addEventListener("abort", abort, { once: true }); signal.throwIfAborted(); return await callback(token); } finally { signal.removeEventListener("abort", abort); } } export async function pprofDump(): Promise { try { return await AddonSymbol.pprofDump(); } catch (err: unknown) { const error = new PFrameError(`PFrame pprofDump failed`); error.cause = ensureError(err); throw error; } } export class PFrame implements PFrameInternal.PFrameV18 { readonly #frame: NodeFrameSymbol; readonly #id: PFrameInternal.PFrameId; get id(): PFrameInternal.PFrameId { return this.#id; } readonly #logger: PFrameInternal.Logger; constructor(options: PFrameInternal.PFrameOptionsV2) { this.#id = options.frameId; this.#logger = options.logger ?? (() => {}); dump( [`${this.id}`, `${this.id}.json`], { timeStamp: Date.now(), requestType: "create", }, this.#logger, ); try { this.#frame = AddonSymbol.pFrameCreate(options.spillPath, options.logger); this.#logger("info", `PFrame ${this.id} created`); } catch (err: unknown) { const error = new PFrameError(`PFrame creation failed`); error.cause = new Error( `PFrame ${this.id} creation failed, ` + `logger: ${this.#logger?.toString()}, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } } async addColumns( columns: PFrameInternal.AddColumnEntryV2[], ops?: { signal?: AbortSignal; }, ): Promise { const requestId = ulid(); dump( [`${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "addColumns", requestData: { columns: columns.map(hashAddColumnEntry), }, }, this.#logger, ); for (const column of columns) { const hashedId = hashColumnId(column.id); dump([`${this.id}`, `data`, `${hashedId}.datainfo`], { ...column.data }, this.#logger); } try { return await callAbortable( (token) => AddonSymbol.pFrameAddColumns(this.#frame, columns, token), ops?.signal, ); } catch (err: unknown) { if (ops?.signal?.aborted) { throw new AbortError(`PFrame ${this.id} addColumns request ${requestId} cancelled`); } const error = new PFrameError(`PFrame addColumns request failed`); error.cause = new Error( `PFrame ${this.id} addColumns request ${requestId} failed, ` + `columns: ${JSON.stringify(columns)}, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } } setDataSource(dataSource: PFrameInternal.PFrameDataSourceV2): void { const requestId = ulid(); dump( [`${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "setDataSource", }, this.#logger, ); const wrappedDataSource = { preloadBlob: async (blobIds: PFrameInternal.PFrameBlobId[]): Promise => { const requestId = ulid(); dump( [`${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "preloadBlob", requestData: { blobIds, }, }, this.#logger, ); this.#logger( "info", `PFrame ${this.id} preloadBlob started, blobIds: ${JSON.stringify(blobIds)}`, ); const timer = PerfTimer.start(); try { return await dataSource.preloadBlob(blobIds); } finally { this.#logger( "info", `PFrame ${this.id} preloadBlob finished, took ${timer.elapsed()} (${blobIds.length} blobs)`, ); } }, resolveBlobContent: async (blobId: PFrameInternal.PFrameBlobId): Promise => { const requestId = ulid(); dump( [`${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "resolveBlobContent", requestData: { blobId, }, }, this.#logger, ); const blob = await dataSource.resolveBlobContent(blobId); this.#logger("info", `PFrame ${this.id} resolved blob ${blobId}`); dump([`${this.id}`, `data`, `${blobId}`], blob, this.#logger); return blob; }, parquetServer: dataSource.parquetServer, }; try { return AddonSymbol.pFrameSetDataSource(this.#frame, wrappedDataSource); } catch (err: unknown) { const error = new PFrameError(`PFrame setDataSource request failed`); error.cause = new Error( `PFrame ${this.id} setDataSource request ${requestId} failed, ` + `dataSource: ${dataSource.toString()}, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } } dispose(): void { const requestId = ulid(); dump( [`${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "dispose", }, this.#logger, ); try { AddonSymbol.pFrameDispose(this.#frame); this.#logger("info", `PFrame ${this.id} disposed`); } catch (err: unknown) { const error = new PFrameError(`PFrame dispose request failed`); error.cause = new Error( `PFrame ${this.id} dispose request ${requestId} failed, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } } [Symbol.dispose](): void { this.dispose(); } createTable(tableId: PFrameInternal.PTableId, dataQuery: DataQuery): PTable { const timer = PerfTimer.start(); const dumpData = { timeStamp: Date.now(), requestType: "createTable", requestData: hashDataQuery(dataQuery), }; dump([`${this.id}`, `${tableId}.json`], dumpData, this.#logger); dump([`${this.id}`, `${tableId}`, `${tableId}.json`], dumpData, this.#logger); try { const boxed = AddonSymbol.pFrameCreateTable(this.#frame, tableId, dataQuery); return new PTable(this, tableId, boxed, this.#logger); } catch (err: unknown) { const error = new PFrameError(`PFrame createTable request failed`); error.cause = new Error( `PFrame ${this.id} createTable request ${tableId} failed, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } finally { this.#logger( "info", `PFrame ${this.id} createTable request ${tableId} took ${timer.elapsed()}`, ); } } async getUniqueValues( request: PFrameInternal.UniqueValuesRequestV2, ops?: { signal?: AbortSignal; }, ): Promise { const requestId = ulid() as PFrameInternal.PTableId; dump( [`${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "getUniqueValues", requestData: hashUniqueValuesRequestColumnId(request), }, this.#logger, ); this.#logger("info", `PFrame ${this.id} getUniqueValues request ${requestId} started`); const timer = PerfTimer.start(); try { return await callAbortable( (token) => AddonSymbol.pFrameGetUniqueValues(this.#frame, requestId, request, token), ops?.signal, ); } catch (err: unknown) { if (ops?.signal?.aborted) { throw new AbortError(`PFrame ${this.id} getUniqueValues request ${requestId} cancelled`); } const error = new PFrameError(`PFrame getUniqueValues request failed`); error.cause = new Error( `PFrame ${this.id} getUniqueValues request ${requestId} failed, ` + `request: ${JSON.stringify(request)}, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } finally { this.#logger( "info", `PFrame ${this.id} getUniqueValues request ${requestId} finished, took ${timer.elapsed()}`, ); } } } class PTable implements PFrameInternal.PTableV12 { readonly #frame: PFrame; readonly #table: NodeTableSymbol; readonly #id: string; get id(): string { return this.#id; } readonly #logger: PFrameInternal.Logger; constructor(frame: PFrame, id: string, table: NodeTableSymbol, logger: PFrameInternal.Logger) { this.#frame = frame; this.#table = table; this.#id = id; this.#logger = logger; this.#logger("info", `PTable ${this.id} created`); } async getFootprint(ops?: { withPredecessors?: boolean; signal?: AbortSignal }): Promise { const requestId = ulid(); dump( [`${this.#frame.id}`, `${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "getFootprint", }, this.#logger, ); try { return await callAbortable( (token) => AddonSymbol.pTableGetFootprint(this.#table, ops?.withPredecessors ?? false, token), ops?.signal, ); } catch (err: unknown) { if (ops?.signal?.aborted) { throw new AbortError(`PTable ${this.id} getFootprint request ${requestId} cancelled`); } const error = new PFrameError(`PTable getFootprint request failed`); error.cause = new Error( `PTable ${this.id} getFootprint request ${requestId} failed, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } } async getShape(ops?: { signal?: AbortSignal }): Promise { const requestId = ulid(); dump( [`${this.#frame.id}`, `${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "getShape", }, this.#logger, ); this.#logger("info", `PTable ${this.id} getShape request ${requestId} started`); const timer = PerfTimer.start(); try { return await callAbortable( (token) => AddonSymbol.pTableGetShape(this.#table, token), ops?.signal, ); } catch (err: unknown) { if (ops?.signal?.aborted) { throw new AbortError(`PTable ${this.id} getShape request ${requestId} cancelled`); } const error = new PFrameError(`PTable getShape request failed`); error.cause = new Error( `PTable ${this.id} getShape request ${requestId} failed, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } finally { this.#logger( "info", `PTable ${this.id} getShape request ${requestId} finished, took ${timer.elapsed()}`, ); } } async getData( columnIndices: number[], ops?: { range?: TableRange | undefined; signal?: AbortSignal | undefined; }, ): Promise { const requestId = ulid() as PFrameInternal.PTableId; dump( [`${this.#frame.id}`, `${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "getData", requestData: { columnIndices, range: ops?.range ?? null, }, }, this.#logger, ); this.#logger("info", `PTable ${this.id} getData request ${requestId} started`); let rowCount = 0; const timer = PerfTimer.start(); try { const result = await callAbortable( (token) => AddonSymbol.pTableGetData(this.#table, columnIndices, ops?.range, token), ops?.signal, ); rowCount = result[0].data.length; return result; } catch (err: unknown) { if (ops?.signal?.aborted) { throw new AbortError(`PTable ${this.id} getData request ${requestId} cancelled`); } const error = new PFrameError(`PTable getData request failed`); error.cause = new Error( `PTable ${this.id} getData request ${requestId} failed, ` + `columnIndices: ${JSON.stringify(columnIndices)}, ` + `range: ${ops?.range ? JSON.stringify(ops.range) : undefined}, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } finally { this.#logger( "info", `PTable ${this.id} getData request ${requestId} finished, took ${timer.elapsed()} (${rowCount} rows)`, ); } } async export( path: string, ops?: { headers: [number, string][]; signal?: AbortSignal; }, ): Promise { const headers = ops?.headers ?? []; const requestId = ulid() as PFrameInternal.PTableId; dump( [`${this.#frame.id}`, `${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "export", requestData: { path, headers, }, }, this.#logger, ); this.#logger("info", `PTable ${this.id} export request ${requestId} started`); const timer = PerfTimer.start(); try { return await callAbortable( (token) => AddonSymbol.pTableExport(this.#table, path, headers, token), ops?.signal, ); } catch (err: unknown) { if (ops?.signal?.aborted) { throw new AbortError(`PTable ${this.id} export request ${requestId} cancelled`); } const error = new PFrameError(`PTable export request failed`); error.cause = new Error( `PTable ${this.id} export request ${requestId} failed, ` + `path: ${path}, ` + `headers: ${JSON.stringify(headers)}, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } finally { this.#logger( "info", `PTable ${this.id} export request ${requestId} finished, took ${timer.elapsed()}`, ); } } dispose() { const requestId = ulid(); dump( [`${this.#frame.id}`, `${this.id}`, `${requestId}.json`], { timeStamp: Date.now(), requestType: "dispose", }, this.#logger, ); try { AddonSymbol.pTableDispose(this.#table); this.#logger("info", `PTable ${this.id} disposed`); } catch (err: unknown) { const error = new PFrameError(`PTable dispose request failed`); error.cause = new Error( `PTable ${this.id} dispose request ${requestId} failed, ` + `error:\n` + `${ensureError(err)}`, ); throw error; } } [Symbol.dispose]() { this.dispose(); } }