Compare commits

..

No commits in common. "feature/migrations" and "master" have entirely different histories.

20 changed files with 137 additions and 1190 deletions

View File

@ -16,7 +16,6 @@ Data transport architecture options: `pf_perspective_options.md`
- **Frontend:** React + Vite + Tailwind CSS in `ui/`; built output lands in `public/app/`
- **Pivot:** [Perspective](https://github.com/perspective-dev/perspective) (`@perspective-dev/*` distribution, **not** FINOS `@finos/perspective`) 4.4.0 loaded from CDN at runtime — see `PERSPECTIVE.md` for config/deploy guidance
- **Dev:** `npm run dev` (nodemon) in root; `npm run build` in `ui/`
- **Schema:** forward-only SQL migrations in `setup_sql/migrations/`, applied by `npm run migrate` and tracked in `pf.schema_version`. `server.js` refuses to start when the database is behind. See `setup_sql/README.md`
---
@ -34,13 +33,7 @@ lib/
sql_generator.js buildFilterClause, token substitution helpers
utils.js
setup_sql/
README.md migration workflow — read before changing the schema
migrations/ ordered .sql, applied once, tracked in pf.schema_version
gen_dim_period.sql parameterised calendar load (not a migration)
schema.generated.sql pg_dump reference snapshot; generated, never edited
scripts/
migrate.js migration CLI (up | status | baseline)
schema-dump.js regenerates schema.generated.sql
01_schema.sql pf schema DDL — run once to install
ui/src/
views/
Setup.jsx DB browser, source registration, col_meta editor
@ -56,13 +49,12 @@ ui/src/
## Database schema (`pf`)
- **`pf.source`** — registered source tables
- **`pf.col_meta`** — column roles: `dimension` | `value` | `units` | `date` | `filter` | `ignore`; `is_key` marks dimensions used in slice WHERE clauses; `dim_group` groups functionally dependent columns (e.g. date + its derived year/month dimensions); `dim_period_col` maps a dimension to a `pf.dim_period` column so date-adjacent values are derived at load time rather than copied raw; `in_grain` flags dimension/date columns that define the **display grain** (see below)
- **`pf.col_meta`** — column roles: `dimension` | `value` | `units` | `date` | `filter` | `ignore`; `is_key` marks dimensions used in slice WHERE clauses; `dim_group` groups functionally dependent columns (e.g. date + its derived year/month dimensions); `dim_period_col` maps a dimension to a `pf.dim_period` column so date-adjacent values are derived at load time rather than copied raw
- **`pf.version`** — named forecast scenarios; `exclude_iters` (default `["reference"]`) blocks those iter values from all operations
- **`pf.fc_{tname}_{version_id}`** — one forecast table per version; contains both operational rows (`pf_iter = baseline|scale|recode|clone`) and reference rows (`pf_iter = reference`)
- **`pf.log`** — audit log; every write gets one entry; `slice` + `params` stored as jsonb
- **`pf.sql`** — generated SQL templates per source/operation; tokens substituted at request time
- **`pf.dim_period`** — calendar lookup table (20182035); one row per month keyed on `sdat` (month start date); provides cal/fiscal year, quarter, and month columns; populated by `setup_sql/gen_dim_period.sql` with a configurable fiscal year start month
- **`pf.schema_version`** — applied migrations (filename + checksum); owned by the migration runner, never edited by hand
### Key token substitution tokens
`{{fc_table}}`, `{{where_clause}}`, `{{exclude_clause}}`, `{{logid}}`, `{{pf_user}}`, `{{value_incr}}`, `{{units_incr}}`, `{{pct}}`, `{{set_clause}}`, `{{scale_factor}}`, `{{date_offset}}`, `{{filter_clause}}`
@ -72,23 +64,15 @@ ui/src/
## Core data flow
### Initial load (Forecast view)
`Forecast.jsx` fetches col_meta first, then picks the endpoint:
- **grain mode** (any `in_grain` column) — `GET /api/versions/:id/agg`, rows pre-aggregated to the grain, table indexed on `pf_gkey`
- **raw mode** (no grain) — `GET /api/versions/:id/data`, raw forecast rows, table indexed on `pf_id`
Either way: Arrow IPC binary stream → `worker.table(buffer)` in Perspective WASM. `fetchArrow()` handles both.
`GET /api/versions/:id/data` → Arrow IPC binary stream → `worker.table(buffer)` in Perspective WASM
**Why one batch (not streaming):** pg returns `bigint`/`numeric` as strings by default — type parsers in `server.js` coerce them to numbers. Per-batch Arrow encoding creates independent dictionaries that cause Perspective WASM to crash on dictionary replacement messages. Server accumulates all rows, emits one record batch.
### Display grain
Aggregating to the grain the pivot actually displays is the load-time fix — measured 534,902 → 6,154 rows on `osm_stack`. It keeps the **native** Perspective engine, so expand/collapse/depth/sort/filter all still work. Set the grain in Setup (`in_grain` per column); it is baked into `pf.sql` at Generate SQL time so load and operations agree. `grainOf()` in `lib/sql_generator.js` is the single definition of what the grain is — `Setup.jsx` and `routes/log.js` mirror it. Full design: `pf_spec.md` → §Display-grain pre-aggregation. Why not a DuckDB virtual server: `pf_perspective_options.md` → §Spike findings.
### Forecast operations
POST to `/api/versions/:id/{scale|recode|clone}` → SQL executed with `RETURNING *` → new rows returned as JSON → `pspTable.update(rows)` — no full reload. In grain mode the operation's final CTE aggregates its own new rows to grain first; since `pf_logid` is part of `pf_gkey` those keys are always new, so `update()` **appends** and the view re-sums.
POST to `/api/versions/:id/{scale|recode|clone}` → SQL executed with `RETURNING *` → new rows returned as JSON → `pspTable.update(rows)` — no full reload.
### Undo
`DELETE /api/log/:logid` → removes rows by logid → `table.remove()` of the affected index values (`pf_gkeys` in grain mode, `pf_ids` in raw mode); the view re-sums. No full reload.
`DELETE /api/log/:logid` → removes rows by logid → **full Perspective reload** (known wart).
---
@ -130,9 +114,6 @@ Theme state lives in `ui/src/theme.jsx` — a React context (`ThemeContext`) wit
- Default pivot layout should be configurable per source (currently hardcodes first 2 dimensions)
- Source/version selection doesn't persist across page reload
- Col_meta / version schema drift: if col_meta roles change after a version's forecast table is created, SQL and DDL go out of sync — workaround is to delete and recreate the version
- Grain drift: changing `in_grain` after a load requires Generate SQL + a page reload, since the loaded table's index and columns are fixed at load time. `routes/log.js` derives the grain from live col_meta, so a grain changed mid-session yields `pf_gkeys` that don't match the loaded table and undo silently removes nothing
- Migrations are forward-only — there are no down migrations. Rolling back a schema change means writing a new migration that reverses it
- Grain is static per source — a dimension left unflagged cannot be pivoted on. Dynamic per-cut grain (intersect the viewer's field set with the eligible set) is the additive next step; see `pf_spec.md` → §Display-grain pre-aggregation
## Deferred (not in v1)
Baseline replay (`replay: true` returns 501), approval workflow, territory filtering, export, version comparison, multi-DB connections. Live server-side aggregation (Path A / DuckDB virtual server) is parked on branch `spike/duckdb-virtual-server`.
Baseline replay (`replay: true` returns 501), approval workflow, territory filtering, export, version comparison, multi-DB connections.

View File

@ -1,158 +0,0 @@
// Migration bookkeeping shared by the runner (scripts/migrate.js) and the
// startup check in server.js.
//
// Migrations are plain .sql files in setup_sql/migrations, applied in filename
// order and recorded in pf.schema_version. They are immutable once applied: the
// runner stores a checksum and refuses to proceed if a file has changed, since
// editing an applied migration means databases silently disagree about what the
// schema is.
//
// pf.dim_period is deliberately not a migration — it is a parameterised data
// load (fiscal year start month), so it stays in setup_sql/gen_dim_period.sql.
const fs = require('fs');
const path = require('path');
const crypto = require('crypto');
const MIGRATIONS_DIR = path.join(__dirname, '..', 'setup_sql', 'migrations');
// Bookkeeping table has to exist before the first migration can be recorded, so
// the runner creates it rather than a migration doing so.
const BOOKKEEPING_DDL = `
CREATE SCHEMA IF NOT EXISTS pf;
CREATE TABLE IF NOT EXISTS pf.schema_version (
filename text PRIMARY KEY,
checksum text NOT NULL,
applied_at timestamptz NOT NULL DEFAULT now(),
applied_by text
);
`;
function checksum(sql) {
return crypto.createHash('sha256').update(sql).digest('hex').slice(0, 16);
}
// every .sql file on disk, in apply order
function readMigrations() {
if (!fs.existsSync(MIGRATIONS_DIR)) return [];
return fs.readdirSync(MIGRATIONS_DIR)
.filter(f => f.endsWith('.sql'))
.sort()
.map(filename => {
const sql = fs.readFileSync(path.join(MIGRATIONS_DIR, filename), 'utf8');
return { filename, sql, checksum: checksum(sql) };
});
}
async function ensureBookkeeping(client) {
await client.query(BOOKKEEPING_DDL);
}
async function readApplied(client) {
const { rows } = await client.query(
`SELECT filename, checksum, applied_at FROM pf.schema_version ORDER BY filename`
);
return rows;
}
// Compare disk against the database. Returns pending migrations plus any
// integrity problems worth refusing to run on.
async function status(client) {
await ensureBookkeeping(client);
const onDisk = readMigrations();
const applied = await readApplied(client);
const appliedBy = new Map(applied.map(r => [r.filename, r]));
const diskBy = new Map(onDisk.map(m => [m.filename, m]));
const pending = onDisk.filter(m => !appliedBy.has(m.filename));
// an applied file whose contents changed — the schema is now undefined
const modified = onDisk
.filter(m => appliedBy.has(m.filename) && appliedBy.get(m.filename).checksum !== m.checksum)
.map(m => m.filename);
// recorded as applied but no longer on disk — someone deleted or renamed it
const missing = applied.filter(r => !diskBy.has(r.filename)).map(r => r.filename);
// a pending migration ordered before an applied one would apply out of sequence
const lastApplied = applied.length ? applied[applied.length - 1].filename : null;
const outOfOrder = lastApplied ? pending.filter(m => m.filename < lastApplied).map(m => m.filename) : [];
return { onDisk, applied, pending, modified, missing, outOfOrder };
}
// Apply pending migrations, each in its own transaction so a failure leaves
// earlier ones committed and the failing one fully rolled back.
async function migrate(client, { user, log = console.log } = {}) {
const st = await status(client);
if (st.modified.length) {
throw new Error(
`Applied migrations have been modified: ${st.modified.join(', ')}. ` +
`Migrations are immutable — revert the edit and add a new migration instead.`
);
}
if (st.missing.length) {
throw new Error(
`Migrations recorded as applied are missing from disk: ${st.missing.join(', ')}.`
);
}
if (st.outOfOrder.length) {
throw new Error(
`Migrations would apply out of order: ${st.outOfOrder.join(', ')} sort before ` +
`already-applied migrations. Renumber them after the latest applied migration.`
);
}
if (!st.pending.length) {
log('Schema is up to date — no migrations to apply.');
return [];
}
const done = [];
for (const m of st.pending) {
log(`applying ${m.filename}`);
try {
await client.query('BEGIN');
await client.query(m.sql);
await client.query(
`INSERT INTO pf.schema_version (filename, checksum, applied_by) VALUES ($1, $2, $3)`,
[m.filename, m.checksum, user || null]
);
await client.query('COMMIT');
done.push(m.filename);
} catch (err) {
await client.query('ROLLBACK').catch(() => {});
throw new Error(`${m.filename} failed: ${err.message}`);
}
}
log(`Applied ${done.length} migration(s).`);
return done;
}
// Record migrations as applied without running them — for a database whose schema
// already matches, from before migrations existed.
async function baseline(client, { user, upTo, log = console.log } = {}) {
const st = await status(client);
const target = upTo
? st.pending.filter(m => m.filename <= upTo)
: st.pending;
if (!target.length) {
log('Nothing to baseline — no pending migrations.');
return [];
}
await client.query('BEGIN');
try {
for (const m of target) {
await client.query(
`INSERT INTO pf.schema_version (filename, checksum, applied_by)
VALUES ($1, $2, $3)
ON CONFLICT (filename) DO NOTHING`,
[m.filename, m.checksum, user ? `${user} (baseline)` : 'baseline']
);
}
await client.query('COMMIT');
} catch (err) {
await client.query('ROLLBACK').catch(() => {});
throw err;
}
log(`Marked ${target.length} migration(s) as applied without running them.`);
return target.map(m => m.filename);
}
module.exports = { MIGRATIONS_DIR, readMigrations, status, migrate, baseline, checksum };

View File

@ -1,10 +1,6 @@
// 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}},
@ -14,37 +10,6 @@
// 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.
const key = (pfx = '') => `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')
@ -80,26 +45,8 @@ function generateSQL(source, colMeta) {
);
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(),
@ -112,43 +59,6 @@ function generateSQL(source, colMeta) {
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() {
return `
SELECT
${grainSelect()}
FROM {{fc_table}}
GROUP BY
${grain.groupCols().join('\n ,')}`.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)
return dataCols.map(c => {
@ -241,7 +151,7 @@ ilog AS (
FROM base
RETURNING *
)
${opTail('ins')}`.trim();
SELECT * FROM ins`.trim();
}
function buildRecode() {
@ -272,12 +182,7 @@ ilog AS (
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();
SELECT * FROM neg UNION ALL SELECT * FROM ins`.trim();
}
function buildClone() {
@ -300,7 +205,7 @@ ilog AS (
{{exclude_clause}}
RETURNING *
)
${opTail('ins')}`.trim();
SELECT * FROM ins`.trim();
}
function buildUndo() {
@ -404,4 +309,4 @@ function esc(val) {
return String(val).replace(/'/g, "''");
}
module.exports = { generateSQL, grainOf, applyTokens, buildWhere, buildExcludeClause, buildSetClause, buildFilterClause, esc };
module.exports = { generateSQL, applyTokens, buildWhere, buildExcludeClause, buildSetClause, buildFilterClause, esc };

View File

@ -7,11 +7,7 @@
"scripts": {
"start": "node server.js",
"dev": "nodemon server.js",
"build": "cd ui && npm run build",
"migrate": "node scripts/migrate.js up",
"migrate:status": "node scripts/migrate.js status",
"migrate:baseline": "node scripts/migrate.js baseline",
"schema:dump": "node scripts/schema-dump.js"
"build": "cd ui && npm run build"
},
"dependencies": {
"apache-arrow": "^21.1.0",

View File

@ -189,22 +189,9 @@ CREATE TABLE pf.sql (
```
setup_sql/
migrations/ -- forward-only .sql, applied once each, tracked in pf.schema_version
0001_initial_schema.sql CREATE SCHEMA pf + metadata tables
0002_source_default_layout.sql
0003_col_meta_dim_group_period.sql
0004_col_meta_in_grain.sql
gen_dim_period.sql -- parameterised calendar load (fiscal year start month); not a migration
schema.generated.sql -- pg_dump reference snapshot; generated, never edited or applied
README.md -- migration workflow
01_schema.sql -- CREATE SCHEMA pf; create all metadata tables (source, col_meta, version, log, sql)
```
Install with `npm run migrate`, then run `gen_dim_period.sql`. For a database that
already matches the schema, `npm run migrate:baseline` records the migrations as
applied without running them. `server.js` refuses to start when the database is
behind, so drift surfaces at boot rather than as a column-not-found error inside
an unrelated request.
Source registration, col_meta configuration, SQL generation, version creation, and forecast table DDL all happen via API.
---
@ -727,26 +714,13 @@ DELETE FROM pf.log WHERE id = {{logid}};
---
## Display-grain pre-aggregation
## Display-grain pre-aggregation (planned)
**Status:** built (static grain). This is **Path B** (pre-aggregated extract →
native Perspective table) of two candidate designs; rationale, the Path A
alternative (live virtual-server aggregation), the spike evidence, and the
A-vs-B trade-off live in `pf_perspective_options.md` (§Two candidate designs,
§Spike findings).
The grain is **static** — set once per source in Setup and baked into the stored
`pf.sql` templates, so load and operations agree by construction. `in_grain`
means *eligible for the grain*, and in this version the grain is exactly the set
of eligible columns. Deriving a narrower grain per pivot at request time (the
dynamic variant) is then additive: intersect the viewer's field set with the
eligible set. Leaving high-cardinality columns (`part`, raw day dates,
currency-level detail) unflagged is what keeps the grain from exploding back
toward raw, regardless of what a user drags into the pivot.
**Measured on `pf.fc_osm_stack_20`** at `pending_rep × customer × smon`:
534,902 → **6,154** rows (≈87×), `pf_gkey` unique across all 6,154, and both
measures reconcile exactly to the raw totals (283,296,087.67 / 962,142,261.46).
**Status:** designed, not yet built. This is the concrete design for **Path B**
(pre-aggregated extract → native Perspective table) of two candidate designs;
rationale, the Path A alternative (live virtual-server aggregation), the spike
evidence, and the A-vs-B trade-off live in `pf_perspective_options.md`
(§Two candidate designs, §Spike findings).
**Problem it solves.** The current transport ships every raw forecast row to the
browser (≈535k rows / ~250 MB / ~2 min on `osm_stack`). Perspective then pivots
@ -777,11 +751,7 @@ the stored `pf.sql` templates, so initial load and operations agree on it.
### Initial load — `GET /api/versions/:id/agg`
Replaces the raw `/data` stream for grain-based versions. Aggregates the forecast
table to the stored grain and returns Arrow IPC. The template is stored in
`pf.sql` as operation `get_agg`, generated only when a grain is defined; clearing
the grain and regenerating removes it, and the client falls back to `/data`.
Both endpoints speak the same protocol (one record batch plus an `X-Row-Count`
header), so the client only chooses the URL:
table to the stored grain and returns Arrow IPC:
```sql
SELECT
@ -813,12 +783,6 @@ the smallest change from today's code, which already appends operation results v
accumulate (rather than replacing a bucket) and a delete can remove exactly that
operation's rows.
As built, the concatenation is
`concat_ws(chr(31), COALESCE(col::text, chr(30)), …, pf_iter, pf_logid::text)`.
The separator and NULL sentinel matter: plain `concat_ws` skips NULLs, so
`('a', NULL)` and `(NULL, 'a')` would produce the same key and silently merge two
groups into one indexed row. `chr(30)` also keeps NULL distinct from `''`.
### Write path (scale / recode / clone) — append the new log entry's rows
Operations INSERT raw rows into `{{fc_table}}` under a new `pf_logid` as today; the
@ -849,20 +813,9 @@ because each row carries the new, unique `pf_logid`.)
A logid's rows are uniquely keyed, so undo just removes them and lets the view
re-sum — no re-aggregation, no emptied-bucket handling, no snapshot caveat:
`RETURNING` does not accept `DISTINCT`, so the delete feeds a CTE that reduces its
output to the distinct grain keys. `rows_deleted` still counts raw rows removed:
```sql
WITH
del AS (
DELETE FROM {{fc_table}}
WHERE pf_logid = {{logid}}
RETURNING {{grain_cols}}, pf_iter, pf_logid
)
SELECT
count(*)::int AS rows_deleted
,array_agg(DISTINCT {{grain_key}}) AS pf_gkeys
FROM del;
DELETE FROM {{fc_table}} WHERE pf_logid = {{logid}}
RETURNING DISTINCT {{grain_cols}}, pf_iter, pf_logid; -- → pf_gkeys to remove
DELETE FROM pf.log WHERE id = {{logid}};
```

View File

@ -1,5 +1,4 @@
const express = require('express');
const { grainOf } = require('../lib/sql_generator');
const { fcTable } = require('../lib/utils');
module.exports = function(pool) {
@ -52,7 +51,7 @@ module.exports = function(pool) {
const logId = parseInt(req.params.logid);
try {
const logResult = await pool.query(`
SELECT l.*, v.status, s.tname, v.id AS version_id, v.source_id
SELECT l.*, v.status, s.tname, v.id AS version_id
FROM pf.log l
JOIN pf.version v ON v.id = l.version_id
JOIN pf.source s ON s.id = v.source_id
@ -62,46 +61,18 @@ module.exports = function(pool) {
const log = logResult.rows[0];
if (log.status === 'closed') return res.status(403).json({ error: 'Version is closed' });
const table = fcTable(log.tname, log.version_id);
// In grain mode the client's table is indexed on pf_gkey, so undo has to
// report the grain keys to remove rather than raw pf_ids. The keys are
// distinct while rows_deleted still counts the raw rows removed.
const colMeta = await pool.query(
`SELECT cname, role, in_grain, opos FROM pf.col_meta WHERE source_id = $1 ORDER BY opos`,
[log.source_id]
);
const grain = grainOf(colMeta.rows);
const client = await pool.connect();
try {
await client.query('BEGIN');
const deleted = grain
? await client.query(`
WITH
del AS (
DELETE FROM ${table}
WHERE pf_logid = $1
RETURNING ${grain.groupCols().join(', ')}
)
SELECT
count(*)::int AS rows_deleted
,array_agg(DISTINCT ${grain.key()}) AS pf_gkeys
FROM del
`, [logId])
: await client.query(
`DELETE FROM ${table} WHERE pf_logid = $1 RETURNING pf_id`, [logId]
);
const deleted = await client.query(
`DELETE FROM ${table} WHERE pf_logid = $1 RETURNING pf_id`, [logId]
);
await client.query('DELETE FROM pf.log WHERE id = $1', [logId]);
await client.query('COMMIT');
res.json(grain
? {
rows_deleted: deleted.rows[0].rows_deleted,
pf_gkeys: deleted.rows[0].pf_gkeys || []
}
: {
rows_deleted: deleted.rowCount,
pf_ids: deleted.rows.map(r => r.pf_id)
});
res.json({
rows_deleted: deleted.rowCount,
pf_ids: deleted.rows.map(r => r.pf_id)
});
} catch (err) {
await client.query('ROLLBACK');
throw err;

View File

@ -127,37 +127,6 @@ module.exports = function(pool) {
}
});
// Aggregate a version to its display grain and return it as Arrow IPC.
// This replaces /data for sources that define a grain (col_meta.in_grain):
// the aggregation collapses the row count by orders of magnitude, so the
// result loads as one small native Perspective table indexed on pf_gkey and
// the WASM view still does all rollup/expand/collapse locally.
router.get('/versions/:id/agg', async (req, res) => {
try {
const ctx = await getContext(parseInt(req.params.id), 'get_agg');
const sql = applyTokens(ctx.sql, { fc_table: ctx.table });
const { rows } = await runSQL(sql);
res.setHeader('Content-Type', 'application/vnd.apache.arrow.stream');
res.setHeader('X-Row-Count', String(rows.length));
if (rows.length === 0) { res.end(); return; }
// column arrays, one Arrow record batch — same constraint as /data:
// per-batch dictionaries crash Perspective's Arrow reader
const colArrays = Object.fromEntries(Object.keys(rows[0]).map(k => [k, []]));
for (const row of rows) {
for (const k of Object.keys(colArrays)) colArrays[k].push(row[k]);
}
const buf = tableToIPC(tableFromArrays(colArrays), 'stream');
res.setHeader('Content-Length', String(buf.byteLength));
res.end(Buffer.from(buf.buffer, buf.byteOffset, buf.byteLength));
} catch (err) {
console.error(err);
if (!res.headersSent) res.status(err.status || 500).json({ error: err.message });
else res.destroy();
}
});
// load baseline rows from source table — additive, no delete
router.post('/versions/:id/baseline', async (req, res) => {
const { where_clause, date_offset, pf_user, note, filters, raw_where } = req.body;

View File

@ -89,15 +89,14 @@ module.exports = function(pool) {
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)
INSERT INTO pf.col_meta (source_id, cname, label, role, is_key, dim_group, dim_period_col, opos)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
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,
@ -107,7 +106,6 @@ module.exports = function(pool) {
col.is_key || false,
col.dim_group || null,
col.dim_period_col || null,
col.in_grain || false,
col.opos || null
]);
}
@ -168,13 +166,6 @@ module.exports = function(pool) {
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');

View File

@ -1,72 +0,0 @@
#!/usr/bin/env node
// Migration CLI.
//
// npm run migrate apply pending migrations
// npm run migrate:status show what is applied and what is pending
// npm run migrate:baseline record pending migrations as applied WITHOUT running
// them — for a database that already matches
//
// Pass --up-to=<filename> to baseline only through a given migration.
require('dotenv').config();
const os = require('os');
const { Client } = require('pg');
const { status, migrate, baseline } = require('../lib/migrations');
const args = process.argv.slice(2);
const command = args.find(a => !a.startsWith('--')) || 'up';
const upTo = (args.find(a => a.startsWith('--up-to=')) || '').split('=')[1] || null;
function connect() {
return new Client({
host: process.env.DB_HOST,
port: parseInt(process.env.DB_PORT) || 5432,
database: process.env.DB_NAME,
user: process.env.DB_USER,
password: process.env.DB_PASSWORD,
ssl: false
});
}
async function showStatus(client) {
const st = await status(client);
console.log(`\n${process.env.DB_NAME} (${st.applied.length} applied, ${st.pending.length} pending)\n`);
for (const m of st.onDisk) {
const hit = st.applied.find(a => a.filename === m.filename);
const mark = !hit ? 'PENDING'
: st.modified.includes(m.filename) ? 'MODIFIED'
: 'applied';
const when = hit ? hit.applied_at.toISOString().slice(0, 19).replace('T', ' ') : '';
console.log(` ${mark.padEnd(9)} ${m.filename.padEnd(40)} ${when}`);
}
for (const f of st.missing) console.log(` MISSING ${f.padEnd(40)} recorded as applied but not on disk`);
if (st.modified.length) console.log(`\n ! modified after being applied: ${st.modified.join(', ')}`);
if (st.outOfOrder.length) console.log(` ! would apply out of order: ${st.outOfOrder.join(', ')}`);
console.log('');
return st;
}
(async () => {
const client = connect();
await client.connect();
try {
const user = `${os.userInfo().username}@${os.hostname()}`;
if (command === 'status') {
const st = await showStatus(client);
process.exitCode = (st.modified.length || st.missing.length) ? 1 : 0;
} else if (command === 'baseline') {
await baseline(client, { user, upTo });
await showStatus(client);
} else if (command === 'up') {
await migrate(client, { user });
} else {
console.error(`Unknown command "${command}" — expected up, status, or baseline.`);
process.exitCode = 2;
}
} catch (err) {
console.error(`\nmigration error: ${err.message}\n`);
process.exitCode = 1;
} finally {
await client.end();
}
})();

View File

@ -1,50 +0,0 @@
#!/usr/bin/env node
// Regenerate setup_sql/schema.generated.sql — a readable snapshot of the current
// pf schema, for reference only. Migrations are the source of truth; this file is
// generated so it cannot drift from the database the way a hand-maintained
// schema file does. Never edit it, and never apply it to create a database.
require('dotenv').config();
const { spawnSync } = require('child_process');
const fs = require('fs');
const path = require('path');
const out = path.join(__dirname, '..', 'setup_sql', 'schema.generated.sql');
const result = spawnSync('pg_dump', [
'--schema-only', '--no-owner', '--no-privileges',
'--schema=pf',
// per-version forecast tables are created at runtime — including them would
// make this file churn every time a version is created or dropped
'--exclude-table=pf.fc_*',
'-f', out
], {
stdio: 'inherit',
env: {
...process.env,
PGHOST: process.env.DB_HOST,
PGPORT: process.env.DB_PORT || '5432',
PGDATABASE: process.env.DB_NAME,
PGUSER: process.env.DB_USER,
PGPASSWORD: process.env.DB_PASSWORD
}
});
if (result.error) {
console.error(`pg_dump failed: ${result.error.message}`);
process.exit(1);
}
if (result.status !== 0) process.exit(result.status);
// Strip the lines pg_dump varies between runs — a random \restrict token and the
// server/client version banner — so regenerating an unchanged schema produces an
// identical file instead of a spurious diff.
const cleaned = fs.readFileSync(out, 'utf8')
.split('\n')
.filter(l => !/^\\(un)?restrict /.test(l))
.filter(l => !/^-- Dumped (from|by) /.test(l))
.join('\n')
.replace(/\n{3,}/g, '\n\n');
fs.writeFileSync(out, cleaned);
console.log(`wrote ${path.relative(process.cwd(), out)} (${cleaned.split('\n').length} lines)`);

View File

@ -33,38 +33,5 @@ app.use('/api', require('./routes/operations')(pool));
app.use('/api', require('./routes/log')(pool));
// Refuse to start against a database whose schema is behind the migrations on
// disk. Without this the mismatch surfaces later as a confusing query error
// ('column "x" does not exist') deep inside an unrelated request.
// Set PF_SKIP_MIGRATION_CHECK=1 to bypass.
async function checkSchema() {
if (process.env.PF_SKIP_MIGRATION_CHECK === '1') return;
const { status } = require('./lib/migrations');
const client = await pool.connect();
try {
const st = await status(client);
const problems = [
st.pending.length && `${st.pending.length} pending: ${st.pending.map(m => m.filename).join(', ')}`,
st.modified.length && `modified after being applied: ${st.modified.join(', ')}`,
st.missing.length && `applied but missing from disk: ${st.missing.join(', ')}`
].filter(Boolean);
if (problems.length) {
console.error('\n Database schema is out of date:');
for (const p of problems) console.error(` - ${p}`);
console.error('\n Run "npm run migrate" (or "npm run migrate:status" for detail).');
console.error(' For a database that already matches, "npm run migrate:baseline".\n');
throw new Error('schema out of date');
}
console.log(`schema up to date (${st.applied.length} migrations applied)`);
} finally {
client.release();
}
}
const port = process.env.PORT || 3010;
checkSchema()
.then(() => app.listen(port, '0.0.0.0', () => console.log(`pf_app started on port ${port}`)))
.catch((err) => {
if (err.message !== 'schema out of date') console.error('startup failed:', err.message);
process.exit(1);
});
app.listen(port, '0.0.0.0', () => console.log(`pf_app started on port ${port}`));

69
setup_sql/01_schema.sql Normal file
View File

@ -0,0 +1,69 @@
-- Pivot Forecast schema install
-- Run once against target database: psql -d <db> -f setup_sql/01_schema.sql
CREATE SCHEMA IF NOT EXISTS pf;
CREATE TABLE IF NOT EXISTS pf.source (
id serial PRIMARY KEY,
schema text NOT NULL,
tname text NOT NULL,
label text,
status text NOT NULL DEFAULT 'active', -- active | archived
default_layout jsonb, -- Perspective view config used as the per-source default
created_at timestamptz NOT NULL DEFAULT now(),
created_by text,
UNIQUE (schema, tname)
);
-- backfill columns for existing installs
ALTER TABLE pf.source ADD COLUMN IF NOT EXISTS default_layout jsonb;
ALTER TABLE pf.col_meta ADD COLUMN IF NOT EXISTS dim_group text;
ALTER TABLE pf.col_meta ADD COLUMN IF NOT EXISTS dim_period_col text;
-- pf.dim_period: run setup_sql/gen_dim_period.sql to create and populate
CREATE TABLE IF NOT EXISTS pf.col_meta (
id serial PRIMARY KEY,
source_id integer NOT NULL REFERENCES pf.source(id) ON DELETE CASCADE,
cname text NOT NULL,
label text,
role text NOT NULL DEFAULT 'ignore', -- dimension | value | units | date | ignore
is_key boolean NOT NULL DEFAULT false, -- true = usable in WHERE slice
opos integer,
UNIQUE (source_id, cname)
);
CREATE TABLE IF NOT EXISTS pf.version (
id serial PRIMARY KEY,
source_id integer NOT NULL REFERENCES pf.source(id) ON DELETE RESTRICT,
name text NOT NULL,
description text,
status text NOT NULL DEFAULT 'open', -- open | closed
exclude_iters jsonb NOT NULL DEFAULT '["reference"]'::jsonb,
created_at timestamptz NOT NULL DEFAULT now(),
created_by text,
closed_at timestamptz,
closed_by text,
UNIQUE (source_id, name)
);
CREATE TABLE IF NOT EXISTS pf.log (
id bigserial PRIMARY KEY,
version_id integer NOT NULL REFERENCES pf.version(id) ON DELETE CASCADE,
pf_user text NOT NULL,
stamp timestamptz NOT NULL DEFAULT now(),
operation text NOT NULL, -- baseline | reference | scale | recode | clone
slice jsonb,
params jsonb,
note text
);
-- generated operation SQL per source, stored after col_meta is configured
CREATE TABLE IF NOT EXISTS pf.sql (
id serial PRIMARY KEY,
source_id integer NOT NULL REFERENCES pf.source(id) ON DELETE CASCADE,
operation text NOT NULL, -- get_data | baseline | reference | scale | recode | clone | undo
sql text NOT NULL,
generated_at timestamptz NOT NULL DEFAULT now(),
UNIQUE (source_id, operation)
);

View File

@ -1,64 +0,0 @@
# Database setup
## Migrations
`migrations/*.sql` are applied in filename order and recorded in
`pf.schema_version`. They are the **source of truth** for the `pf` schema — there
is no hand-maintained current-state file to drift out of sync.
```bash
npm run migrate # apply pending migrations
npm run migrate:status # what is applied, what is pending
npm run migrate:baseline # record pending as applied WITHOUT running them
```
`server.js` refuses to start when the database is behind, so drift surfaces at
boot rather than as `column "x" does not exist` inside an unrelated request. Set
`PF_SKIP_MIGRATION_CHECK=1` to bypass.
### Fresh database
```bash
npm run migrate
psql -d <db> -f setup_sql/gen_dim_period.sql
```
### Existing database that already matches
Use `baseline` so the runner does not try to re-create tables that exist:
```bash
npm run migrate:baseline
```
To baseline only part of the way — the schema matches through `0003` but not
`0004` — pass `--up-to` and then migrate the rest:
```bash
node scripts/migrate.js baseline --up-to=0003_col_meta_dim_group_period.sql
npm run migrate
```
## Writing a migration
- Name it `NNNN_short_description.sql`, numbered after the highest existing file.
- One concern per file. Keep it forward-only; there are no down migrations.
- **Applied migrations are immutable.** The runner stores a checksum and refuses
to proceed if a file changes after being applied, because editing one means
databases silently disagree about what the schema is. To fix a mistake, add a
new migration.
- No `IF NOT EXISTS` guards on new migrations. The bookkeeping already guarantees
each runs once, and the guards hide ordering mistakes — the reason the old
`01_schema.sql` had `ALTER`s sitting above the `CREATE TABLE` they depended on,
broken for anyone installing from scratch. `0004` is the one exception, since it
was applied by hand before migrations existed.
## Not migrations
- **`gen_dim_period.sql`** — creates and populates `pf.dim_period`. It is a
parameterised data load (configurable fiscal year start month), not a schema
change, so it stays a script you run deliberately.
- **`pf.fc_{tname}_{version_id}`** — per-version forecast tables, created and
dropped at runtime by `routes/versions.js` from `col_meta`. Never migrated.
- **`schema.generated.sql`** — a `pg_dump` snapshot for reading, refreshed with
`npm run schema:dump`. Generated, never edited, never applied.

View File

@ -1,63 +0,0 @@
-- Initial pf schema: sources, column metadata, versions, audit log, generated SQL.
--
-- This is the schema as it stood before the additive columns in later migrations.
-- Columns added afterwards are declared once, in their own migration — not here.
CREATE SCHEMA IF NOT EXISTS pf;
CREATE TABLE pf.source (
id serial PRIMARY KEY,
schema text NOT NULL,
tname text NOT NULL,
label text,
status text NOT NULL DEFAULT 'active', -- active | archived
created_at timestamptz NOT NULL DEFAULT now(),
created_by text,
UNIQUE (schema, tname)
);
CREATE TABLE pf.col_meta (
id serial PRIMARY KEY,
source_id integer NOT NULL REFERENCES pf.source(id) ON DELETE CASCADE,
cname text NOT NULL,
label text,
role text NOT NULL DEFAULT 'ignore', -- dimension | value | units | date | filter | ignore
is_key boolean NOT NULL DEFAULT false, -- true = usable in WHERE slice
opos integer,
UNIQUE (source_id, cname)
);
CREATE TABLE pf.version (
id serial PRIMARY KEY,
source_id integer NOT NULL REFERENCES pf.source(id) ON DELETE RESTRICT,
name text NOT NULL,
description text,
status text NOT NULL DEFAULT 'open', -- open | closed
exclude_iters jsonb NOT NULL DEFAULT '["reference"]'::jsonb,
created_at timestamptz NOT NULL DEFAULT now(),
created_by text,
closed_at timestamptz,
closed_by text,
UNIQUE (source_id, name)
);
CREATE TABLE pf.log (
id bigserial PRIMARY KEY,
version_id integer NOT NULL REFERENCES pf.version(id) ON DELETE CASCADE,
pf_user text NOT NULL,
stamp timestamptz NOT NULL DEFAULT now(),
operation text NOT NULL, -- baseline | reference | scale | recode | clone
slice jsonb,
params jsonb,
note text
);
-- generated operation SQL per source, stored after col_meta is configured
CREATE TABLE pf.sql (
id serial PRIMARY KEY,
source_id integer NOT NULL REFERENCES pf.source(id) ON DELETE CASCADE,
operation text NOT NULL, -- get_data | get_agg | baseline | reference | scale | recode | clone | undo
sql text NOT NULL,
generated_at timestamptz NOT NULL DEFAULT now(),
UNIQUE (source_id, operation)
);

View File

@ -1,4 +0,0 @@
-- Per-source default Perspective view config, applied when a version has no
-- saved layout of its own.
ALTER TABLE pf.source ADD COLUMN default_layout jsonb;

View File

@ -1,7 +0,0 @@
-- dim_group groups functionally dependent columns (e.g. a date and the year/month
-- dimensions derived from it). dim_period_col maps such a dimension to a
-- pf.dim_period column, so date-adjacent values are derived by JOIN at load time
-- rather than copied raw from the source.
ALTER TABLE pf.col_meta ADD COLUMN dim_group text;
ALTER TABLE pf.col_meta ADD COLUMN dim_period_col text;

View File

@ -1,8 +0,0 @@
-- in_grain flags dimension/date columns that define the display grain: the
-- forecast load is pre-aggregated to the flagged columns instead of shipping raw
-- rows. See pf_spec.md -> Display-grain pre-aggregation.
--
-- Already applied by hand on the original dev database before migrations existed;
-- IF NOT EXISTS keeps replay safe there. New migrations should not need the guard.
ALTER TABLE pf.col_meta ADD COLUMN IF NOT EXISTS in_grain boolean NOT NULL DEFAULT false;

View File

@ -1,384 +0,0 @@
--
-- PostgreSQL database dump
--
SET statement_timeout = 0;
SET lock_timeout = 0;
SET idle_in_transaction_session_timeout = 0;
SET transaction_timeout = 0;
SET client_encoding = 'UTF8';
SET standard_conforming_strings = on;
SELECT pg_catalog.set_config('search_path', '', false);
SET check_function_bodies = false;
SET xmloption = content;
SET client_min_messages = warning;
SET row_security = off;
--
-- Name: pf; Type: SCHEMA; Schema: -; Owner: -
--
CREATE SCHEMA pf;
SET default_tablespace = '';
SET default_table_access_method = heap;
--
-- Name: col_meta; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.col_meta (
id integer NOT NULL,
source_id integer NOT NULL,
cname text NOT NULL,
label text,
role text DEFAULT 'ignore'::text NOT NULL,
is_key boolean DEFAULT false NOT NULL,
opos integer,
dim_group text,
dim_period_col text,
in_grain boolean DEFAULT false NOT NULL
);
--
-- Name: col_meta_id_seq; Type: SEQUENCE; Schema: pf; Owner: -
--
CREATE SEQUENCE pf.col_meta_id_seq
AS integer
START WITH 1
INCREMENT BY 1
NO MINVALUE
NO MAXVALUE
CACHE 1;
--
-- Name: col_meta_id_seq; Type: SEQUENCE OWNED BY; Schema: pf; Owner: -
--
ALTER SEQUENCE pf.col_meta_id_seq OWNED BY pf.col_meta.id;
--
-- Name: dim_period; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.dim_period (
sdat date NOT NULL,
edat date NOT NULL,
drange daterange NOT NULL,
ndays integer NOT NULL,
cal_year integer NOT NULL,
cal_quarter integer NOT NULL,
cal_month integer NOT NULL,
cal_month_abbr text NOT NULL,
cal_month_name text NOT NULL,
cal_label text NOT NULL,
fisc_year integer NOT NULL,
fisc_quarter integer NOT NULL,
fisc_quarter_label text NOT NULL,
fisc_month integer NOT NULL,
fisc_month_abbr text NOT NULL,
fisc_month_name text NOT NULL,
fisc_label text NOT NULL,
period_key text NOT NULL
);
--
-- Name: log; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.log (
id bigint NOT NULL,
version_id integer NOT NULL,
pf_user text NOT NULL,
stamp timestamp with time zone DEFAULT now() NOT NULL,
operation text NOT NULL,
slice jsonb,
params jsonb,
note text
);
--
-- Name: log_id_seq; Type: SEQUENCE; Schema: pf; Owner: -
--
CREATE SEQUENCE pf.log_id_seq
START WITH 1
INCREMENT BY 1
NO MINVALUE
NO MAXVALUE
CACHE 1;
--
-- Name: log_id_seq; Type: SEQUENCE OWNED BY; Schema: pf; Owner: -
--
ALTER SEQUENCE pf.log_id_seq OWNED BY pf.log.id;
--
-- Name: schema_version; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.schema_version (
filename text NOT NULL,
checksum text NOT NULL,
applied_at timestamp with time zone DEFAULT now() NOT NULL,
applied_by text
);
--
-- Name: source; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.source (
id integer NOT NULL,
schema text NOT NULL,
tname text NOT NULL,
label text,
status text DEFAULT 'active'::text NOT NULL,
created_at timestamp with time zone DEFAULT now() NOT NULL,
created_by text,
default_layout jsonb
);
--
-- Name: source_id_seq; Type: SEQUENCE; Schema: pf; Owner: -
--
CREATE SEQUENCE pf.source_id_seq
AS integer
START WITH 1
INCREMENT BY 1
NO MINVALUE
NO MAXVALUE
CACHE 1;
--
-- Name: source_id_seq; Type: SEQUENCE OWNED BY; Schema: pf; Owner: -
--
ALTER SEQUENCE pf.source_id_seq OWNED BY pf.source.id;
--
-- Name: sql; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.sql (
id integer NOT NULL,
source_id integer NOT NULL,
operation text NOT NULL,
sql text NOT NULL,
generated_at timestamp with time zone DEFAULT now() NOT NULL
);
--
-- Name: sql_id_seq; Type: SEQUENCE; Schema: pf; Owner: -
--
CREATE SEQUENCE pf.sql_id_seq
AS integer
START WITH 1
INCREMENT BY 1
NO MINVALUE
NO MAXVALUE
CACHE 1;
--
-- Name: sql_id_seq; Type: SEQUENCE OWNED BY; Schema: pf; Owner: -
--
ALTER SEQUENCE pf.sql_id_seq OWNED BY pf.sql.id;
--
-- Name: version; Type: TABLE; Schema: pf; Owner: -
--
CREATE TABLE pf.version (
id integer NOT NULL,
source_id integer NOT NULL,
name text NOT NULL,
description text,
status text DEFAULT 'open'::text NOT NULL,
exclude_iters jsonb DEFAULT '["reference"]'::jsonb NOT NULL,
created_at timestamp with time zone DEFAULT now() NOT NULL,
created_by text,
closed_at timestamp with time zone,
closed_by text
);
--
-- Name: version_id_seq; Type: SEQUENCE; Schema: pf; Owner: -
--
CREATE SEQUENCE pf.version_id_seq
AS integer
START WITH 1
INCREMENT BY 1
NO MINVALUE
NO MAXVALUE
CACHE 1;
--
-- Name: version_id_seq; Type: SEQUENCE OWNED BY; Schema: pf; Owner: -
--
ALTER SEQUENCE pf.version_id_seq OWNED BY pf.version.id;
--
-- Name: col_meta id; Type: DEFAULT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.col_meta ALTER COLUMN id SET DEFAULT nextval('pf.col_meta_id_seq'::regclass);
--
-- Name: log id; Type: DEFAULT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.log ALTER COLUMN id SET DEFAULT nextval('pf.log_id_seq'::regclass);
--
-- Name: source id; Type: DEFAULT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.source ALTER COLUMN id SET DEFAULT nextval('pf.source_id_seq'::regclass);
--
-- Name: sql id; Type: DEFAULT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.sql ALTER COLUMN id SET DEFAULT nextval('pf.sql_id_seq'::regclass);
--
-- Name: version id; Type: DEFAULT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.version ALTER COLUMN id SET DEFAULT nextval('pf.version_id_seq'::regclass);
--
-- Name: col_meta col_meta_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.col_meta
ADD CONSTRAINT col_meta_pkey PRIMARY KEY (id);
--
-- Name: col_meta col_meta_source_id_cname_key; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.col_meta
ADD CONSTRAINT col_meta_source_id_cname_key UNIQUE (source_id, cname);
--
-- Name: dim_period dim_period_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.dim_period
ADD CONSTRAINT dim_period_pkey PRIMARY KEY (sdat);
--
-- Name: log log_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.log
ADD CONSTRAINT log_pkey PRIMARY KEY (id);
--
-- Name: schema_version schema_version_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.schema_version
ADD CONSTRAINT schema_version_pkey PRIMARY KEY (filename);
--
-- Name: source source_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.source
ADD CONSTRAINT source_pkey PRIMARY KEY (id);
--
-- Name: source source_schema_tname_key; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.source
ADD CONSTRAINT source_schema_tname_key UNIQUE (schema, tname);
--
-- Name: sql sql_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.sql
ADD CONSTRAINT sql_pkey PRIMARY KEY (id);
--
-- Name: sql sql_source_id_operation_key; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.sql
ADD CONSTRAINT sql_source_id_operation_key UNIQUE (source_id, operation);
--
-- Name: version version_pkey; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.version
ADD CONSTRAINT version_pkey PRIMARY KEY (id);
--
-- Name: version version_source_id_name_key; Type: CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.version
ADD CONSTRAINT version_source_id_name_key UNIQUE (source_id, name);
--
-- Name: dim_period_cal_idx; Type: INDEX; Schema: pf; Owner: -
--
CREATE INDEX dim_period_cal_idx ON pf.dim_period USING btree (cal_year, cal_month);
--
-- Name: dim_period_drange_idx; Type: INDEX; Schema: pf; Owner: -
--
CREATE INDEX dim_period_drange_idx ON pf.dim_period USING gist (drange);
--
-- Name: dim_period_fisc_idx; Type: INDEX; Schema: pf; Owner: -
--
CREATE INDEX dim_period_fisc_idx ON pf.dim_period USING btree (fisc_year, fisc_month);
--
-- Name: col_meta col_meta_source_id_fkey; Type: FK CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.col_meta
ADD CONSTRAINT col_meta_source_id_fkey FOREIGN KEY (source_id) REFERENCES pf.source(id) ON DELETE CASCADE;
--
-- Name: log log_version_id_fkey; Type: FK CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.log
ADD CONSTRAINT log_version_id_fkey FOREIGN KEY (version_id) REFERENCES pf.version(id) ON DELETE CASCADE;
--
-- Name: sql sql_source_id_fkey; Type: FK CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.sql
ADD CONSTRAINT sql_source_id_fkey FOREIGN KEY (source_id) REFERENCES pf.source(id) ON DELETE CASCADE;
--
-- Name: version version_source_id_fkey; Type: FK CONSTRAINT; Schema: pf; Owner: -
--
ALTER TABLE ONLY pf.version
ADD CONSTRAINT version_source_id_fkey FOREIGN KEY (source_id) REFERENCES pf.source(id) ON DELETE RESTRICT;
--
-- PostgreSQL database dump complete
--

View File

@ -148,37 +148,6 @@ export default function Forecast({ sources = [], sourceId, versionId, refreshSou
}
}
// Stream an Arrow IPC endpoint into one buffer, reporting download progress.
// Both /data and /agg speak the same protocol a single record batch plus an
// X-Row-Count header so the caller only picks the URL.
async function fetchArrow(url) {
const r = await fetch(url)
if (!r.ok) { const { error } = await r.json(); throw new Error(error || 'Failed to load data') }
const rowCount = parseInt(r.headers.get('X-Row-Count') || '0')
const total = parseInt(r.headers.get('Content-Length') || '0') || null
const reader = r.body.getReader()
const chunks = []
let received = 0
let lastUpdate = 0
setLoadProgress({ received: 0, total })
while (true) {
const { done, value } = await reader.read()
if (done) break
chunks.push(value)
received += value.byteLength
const now = Date.now()
if (now - lastUpdate >= 100) {
setLoadProgress({ received, total })
lastUpdate = now
}
}
setLoadProgress({ received, total })
const merged = new Uint8Array(received)
let pos = 0
for (const c of chunks) { merged.set(c, pos); pos += c.byteLength }
return { buffer: merged.buffer, rowCount }
}
function loadLayouts(vid) {
const stored = localStorage.getItem(LAYOUTS_KEY(vid))
setLayouts(stored ? JSON.parse(stored) : [])
@ -195,34 +164,43 @@ export default function Forecast({ sources = [], sourceId, versionId, refreshSou
setSlice({})
expandDepthRef.current = null
try {
// col_meta first it decides which endpoint to load from, and it is a tiny
// query next to the data fetch it gates.
const meta = await fetch(`/api/sources/${sid}/cols`).then(r => r.json())
colMetaRef.current = meta
// Grain mode: the source declares a display grain, so the server ships rows
// already aggregated to it and the table is indexed on pf_gkey. Without a
// grain we load raw forecast rows indexed on pf_id, as before.
const grainMeta = meta.filter(c => c.in_grain && ['dimension','date'].includes(c.role))
const grainMode = grainMeta.length > 0
const indexCol = grainMode ? 'pf_gkey' : 'pf_id'
const [perspective, dataResult] = await Promise.all([
const [perspective, dataResult, meta] = await Promise.all([
loadPerspective(),
fetchArrow(`/api/versions/${vid}/${grainMode ? 'agg' : 'data'}`),
fetch(`/api/versions/${vid}/data`).then(async r => {
if (!r.ok) { const { error } = await r.json(); throw new Error(error || 'Failed to load data') }
const rowCount = parseInt(r.headers.get('X-Row-Count') || '0')
const total = parseInt(r.headers.get('Content-Length') || '0') || null
const reader = r.body.getReader()
const chunks = []
let received = 0
let lastUpdate = 0
setLoadProgress({ received: 0, total })
while (true) {
const { done, value } = await reader.read()
if (done) break
chunks.push(value)
received += value.byteLength
const now = Date.now()
if (now - lastUpdate >= 100) {
setLoadProgress({ received, total })
lastUpdate = now
}
}
setLoadProgress({ received, total })
const merged = new Uint8Array(received)
let pos = 0
for (const c of chunks) { merged.set(c, pos); pos += c.byteLength }
return { buffer: merged.buffer, rowCount }
}),
fetch(`/api/sources/${sid}/cols`).then(r => r.json()),
])
const { buffer, rowCount } = dataResult
const validCols = new Set(grainMode
? [
...grainMeta.map(c => c.cname),
...meta.filter(c => ['value','units'].includes(c.role)).map(c => c.cname),
'pf_gkey', 'pf_iter', 'pf_logid',
]
: [
...meta.filter(c => ['dimension','value','units','date'].includes(c.role)).map(c => c.cname),
'pf_id', 'pf_iter', 'pf_logid', 'pf_user', 'created_at',
])
colMetaRef.current = meta
const validCols = new Set([
...meta.filter(c => ['dimension','value','units','date'].includes(c.role)).map(c => c.cname),
'pf_id', 'pf_iter', 'pf_logid', 'pf_user', 'created_at',
])
const tableName = `fc_${vid}`
if (rowCount >= 500000) setLargeDataset(true)
@ -243,7 +221,7 @@ export default function Forecast({ sources = [], sourceId, versionId, refreshSou
if (stale) await stale.delete()
} catch {}
const opts = { name: tableName, index: indexCol }
const opts = { name: tableName, index: 'pf_id' }
tableRef.current = await (rowCount > 0 ? worker.table(buffer, opts) : worker.table([], opts))
if (myId !== initIdRef.current) {
@ -545,11 +523,8 @@ export default function Forecast({ sources = [], sourceId, versionId, refreshSou
const data = await res.json()
if (!res.ok) { flash(data.error, 'error'); return }
setLogEntries(prev => prev.filter(e => e.id !== logId))
// grain versions report pf_gkeys, raw versions pf_ids either way these are
// the index values of the rows to drop, and the view re-sums what remains
const removed = data.pf_gkeys ?? data.pf_ids
if (removed?.length && tableRef.current) {
await tableRef.current.remove(removed)
if (data.pf_ids?.length && tableRef.current) {
await tableRef.current.remove(data.pf_ids)
}
flash(`Undone — ${data.rows_deleted} rows removed`)
} catch (err) {

View File

@ -173,11 +173,6 @@ export default function Setup({ refreshSources }) {
const registeredKeys = new Set(sources.map(s => `${s.schema}.${s.tname}`))
// display grain must match grainOf() in lib/sql_generator.js
const grainCols = editedCols
.filter(c => c.in_grain && (c.role === 'dimension' || c.role === 'date'))
.map(c => c.cname)
return (
<div className="h-full flex overflow-hidden text-sm">
@ -283,11 +278,6 @@ export default function Setup({ refreshSources }) {
<div className="px-3 py-2 border-b border-gray-100 flex items-center justify-between shrink-0">
<span className="text-xs font-medium text-gray-500 uppercase tracking-wide">
Col Meta <span className="text-gray-700 normal-case">{selectedSource.schema}.{selectedSource.tname}</span>
<span className="ml-3 normal-case font-normal text-gray-400" title="Columns the forecast load is pre-aggregated to">
grain: {grainCols.length
? <span className="font-mono text-gray-600">{grainCols.join(' × ')}</span>
: <span className="italic">none raw rows</span>}
</span>
</span>
<div className="flex items-center gap-2">
{colsDirty && (
@ -312,7 +302,6 @@ export default function Setup({ refreshSources }) {
<th className="px-3 py-1.5 font-medium">column</th>
<th className="px-3 py-1.5 font-medium">role</th>
<th className="px-3 py-1.5 font-medium text-center">key</th>
<th className="px-3 py-1.5 font-medium text-center" title="Include this column in the display grain — the load is pre-aggregated to the flagged columns">grain</th>
<th className="px-3 py-1.5 font-medium">group</th>
<th className="px-3 py-1.5 font-medium">period col</th>
<th className="px-3 py-1.5 font-medium">label</th>
@ -340,15 +329,6 @@ export default function Setup({ refreshSources }) {
className="cursor-pointer disabled:opacity-20"
/>
</td>
<td className="px-3 py-1.5 text-center">
<input
type="checkbox"
checked={!!col.in_grain}
onChange={e => updateCol(i, 'in_grain', e.target.checked)}
disabled={col.role !== 'dimension' && col.role !== 'date'}
className="cursor-pointer disabled:opacity-20"
/>
</td>
<td className="px-3 py-1.5">
<input
type="text"