import * as duration from 'duration-fns'; import stableStringify from 'fast-json-stable-stringify'; import { SqupContext } from '../types/SqupContext.js'; export type DateFormat = 'ISO-8601' | 'Ticks' | 'Unix'; export type CompressTimeSeriesOptions = { dateFormat?: DateFormat; interval?: string; valueColumn?: string | string[]; trimBefore?: Date; trimAfter?: Date; }; type AggregateValue = { count: number; sum: number; }; export function CompressTimeSeries( data: Record[], timestampColumn: string, squpContext: SqupContext, opts: CompressTimeSeriesOptions = {} ): Record[] { squpContext.log.debug(`CompressTimeSeries called: interval:${opts.interval}, valueColumn:${JSON.stringify(opts.valueColumn)}`); let intervalSeconds = 0; if (typeof opts.interval === 'string' && opts.interval) { intervalSeconds = duration.toSeconds(opts.interval); if (isNaN(intervalSeconds) || intervalSeconds < 1) { throw new Error(`Invalid interval: "${opts.interval}" specified`); } } const intervalMilliSeconds = 1000 * intervalSeconds; const trimBeforeMilliSeconds: number | null = opts.trimBefore ? opts.trimBefore.getTime() : null; const trimAfterMilliSeconds: number | null = opts.trimAfter ? opts.trimAfter.getTime() : null; const valueColumns: string[] = opts.valueColumn ? Array.isArray(opts.valueColumn) ? opts.valueColumn : [opts.valueColumn] : []; const aggResult = new Map>>(); for (const row of data) { const ts = getTicks(row[timestampColumn], opts.dateFormat ?? 'ISO-8601'); if (isNaN(ts)) { continue; } if (trimBeforeMilliSeconds !== null && ts < trimBeforeMilliSeconds) { continue; } if (trimAfterMilliSeconds !== null && ts > trimAfterMilliSeconds) { continue; } const bucketTs = intervalMilliSeconds ? Math.floor(ts / intervalMilliSeconds) * intervalMilliSeconds : ts; const groupObject = Object.entries(row).reduce( (acc, [key, val]) => { if (key !== timestampColumn && !valueColumns.includes(key)) { acc[key] = val; } return acc; }, {} as Record ); //groupObject[timestampColumn] = bucketTs; const groupKey = stableStringify(groupObject); let bucket = aggResult.get(bucketTs); if (!bucket) { bucket = new Map>(); aggResult.set(bucketTs, bucket); } if (intervalMilliSeconds) { let aggRow = bucket.get(groupKey); if (!aggRow) { aggRow = { ...row }; aggRow[timestampColumn] = bucketTs; for (const valueColumn of valueColumns) { aggRow[valueColumn] = { count: 0, sum: 0 } as AggregateValue; } bucket.set(groupKey, aggRow); } for (const valueColumn of valueColumns) { const val = Number(row[valueColumn]); if (!isNaN(val)) { const agg = aggRow[valueColumn] as AggregateValue; agg.count++; agg.sum += val; } } } else { bucket.set(groupKey, row); } } const result = [] as Record[]; for (const bucketTs of Array.from(aggResult.keys()).sort()) { const bucket = aggResult.get(bucketTs); if (bucket) { for (const aggRow of bucket.values()) { aggRow[timestampColumn] = getTimeInFormat(bucketTs, opts.dateFormat ?? 'ISO-8601'); if (intervalMilliSeconds) { for (const valueColumn of valueColumns) { const agg = aggRow[valueColumn] as AggregateValue; aggRow[valueColumn] = agg.sum / agg.count; } } result.push(aggRow); } } } return result; } function getTicks(date: unknown, dateFormat: DateFormat): number { switch (dateFormat) { case 'ISO-8601': { if (typeof date === 'string') { try { return Date.parse(date); } catch { return NaN; } } else { return NaN; } } case 'Ticks': { return Number(date); } case 'Unix': { return Number(date) * 1000; } default: { throw new Error(`Bad dateFormat: "${dateFormat}"`); } } } function getTimeInFormat(ticks: number, dateFormat: DateFormat) { switch (dateFormat) { case 'ISO-8601': { return new Date(ticks).toISOString(); } case 'Ticks': { return ticks; } case 'Unix': { return Math.floor(ticks / 1000); } default: { throw new Error(`Bad dateFormat: "${dateFormat}"`); } } }