Territory filtering was deferred from v1, so every account saw and could write every row. With the sales team about to adjust their own territories that is the thing standing in the way, and it is also why a rep would wait fifteen seconds to load 2.7M rows to work on a few thousand. The list lives on pf.app_user.territory with is_admin beside it, and col_meta.is_territory marks which column of a source the values belong to -- flagged rather than named in code, so a second source can be divided by something other than a sales rep. Fail closed: buildTerritoryClause returns FALSE for an empty list or an unflagged source. An account nobody configured sees nothing, rather than everything because a column was left null. Built from the session, never the request. That is what separates it from `scope`, which the browser sends and should: a filter the user chose belongs in the payload, a permission cannot come from the thing it restrains. It is ANDed on last, where nothing in the request can undo it. Enforced on /data (the cursor and the count behind X-Row-Count), on /agg before the GROUP BY since the territory column need not be in the grain, on every operation through sliceUnits, and on the value completion endpoint -- which reads the source table, so without it a dropdown enumerates every customer and rep in the business to someone shown none of their rows. Undo is gated by owner rather than territory: it removes an entry's rows wholesale, so half-undoing one would leave a state nothing describes. Recode refuses to set the territory column unless you are an admin, since moving a row between territories is reassignment, not forecasting. ./pf.sh gains set-territory, set-admin and orphan-territory. The last lists territory values no account owns -- work under one is invisible to everybody but an admin, which a typo causes easily and nothing in the app reveals. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
496 lines
21 KiB
JavaScript
496 lines
21 KiB
JavaScript
const express = require('express');
|
|
const { generateSQL, buildTerritoryClause } = require('../lib/sql_generator');
|
|
const { sessionTerritory } = require('../lib/auth');
|
|
const { sessionUser } = require('../lib/auth');
|
|
|
|
module.exports = function(pool) {
|
|
const router = express.Router();
|
|
|
|
// list all registered sources
|
|
router.get('/sources', async (req, res) => {
|
|
try {
|
|
const result = await pool.query(
|
|
`SELECT * FROM pf.source ORDER BY schema, tname`
|
|
);
|
|
res.json(result.rows);
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// register a source table
|
|
// auto-populates col_meta from information_schema with role='ignore'
|
|
router.post('/sources', async (req, res) => {
|
|
const { schema, tname, label } = req.body;
|
|
const created_by = sessionUser(req);
|
|
if (!schema || !tname) {
|
|
return res.status(400).json({ error: 'schema and tname are required' });
|
|
}
|
|
if (!/^\w+$/.test(schema) || !/^\w+$/.test(tname)) {
|
|
return res.status(400).json({ error: 'Invalid schema or table name' });
|
|
}
|
|
const client = await pool.connect();
|
|
try {
|
|
await client.query('BEGIN');
|
|
|
|
const src = await client.query(
|
|
`INSERT INTO pf.source (schema, tname, label, created_by)
|
|
VALUES ($1, $2, $3, $4)
|
|
RETURNING *`,
|
|
[schema, tname, label || null, created_by || null]
|
|
);
|
|
const source = src.rows[0];
|
|
|
|
// seed col_meta from information_schema
|
|
await client.query(`
|
|
INSERT INTO pf.col_meta (source_id, cname, role, opos)
|
|
SELECT $1, column_name, 'dimension', ordinal_position
|
|
FROM information_schema.columns
|
|
WHERE table_schema = $2 AND table_name = $3
|
|
ORDER BY ordinal_position
|
|
ON CONFLICT (source_id, cname) DO NOTHING
|
|
`, [source.id, schema, tname]);
|
|
|
|
await client.query('COMMIT');
|
|
res.status(201).json(source);
|
|
} catch (err) {
|
|
await client.query('ROLLBACK');
|
|
console.error(err);
|
|
if (err.code === '23505') {
|
|
return res.status(409).json({ error: 'Source already registered' });
|
|
}
|
|
res.status(500).json({ error: err.message });
|
|
} finally {
|
|
client.release();
|
|
}
|
|
});
|
|
|
|
// get col_meta for a source
|
|
router.get('/sources/:id/cols', async (req, res) => {
|
|
try {
|
|
const result = await pool.query(
|
|
`SELECT * FROM pf.col_meta WHERE source_id = $1 ORDER BY opos`,
|
|
[req.params.id]
|
|
);
|
|
res.json(result.rows);
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// save col_meta — accepts full array, upserts each row
|
|
router.put('/sources/:id/cols', async (req, res) => {
|
|
const sourceId = parseInt(req.params.id);
|
|
const cols = req.body;
|
|
if (!Array.isArray(cols)) {
|
|
return res.status(400).json({ error: 'body must be an array' });
|
|
}
|
|
const client = await pool.connect();
|
|
try {
|
|
await client.query('BEGIN');
|
|
for (const col of cols) {
|
|
await client.query(`
|
|
INSERT INTO pf.col_meta (source_id, cname, label, role, is_key, dim_group, dim_period_col, in_grain, opos)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
|
ON CONFLICT (source_id, cname) DO UPDATE SET
|
|
label = EXCLUDED.label,
|
|
role = EXCLUDED.role,
|
|
is_key = EXCLUDED.is_key,
|
|
dim_group = EXCLUDED.dim_group,
|
|
dim_period_col = EXCLUDED.dim_period_col,
|
|
in_grain = EXCLUDED.in_grain,
|
|
opos = EXCLUDED.opos
|
|
`, [
|
|
sourceId,
|
|
col.cname,
|
|
col.label || null,
|
|
col.role || 'ignore',
|
|
col.is_key || false,
|
|
col.dim_group || null,
|
|
col.dim_period_col || null,
|
|
col.in_grain || false,
|
|
col.opos || null
|
|
]);
|
|
}
|
|
await client.query('COMMIT');
|
|
const result = await pool.query(
|
|
`SELECT * FROM pf.col_meta WHERE source_id = $1 ORDER BY opos`,
|
|
[sourceId]
|
|
);
|
|
res.json(result.rows);
|
|
} catch (err) {
|
|
await client.query('ROLLBACK');
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
} finally {
|
|
client.release();
|
|
}
|
|
});
|
|
|
|
// generate SQL for all operations from current col_meta and store in pf.sql
|
|
router.post('/sources/:id/generate-sql', async (req, res) => {
|
|
const sourceId = parseInt(req.params.id);
|
|
try {
|
|
const srcResult = await pool.query(
|
|
`SELECT * FROM pf.source WHERE id = $1`, [sourceId]
|
|
);
|
|
if (srcResult.rows.length === 0) {
|
|
return res.status(404).json({ error: 'Source not found' });
|
|
}
|
|
|
|
const colResult = await pool.query(
|
|
`SELECT * FROM pf.col_meta WHERE source_id = $1 ORDER BY opos`,
|
|
[sourceId]
|
|
);
|
|
|
|
// validate required roles
|
|
const colMeta = colResult.rows;
|
|
const roles = new Set(colMeta.map(c => c.role));
|
|
const missing = ['value', 'date'].filter(r => !roles.has(r));
|
|
if (missing.length > 0) {
|
|
return res.status(400).json({
|
|
error: `col_meta is missing required roles: ${missing.join(', ')}`
|
|
});
|
|
}
|
|
if (!colMeta.some(c => c.role === 'dimension')) {
|
|
return res.status(400).json({ error: 'col_meta has no dimension columns' });
|
|
}
|
|
|
|
const sqls = generateSQL(srcResult.rows[0], colMeta);
|
|
const client = await pool.connect();
|
|
try {
|
|
await client.query('BEGIN');
|
|
for (const [operation, sql] of Object.entries(sqls)) {
|
|
await client.query(`
|
|
INSERT INTO pf.sql (source_id, operation, sql, generated_at)
|
|
VALUES ($1, $2, $3, now())
|
|
ON CONFLICT (source_id, operation) DO UPDATE SET
|
|
sql = EXCLUDED.sql,
|
|
generated_at = EXCLUDED.generated_at
|
|
`, [sourceId, operation, sql]);
|
|
}
|
|
// drop operations this generation no longer produces — e.g. get_agg
|
|
// after the grain has been cleared, which would otherwise leave a
|
|
// stale template the load path would still pick up
|
|
await client.query(
|
|
`DELETE FROM pf.sql WHERE source_id = $1 AND operation <> ALL($2::text[])`,
|
|
[sourceId, Object.keys(sqls)]
|
|
);
|
|
await client.query('COMMIT');
|
|
} catch (err) {
|
|
await client.query('ROLLBACK');
|
|
throw err;
|
|
} finally {
|
|
client.release();
|
|
}
|
|
|
|
res.json({ message: 'SQL generated', operations: Object.keys(sqls) });
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// view generated SQL for a source (for inspection / debug)
|
|
router.get('/sources/:id/sql', async (req, res) => {
|
|
try {
|
|
const result = await pool.query(
|
|
`SELECT operation, sql, generated_at
|
|
FROM pf.sql WHERE source_id = $1 ORDER BY operation`,
|
|
[req.params.id]
|
|
);
|
|
res.json(result.rows);
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// get distinct values for a key column (used to populate operation panel dropdowns)
|
|
router.get('/sources/:id/values/:col', async (req, res) => {
|
|
const col = req.params.col;
|
|
try {
|
|
const srcResult = await pool.query(
|
|
`SELECT schema, tname FROM pf.source WHERE id = $1`, [req.params.id]
|
|
);
|
|
if (srcResult.rows.length === 0) return res.status(404).json({ error: 'Source not found' });
|
|
|
|
// validate col is a key dimension on this source
|
|
const metaResult = await pool.query(
|
|
`SELECT 1 FROM pf.col_meta WHERE source_id = $1 AND cname = $2 AND is_key = true`,
|
|
[req.params.id, col]
|
|
);
|
|
if (metaResult.rows.length === 0) {
|
|
return res.status(400).json({ error: `"${col}" is not a key column` });
|
|
}
|
|
|
|
// ?q= narrows, ?limit= caps. A key column can be very wide -- part on
|
|
// osm_skinny has 11,290 distinct values -- so returning the lot to fill a
|
|
// completion list is both a slow query and a large response for a control
|
|
// that can only usefully show a handful.
|
|
const { schema, tname } = srcResult.rows[0];
|
|
const q = (req.query.q || '').trim();
|
|
const limit = Math.min(parseInt(req.query.limit) || 5000, 5000);
|
|
|
|
const params = [];
|
|
let filter = `WHERE "${col}" IS NOT NULL`;
|
|
|
|
// Completion reads the *source* table, which no territory scope has
|
|
// touched -- so without this a scoped account could enumerate every
|
|
// customer, part and rep in the business from a dropdown, having
|
|
// been shown none of their rows.
|
|
const terrRow = (await pool.query(
|
|
`SELECT cname FROM pf.col_meta WHERE source_id = $1 AND is_territory LIMIT 1`,
|
|
[req.params.id]
|
|
)).rows[0];
|
|
const terrClause = buildTerritoryClause(sessionTerritory(req), terrRow?.cname || null);
|
|
if (terrClause) filter += ` AND ${terrClause}`;
|
|
|
|
if (q) {
|
|
params.push(`%${q}%`);
|
|
filter += ` AND "${col}"::text ILIKE $${params.length}`;
|
|
}
|
|
params.push(limit);
|
|
|
|
const result = await pool.query(
|
|
`SELECT DISTINCT "${col}"::text AS val FROM ${schema}.${tname}
|
|
${filter} ORDER BY 1 LIMIT $${params.length}`,
|
|
params
|
|
);
|
|
res.json(result.rows.map(r => r.val));
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// Resolve a dim_group to its key column and siblings, or explain why it cannot be.
|
|
async function resolveGroup(sourceId, group) {
|
|
const { rows: meta } = await pool.query(
|
|
`SELECT * FROM pf.col_meta WHERE source_id = $1 ORDER BY opos`, [sourceId]);
|
|
const members = meta.filter(c => c.dim_group === group);
|
|
if (!members.length) {
|
|
const err = new Error(`No columns are grouped as "${group}" on this source`);
|
|
err.status = 404; throw err;
|
|
}
|
|
const keyCol = members.find(c => c.is_key);
|
|
if (!keyCol) {
|
|
const err = new Error(
|
|
`Group "${group}" has no is_key column, so its members have nothing to be keyed on`);
|
|
err.status = 400; throw err;
|
|
}
|
|
return {
|
|
keyCol,
|
|
siblings: members.filter(c => c.cname !== keyCol.cname),
|
|
// recency column: the source's primary date, the same one the generator
|
|
// treats as the date for loads
|
|
dateCol: meta.find(c => c.role === 'date')?.cname || null,
|
|
};
|
|
}
|
|
|
|
// The member list for a group, as one array. Small enough to send whole --
|
|
// 11,290 parts on osm_skinny -- so the client holds it and filters locally
|
|
// instead of querying per keystroke.
|
|
router.get('/sources/:id/dim/:group', async (req, res) => {
|
|
try {
|
|
const sourceId = parseInt(req.params.id);
|
|
const { keyCol, siblings } = await resolveGroup(sourceId, req.params.group);
|
|
const includeInactive = req.query.all === '1';
|
|
|
|
const { rows } = await pool.query(`
|
|
SELECT key_value, attrs, is_active, source_seen
|
|
FROM pf.dim_member
|
|
WHERE source_id = $1 AND dim_group = $2
|
|
${includeInactive ? '' : 'AND is_active'}
|
|
ORDER BY key_value
|
|
`, [sourceId, req.params.group]);
|
|
|
|
res.json({
|
|
group: req.params.group,
|
|
key_col: keyCol.cname,
|
|
siblings: siblings.map(c => c.cname),
|
|
members: rows,
|
|
});
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(err.status || 500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// Rebuild a group's members from the source. A merge, not a replace: curation
|
|
// (a member deactivated by hand, or added before it ever sold) has to survive a
|
|
// refresh, so absent members are marked source_seen = false rather than deleted.
|
|
//
|
|
// Slow by nature -- it reads the whole source, which for a view over a
|
|
// transaction table is millions of rows -- so it is a deliberate action rather
|
|
// than something that happens on a page load.
|
|
router.post('/sources/:id/dim/:group/refresh', async (req, res) => {
|
|
const sourceId = parseInt(req.params.id);
|
|
const group = req.params.group;
|
|
try {
|
|
const srcResult = await pool.query(
|
|
`SELECT schema, tname FROM pf.source WHERE id = $1`, [sourceId]);
|
|
if (!srcResult.rows.length) return res.status(404).json({ error: 'Source not found' });
|
|
const { schema, tname } = srcResult.rows[0];
|
|
|
|
const { keyCol, siblings, dateCol } = await resolveGroup(sourceId, group);
|
|
if (!siblings.length) {
|
|
return res.status(400).json({ error: `Group "${group}" has no sibling columns to store` });
|
|
}
|
|
|
|
const q = (n) => `"${n}"`;
|
|
const attrs = siblings.map(c => `'${c.cname}', s.${q(c.cname)}::text`).join(', ');
|
|
// A key can carry more than one attribute set across history -- 11,290
|
|
// parts against 13,662 combinations on osm_skinny. Take the most recent
|
|
// by the source's date column, which is the live definition.
|
|
const recency = dateCol ? `s.${q(dateCol)} DESC NULLS LAST` : `1`;
|
|
|
|
const started = Date.now();
|
|
// An explicit stamp rather than now(): inside a transaction now() is the
|
|
// transaction's start time, so "refreshed in this run" and "refreshed in
|
|
// a run that began at the same instant" would be indistinguishable.
|
|
const runAt = new Date();
|
|
const client = await pool.connect();
|
|
try {
|
|
await client.query('BEGIN');
|
|
const { rows: [{ n }] } = await client.query(`
|
|
WITH ranked AS (
|
|
SELECT s.${q(keyCol.cname)}::text AS key_value,
|
|
jsonb_build_object(${attrs}) AS attrs,
|
|
row_number() OVER (
|
|
PARTITION BY s.${q(keyCol.cname)} ORDER BY ${recency}
|
|
) AS rn
|
|
FROM ${q(schema)}.${q(tname)} s
|
|
WHERE s.${q(keyCol.cname)} IS NOT NULL
|
|
)
|
|
,upserted AS (
|
|
INSERT INTO pf.dim_member
|
|
(source_id, dim_group, key_value, attrs, refreshed_at, source_seen)
|
|
SELECT $1, $2, key_value, attrs, $3, true
|
|
FROM ranked WHERE rn = 1
|
|
ON CONFLICT (source_id, dim_group, key_value) DO UPDATE SET
|
|
attrs = EXCLUDED.attrs,
|
|
source_seen = true,
|
|
refreshed_at = $3,
|
|
updated_at = now()
|
|
RETURNING 1
|
|
)
|
|
SELECT count(*)::int AS n FROM upserted
|
|
`, [sourceId, group, runAt]);
|
|
|
|
// anything this run did not touch is no longer in the source
|
|
const { rowCount: dropped } = await client.query(`
|
|
UPDATE pf.dim_member
|
|
SET source_seen = false, updated_at = now()
|
|
WHERE source_id = $1 AND dim_group = $2 AND source_seen
|
|
AND refreshed_at IS DISTINCT FROM $3
|
|
`, [sourceId, group, runAt]);
|
|
|
|
await client.query('COMMIT');
|
|
res.json({ group, members: n, no_longer_in_source: dropped, ms: Date.now() - started });
|
|
} catch (err) {
|
|
await client.query('ROLLBACK');
|
|
throw err;
|
|
} finally {
|
|
client.release();
|
|
}
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(err.status || 500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// given a key column value, look up sibling dim_group column values from source
|
|
// returns { sibling_col: value, ... } if exactly one match, null if none or ambiguous
|
|
router.get('/sources/:id/lookup', async (req, res) => {
|
|
const { col, value } = req.query;
|
|
if (!col || value == null || value === '') return res.json(null);
|
|
try {
|
|
const [srcResult, metaResult] = await Promise.all([
|
|
pool.query(`SELECT schema, tname FROM pf.source WHERE id = $1`, [req.params.id]),
|
|
pool.query(`SELECT * FROM pf.col_meta WHERE source_id = $1 ORDER BY opos`, [req.params.id])
|
|
]);
|
|
if (srcResult.rows.length === 0) return res.status(404).json({ error: 'Source not found' });
|
|
|
|
const keyCol = metaResult.rows.find(c => c.cname === col && c.is_key && c.dim_group);
|
|
if (!keyCol) return res.json(null);
|
|
|
|
const siblings = metaResult.rows.filter(c =>
|
|
c.dim_group === keyCol.dim_group && c.cname !== col
|
|
);
|
|
if (!siblings.length) return res.json(null);
|
|
|
|
const { schema, tname } = srcResult.rows[0];
|
|
const sibCols = siblings.map(c => `"${c.cname}"`).join(', ');
|
|
const result = await pool.query(
|
|
`SELECT DISTINCT ${sibCols} FROM "${schema}"."${tname}" WHERE "${col}" = $1 LIMIT 2`,
|
|
[value]
|
|
);
|
|
|
|
if (result.rows.length !== 1) return res.json(null);
|
|
res.json(result.rows[0]);
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// set or clear the default Perspective layout for a source.
|
|
// Body: a Perspective view config (group_by, split_by, columns, plugin_config, …).
|
|
// Pass null or {} to clear.
|
|
router.put('/sources/:id/default-layout', async (req, res) => {
|
|
try {
|
|
const layout = req.body && Object.keys(req.body).length > 0 ? req.body : null;
|
|
const result = await pool.query(
|
|
`UPDATE pf.source SET default_layout = $1 WHERE id = $2 RETURNING *`,
|
|
[layout, req.params.id]
|
|
);
|
|
if (result.rows.length === 0) return res.status(404).json({ error: 'Source not found' });
|
|
res.json(result.rows[0]);
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
// deregister a source — does not drop existing forecast tables
|
|
router.get('/dim-period/cols', async (req, res) => {
|
|
try {
|
|
const result = await pool.query(`
|
|
SELECT column_name
|
|
FROM information_schema.columns
|
|
WHERE table_schema = 'pf' AND table_name = 'dim_period'
|
|
AND column_name NOT IN ('sdat', 'edat', 'drange', 'ndays')
|
|
ORDER BY ordinal_position
|
|
`);
|
|
res.json(result.rows.map(r => r.column_name));
|
|
} catch (err) {
|
|
console.error(err);
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
router.delete('/sources/:id', async (req, res) => {
|
|
try {
|
|
const result = await pool.query(
|
|
`DELETE FROM pf.source WHERE id = $1 RETURNING *`,
|
|
[req.params.id]
|
|
);
|
|
if (result.rows.length === 0) {
|
|
return res.status(404).json({ error: 'Source not found' });
|
|
}
|
|
res.json({ message: 'Source deregistered', source: result.rows[0] });
|
|
} catch (err) {
|
|
console.error(err);
|
|
if (err.code === '23001' || err.code === '23503') {
|
|
return res.status(409).json({ error: 'Source has existing versions — delete them first.' });
|
|
}
|
|
res.status(500).json({ error: err.message });
|
|
}
|
|
});
|
|
|
|
return router;
|
|
};
|