pf_gkey is the Perspective table index -- an opaque handle read only by table.update() and table.remove(). It was the raw concat_ws of every grain column, which on the 24-column grain of fc_osm_skinny_29 averaged 233 characters. Being unique per row by construction, it also defeated Arrow's dictionary encoding, so it alone accounted for 65.6 MB of a 109 MB payload -- more than the other 31 columns combined, all of which do dictionary cleanly. md5 of the same string keeps the determinism undo depends on (routes/log.js recomputes the key through this same grainOf) and keeps collisions out of reach at 128 bits, while fixing the width at 32 characters. Measured on the live 285,685-row aggregate: payload 109.0 MB -> 54.2 MB, all 285,685 keys still distinct. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
478 lines
19 KiB
JavaScript
478 lines
19 KiB
JavaScript
// Generates operation SQL for a source table, baking in column names from col_meta.
|
|
// Runtime values are left as {{token}} substitution points.
|
|
//
|
|
// Columns flagged col_meta.in_grain define a display grain. When one is set the
|
|
// initial load (get_agg) and every operation return rows pre-aggregated to that
|
|
// grain and keyed on pf_gkey, instead of raw forecast rows keyed on pf_id.
|
|
//
|
|
// Tokens baked in at generation time: column names, source schema.table
|
|
// Tokens substituted at request time: {{fc_table}}, {{where_clause}}, {{exclude_clause}},
|
|
// {{version_id}}, {{logid}}, {{pf_user}}, {{note}},
|
|
// {{params}}, {{slice}}, {{date_from}}, {{date_to}},
|
|
// {{value_incr}}, {{units_incr}}, {{set_clause}}, {{scale_factor}}
|
|
|
|
// wrap a column name in double quotes for safe use in SQL
|
|
function q(name) { return `"${name}"`; }
|
|
|
|
// The display grain: dimension/date columns flagged in_grain, plus pf_iter and
|
|
// pf_logid which are always part of it. Returns null when nothing is flagged —
|
|
// that is raw-row mode, where operations return whole rows and the client
|
|
// indexes on pf_id (the pre-grain behaviour).
|
|
//
|
|
// Keeping pf_logid in the grain is what makes the append model work: each
|
|
// operation's contribution stays a distinct row, so table.update() accumulates
|
|
// rather than replacing a bucket total, and undo can remove exactly that
|
|
// operation's rows.
|
|
function grainOf(colMeta) {
|
|
const cols = colMeta
|
|
.filter(c => c.in_grain && (c.role === 'dimension' || c.role === 'date'))
|
|
.sort((a, b) => (a.opos || 0) - (b.opos || 0))
|
|
.map(c => c.cname);
|
|
if (cols.length === 0) return null;
|
|
|
|
// pf_gkey must be unique per grain tuple. chr(31) (unit separator) joins the
|
|
// parts and chr(30) stands in for NULL, so ('a', NULL) cannot collide with
|
|
// (NULL, 'a') and a NULL stays distinct from an empty string — a collision
|
|
// would silently merge two groups into one indexed row.
|
|
//
|
|
// md5 of that, rather than the concatenation itself, because the key is an
|
|
// opaque handle -- nothing reads it but table.update() and table.remove().
|
|
// The raw form averaged 233 chars on a 24-column grain and, being unique per
|
|
// row, defeated Arrow's dictionary encoding: 65.6 MB of a 109 MB payload,
|
|
// more than every other column combined. 128 bits keeps collisions unreachable.
|
|
const key = (pfx = '') => `md5(concat_ws(chr(31), ${[
|
|
...cols.map(c => `COALESCE(${pfx}${q(c)}::text, chr(30))`),
|
|
`${pfx}pf_iter`,
|
|
`${pfx}pf_logid::text`
|
|
].join(', ')}))`;
|
|
|
|
const groupCols = (pfx = '') => [...cols.map(c => `${pfx}${q(c)}`), `${pfx}pf_iter`, `${pfx}pf_logid`];
|
|
|
|
return { cols, key, groupCols };
|
|
}
|
|
|
|
function generateSQL(source, colMeta) {
|
|
const dims = colMeta
|
|
.filter(c => c.role === 'dimension')
|
|
.sort((a, b) => (a.opos || 0) - (b.opos || 0))
|
|
.map(c => c.cname);
|
|
|
|
// Every column of each measure/date role, in col_meta order. Loads carry all of
|
|
// them; the adjustment operations are single-measure (scale distributes one
|
|
// {{value_incr}}) and use only the first of each, below.
|
|
const byRole = role => colMeta
|
|
.filter(c => c.role === role)
|
|
.sort((a, b) => (a.opos || 0) - (b.opos || 0))
|
|
.map(c => c.cname);
|
|
|
|
const valueCols = byRole('value');
|
|
const unitsCols = byRole('units');
|
|
const dateCols = byRole('date');
|
|
|
|
const valueCol = valueCols[0];
|
|
const unitsCol = unitsCols[0];
|
|
const dateCol = dateCols[0];
|
|
|
|
if (!valueCol) throw new Error('No value column defined in col_meta');
|
|
if (!dateCol) throw new Error('No date column defined in col_meta');
|
|
if (dims.length === 0) throw new Error('No dimension columns defined in col_meta');
|
|
|
|
const srcTable = `"${source.schema}"."${source.tname}"`;
|
|
const dataCols = [...dims, dateCol, valueCol, unitsCol].filter(Boolean);
|
|
const effectiveValue = dataCols.includes(valueCol) ? valueCol : null;
|
|
const effectiveUnits = dataCols.includes(unitsCol) ? unitsCol : null;
|
|
const insertCols = [...dataCols.map(q), 'pf_iter', 'pf_logid', 'pf_user', 'pf_created_at'].join(', ');
|
|
const selectData = dataCols.map(q).join(', ');
|
|
const dimsJoined = dims.map(q).join(', ');
|
|
|
|
// Baseline and reference copy the source row wholesale, so they carry every
|
|
// measure and every date — not just the primary one the operations act on.
|
|
// Dropping the others would leave those columns null for the life of the version.
|
|
const loadCols = [...dims, ...dateCols, ...valueCols, ...unitsCols];
|
|
const loadInsertCols = [...loadCols.map(q), 'pf_iter', 'pf_logid', 'pf_user', 'pf_created_at'].join(', ');
|
|
const dateColSet = new Set(dateCols);
|
|
|
|
// dim_period JOIN support: if the date column is the is_key of a dim_group,
|
|
// dimension siblings with dim_period_col set are derived from pf.dim_period
|
|
// instead of being copied raw from the source on baseline/reference load.
|
|
const dateKeyGroup = colMeta.find(c => c.role === 'date' && c.is_key && c.dim_group)?.dim_group;
|
|
const dimPeriodMap = new Map(
|
|
dateKeyGroup
|
|
? colMeta
|
|
.filter(c => c.role === 'dimension' && c.dim_group === dateKeyGroup && c.dim_period_col)
|
|
.map(c => [c.cname, c.dim_period_col])
|
|
: []
|
|
);
|
|
const hasDimPeriod = dimPeriodMap.size > 0;
|
|
|
|
// display grain — when set, initial load and operations both return rows
|
|
// pre-aggregated to it instead of raw forecast rows
|
|
const grain = grainOf(colMeta);
|
|
// A flag on anything other than a dimension/date column is ignored by grainOf,
|
|
// which is what we want — the role change is the source of truth, not a stale flag.
|
|
if (grain) {
|
|
const missing = grain.cols.filter(c => !dataCols.includes(c));
|
|
if (missing.length > 0) {
|
|
throw new Error(
|
|
`Grain columns are never populated in the forecast table: ${missing.join(', ')}`
|
|
);
|
|
}
|
|
if (!effectiveValue && !effectiveUnits) {
|
|
throw new Error('A grain requires at least one value or units column to aggregate');
|
|
}
|
|
}
|
|
|
|
return {
|
|
get_data: buildGetData(),
|
|
...(grain ? { get_agg: buildGetAgg() } : {}),
|
|
baseline: buildBaseline(),
|
|
reference: buildReference(),
|
|
scale: buildScale(),
|
|
recode: buildRecode(),
|
|
clone: buildClone(),
|
|
undo: buildUndo()
|
|
};
|
|
|
|
function buildGetData() {
|
|
return `SELECT * FROM {{fc_table}}`;
|
|
}
|
|
|
|
// Aggregate the whole forecast table to the display grain. This is the initial
|
|
// load for grain sources — the client loads the result into a native Perspective
|
|
// table indexed on pf_gkey and its view sums across these rows, exactly as an
|
|
// Excel pivot cache sums its data tab.
|
|
function buildGetAgg() {
|
|
// pf_logid is part of the grain, so joining pf.log adds no rows — each group
|
|
// already belongs to exactly one log entry. Without this the segment labels
|
|
// that /data surfaces would vanish the moment a source declares a grain.
|
|
return `
|
|
SELECT
|
|
${grainSelect('t.')}
|
|
,CASE WHEN l.operation IN ('baseline','reference')
|
|
THEN COALESCE(NULLIF(l.tag, ''), NULLIF(l.note, ''), '(unlabeled load)')
|
|
ELSE '(adjustment)' END AS pf_segment
|
|
,CASE WHEN l.operation IN ('baseline','reference')
|
|
THEN NULL
|
|
ELSE COALESCE(NULLIF(l.tag, ''), NULLIF(l.note, '')) END AS pf_note
|
|
,l.operation AS pf_op
|
|
FROM {{fc_table}} t
|
|
LEFT JOIN pf.log l
|
|
ON l.id = t.pf_logid
|
|
GROUP BY
|
|
${grain.groupCols('t.').join('\n ,')}
|
|
,l.operation
|
|
,l.tag
|
|
,l.note`.trim();
|
|
}
|
|
|
|
// grain columns + pf_gkey + summed measures, in the leading-comma style the
|
|
// rest of the generated SQL uses
|
|
function grainSelect(pfx = '') {
|
|
return [
|
|
...grain.groupCols(pfx),
|
|
`${grain.key(pfx)} AS pf_gkey`,
|
|
effectiveValue ? `SUM(${pfx}${q(effectiveValue)}) AS ${q(effectiveValue)}` : null,
|
|
effectiveUnits ? `SUM(${pfx}${q(effectiveUnits)}) AS ${q(effectiveUnits)}` : null
|
|
].filter(Boolean).join('\n ,');
|
|
}
|
|
|
|
// Tail of an operation statement: in grain mode the inserted rows come back
|
|
// aggregated to grain (the client appends them and lets the view re-sum);
|
|
// otherwise whole rows come back as before.
|
|
function opTail(cte) {
|
|
if (!grain) return `SELECT * FROM ${cte}`;
|
|
return `
|
|
SELECT
|
|
${grainSelect()}
|
|
FROM ${cte}
|
|
GROUP BY
|
|
${grain.groupCols().join('\n ,')}`.trim();
|
|
}
|
|
|
|
function buildLoadSelect(pfx) {
|
|
// pfx: table alias prefix ('s.' when joining dim_period, '' otherwise)
|
|
// The offset shifts every date column, so order date and ship date stay in step.
|
|
return loadCols.map(c => {
|
|
if (dateColSet.has(c)) return `(${pfx}${q(c)} + '{{date_offset}}'::interval)::date`;
|
|
if (dimPeriodMap.has(c)) return `dp.${q(dimPeriodMap.get(c))} AS ${q(c)}`;
|
|
return `${pfx}${q(c)}`;
|
|
}).join(',\n ');
|
|
}
|
|
|
|
function buildFromClause() {
|
|
if (!hasDimPeriod) return srcTable;
|
|
return `${srcTable} s\n JOIN pf.dim_period dp`
|
|
+ ` ON dp.drange @> (s.${q(dateCol)} + '{{date_offset}}'::interval)::date`;
|
|
}
|
|
|
|
function buildBaseline() {
|
|
return `
|
|
WITH
|
|
ilog AS (
|
|
INSERT INTO pf.log (version_id, pf_user, operation, slice, params, note)
|
|
VALUES ({{version_id}}, '{{pf_user}}', 'baseline', NULL, '{{params}}'::jsonb, '{{note}}')
|
|
RETURNING id
|
|
)
|
|
,ins AS (
|
|
INSERT INTO {{fc_table}} (${loadInsertCols})
|
|
SELECT
|
|
${buildLoadSelect(hasDimPeriod ? 's.' : '')},
|
|
'baseline', (SELECT id FROM ilog), '{{pf_user}}', now()
|
|
FROM ${buildFromClause()}
|
|
WHERE {{filter_clause}}
|
|
RETURNING *
|
|
)
|
|
SELECT count(*) AS rows_affected FROM ins`.trim();
|
|
}
|
|
|
|
function buildReference() {
|
|
return `
|
|
WITH
|
|
ilog AS (
|
|
INSERT INTO pf.log (version_id, pf_user, operation, slice, params, note)
|
|
VALUES ({{version_id}}, '{{pf_user}}', 'reference', NULL, '{{params}}'::jsonb, '{{note}}')
|
|
RETURNING id
|
|
)
|
|
,ins AS (
|
|
INSERT INTO {{fc_table}} (${loadInsertCols})
|
|
SELECT
|
|
${buildLoadSelect(hasDimPeriod ? 's.' : '')},
|
|
'reference', (SELECT id FROM ilog), '{{pf_user}}', now()
|
|
FROM ${buildFromClause()}
|
|
WHERE {{filter_clause}}
|
|
RETURNING *
|
|
)
|
|
SELECT count(*) AS rows_affected FROM ins`.trim();
|
|
}
|
|
|
|
function buildScale() {
|
|
const vSel = effectiveValue
|
|
? `round((${q(effectiveValue)} / NULLIF(total_value, 0)) * {{value_incr}}, 2)`
|
|
: `0`;
|
|
const uSel = effectiveUnits
|
|
? `round((${q(effectiveUnits)} / NULLIF(total_units, 0)) * {{units_incr}}, 5)`
|
|
: `0`;
|
|
const baseSelectParts = [
|
|
...dimsJoined ? [dimsJoined] : [],
|
|
q(dateCol),
|
|
effectiveValue ? q(effectiveValue) : null,
|
|
effectiveUnits ? q(effectiveUnits) : null,
|
|
effectiveValue ? `sum(${q(effectiveValue)}) OVER () AS total_value` : null,
|
|
effectiveUnits ? `sum(${q(effectiveUnits)}) OVER () AS total_units` : null
|
|
].filter(Boolean).join(',\n ');
|
|
return `
|
|
WITH
|
|
ilog AS (
|
|
INSERT INTO pf.log (version_id, pf_user, operation, slice, params, note)
|
|
VALUES ({{version_id}}, '{{pf_user}}', 'scale', '{{slice}}'::jsonb, '{{params}}'::jsonb, '{{note}}')
|
|
RETURNING id
|
|
)
|
|
,base AS (
|
|
SELECT
|
|
${baseSelectParts}
|
|
FROM {{fc_table}}
|
|
WHERE {{where_clause}}
|
|
{{exclude_clause}}
|
|
)
|
|
,ins AS (
|
|
INSERT INTO {{fc_table}} (${insertCols})
|
|
SELECT
|
|
${[dimsJoined, q(dateCol), ...(effectiveValue ? [vSel] : []), ...(effectiveUnits ? [uSel] : [])].join(',\n ')},
|
|
'scale', (SELECT id FROM ilog), '{{pf_user}}', now()
|
|
FROM base
|
|
RETURNING *
|
|
)
|
|
${opTail('ins')}`.trim();
|
|
}
|
|
|
|
function buildRecode() {
|
|
return `
|
|
WITH
|
|
ilog AS (
|
|
INSERT INTO pf.log (version_id, pf_user, operation, slice, params, note)
|
|
VALUES ({{version_id}}, '{{pf_user}}', 'recode', '{{slice}}'::jsonb, '{{params}}'::jsonb, '{{note}}')
|
|
RETURNING id
|
|
)
|
|
,src AS (
|
|
SELECT ${selectData}
|
|
FROM {{fc_table}}
|
|
WHERE {{where_clause}}
|
|
{{exclude_clause}}
|
|
)
|
|
,neg AS (
|
|
INSERT INTO {{fc_table}} (${insertCols})
|
|
SELECT ${dimsJoined}, ${q(dateCol)}, ${effectiveValue ? `-${q(effectiveValue)}` : '0'}${effectiveUnits ? `, -${q(effectiveUnits)}` : ''},
|
|
'recode', (SELECT id FROM ilog), '{{pf_user}}', now()
|
|
FROM src
|
|
RETURNING *
|
|
)
|
|
,ins AS (
|
|
INSERT INTO {{fc_table}} (${insertCols})
|
|
SELECT {{set_clause}}, ${q(dateCol)}, ${effectiveValue ? q(effectiveValue) : '0'}${effectiveUnits ? `, ${q(effectiveUnits)}` : ''},
|
|
'recode', (SELECT id FROM ilog), '{{pf_user}}', now()
|
|
FROM src
|
|
RETURNING *
|
|
)
|
|
${grain ? `,allrows AS (
|
|
SELECT * FROM neg
|
|
UNION ALL
|
|
SELECT * FROM ins
|
|
)
|
|
${opTail('allrows')}` : 'SELECT * FROM neg UNION ALL SELECT * FROM ins'}`.trim();
|
|
}
|
|
|
|
function buildClone() {
|
|
return `
|
|
WITH
|
|
ilog AS (
|
|
INSERT INTO pf.log (version_id, pf_user, operation, slice, params, note)
|
|
VALUES ({{version_id}}, '{{pf_user}}', 'clone', '{{slice}}'::jsonb, '{{params}}'::jsonb, '{{note}}')
|
|
RETURNING id
|
|
)
|
|
,ins AS (
|
|
INSERT INTO {{fc_table}} (${insertCols})
|
|
SELECT
|
|
{{set_clause}},
|
|
${q(dateCol)},
|
|
${effectiveValue ? `round(${q(effectiveValue)} * {{scale_factor}}, 2)` : '0'}${effectiveUnits ? `,\n round(${q(effectiveUnits)} * {{scale_factor}}, 5)` : ''},
|
|
'clone', (SELECT id FROM ilog), '{{pf_user}}', now()
|
|
FROM {{fc_table}}
|
|
WHERE {{where_clause}}
|
|
{{exclude_clause}}
|
|
RETURNING *
|
|
)
|
|
${opTail('ins')}`.trim();
|
|
}
|
|
|
|
function buildUndo() {
|
|
// undo is executed as two separate queries in the route handler
|
|
// (delete from fc_table first, then delete from pf.log) to avoid
|
|
// FK constraint ordering issues within a single CTE statement.
|
|
// This entry is a placeholder — the undo route uses it as a template reference.
|
|
return `
|
|
-- step 1 (run first):
|
|
DELETE FROM {{fc_table}} WHERE pf_logid = {{logid}};
|
|
-- step 2 (run after step 1):
|
|
DELETE FROM pf.log WHERE id = {{logid}};`.trim();
|
|
}
|
|
}
|
|
|
|
// substitute {{token}} placeholders in a SQL string
|
|
function applyTokens(sql, tokens) {
|
|
let result = sql;
|
|
for (const [key, value] of Object.entries(tokens)) {
|
|
result = result.replace(new RegExp(`\\{\\{${key}\\}\\}`, 'g'), value ?? '');
|
|
}
|
|
return result;
|
|
}
|
|
|
|
// build a SQL WHERE clause string from a slice object
|
|
// only dimension columns are included; unrecognised keys are silently skipped
|
|
function buildWhere(slice, dimCols) {
|
|
if (!slice || Object.keys(slice).length === 0) return 'TRUE';
|
|
|
|
const allowed = new Set(dimCols);
|
|
const parts = [];
|
|
|
|
for (const [col, val] of Object.entries(slice)) {
|
|
if (!allowed.has(col)) continue;
|
|
if (Array.isArray(val)) {
|
|
const escaped = val.map(v => esc(v));
|
|
parts.push(`"${col}" IN ('${escaped.join("', '")}')`);
|
|
} else {
|
|
parts.push(`"${col}" = '${esc(val)}'`);
|
|
}
|
|
}
|
|
|
|
return parts.length ? parts.join('\nAND ') : 'TRUE';
|
|
}
|
|
|
|
// build a WHERE clause spanning several slices — an OR of AND-groups.
|
|
// A union of slices cannot be flattened into one IN list per column: slices
|
|
// {Region:East, State:NY} and {Region:West, State:CA} would become
|
|
// Region IN (East,West) AND State IN (NY,CA), which also matches East/CA.
|
|
function buildWhereAny(slices, dimCols) {
|
|
const list = (slices || []).filter(s => s && Object.keys(s).length > 0);
|
|
if (list.length === 0) return 'TRUE';
|
|
if (list.length === 1) return buildWhere(list[0], dimCols);
|
|
|
|
const groups = list
|
|
.map(s => buildWhere(s, dimCols))
|
|
.filter(w => w !== 'TRUE');
|
|
|
|
// any slice that reduced to TRUE selects everything, so the union does too
|
|
if (groups.length !== list.length) return 'TRUE';
|
|
|
|
// outer parens matter: the caller appends `AND pf_iter NOT IN (...)`,
|
|
// and AND binds tighter than OR
|
|
return `(${groups.map(g => `(${g.replace(/\n/g, ' ')})`).join('\n OR ')})`;
|
|
}
|
|
|
|
// the bare predicate for "this row participates in operations", for use in a
|
|
// FILTER clause where the excluded rows still need to be counted separately
|
|
function buildExcludePredicate(excludeIters) {
|
|
if (!excludeIters || excludeIters.length === 0) return 'TRUE';
|
|
const list = excludeIters.map(i => `'${esc(i)}'`).join(', ');
|
|
return `pf_iter NOT IN (${list})`;
|
|
}
|
|
|
|
// build AND iter NOT IN (...) from a version's exclude_iters array
|
|
function buildExcludeClause(excludeIters) {
|
|
if (!excludeIters || excludeIters.length === 0) return '';
|
|
const list = excludeIters.map(i => `'${esc(i)}'`).join(', ');
|
|
return `AND pf_iter NOT IN (${list})`;
|
|
}
|
|
|
|
// build the dimension columns portion of a SELECT for recode/clone
|
|
// replaces named dimensions with literal values, passes others through unchanged
|
|
function buildSetClause(dimCols, setObj) {
|
|
return dimCols.map(col => {
|
|
if (setObj && setObj[col] !== undefined) {
|
|
return `'${esc(setObj[col])}' AS "${col}"`;
|
|
}
|
|
return `"${col}"`;
|
|
}).join(', ');
|
|
}
|
|
|
|
// build a SQL WHERE clause from an array of filter objects { col, op, values }
|
|
// only allows columns with role 'date' or 'filter'
|
|
function buildFilterClause(filters, colMeta) {
|
|
if (!filters || filters.length === 0) {
|
|
const err = new Error('At least one filter is required');
|
|
err.status = 400; throw err;
|
|
}
|
|
const allowed = new Set(
|
|
colMeta.filter(c => c.role !== 'ignore').map(c => c.cname)
|
|
);
|
|
const parts = filters.map(({ col, op, values = [] }) => {
|
|
if (!allowed.has(col)) {
|
|
const err = new Error(`Column "${col}" is not available for baseline filtering`);
|
|
err.status = 400; throw err;
|
|
}
|
|
const c = `"${col}"`;
|
|
const v = values.map(x => `'${esc(String(x))}'`);
|
|
switch (op) {
|
|
case '=': return `${c} = ${v[0]}`;
|
|
case '!=': return `${c} != ${v[0]}`;
|
|
case 'IN': return `${c} IN (${v.join(', ')})`;
|
|
case 'NOT IN': return `${c} NOT IN (${v.join(', ')})`;
|
|
case 'BETWEEN': return `${c} BETWEEN ${v[0]} AND ${v[1]}`;
|
|
case 'IS NULL': return `${c} IS NULL`;
|
|
case 'IS NOT NULL': return `${c} IS NOT NULL`;
|
|
default: {
|
|
const err = new Error(`Unsupported operator "${op}"`);
|
|
err.status = 400; throw err;
|
|
}
|
|
}
|
|
});
|
|
return parts.join('\nAND ');
|
|
}
|
|
|
|
// escape a value for safe SQL string substitution
|
|
function esc(val) {
|
|
if (val === null || val === undefined) return '';
|
|
return String(val).replace(/'/g, "''");
|
|
}
|
|
|
|
module.exports = { generateSQL, grainOf, applyTokens, buildWhere, buildWhereAny, buildExcludeClause, buildExcludePredicate, buildSetClause, buildFilterClause, esc };
|