← back to Nationalrealestate
src/server/watchlist_admin.ts
113 lines
/**
* M5 routes: watchlist CRUD, alerts feed, admin ingest-runs provenance + run-now.
* Mounted from index.ts via mountWatchlistAdmin(app) — kept separate so the
* broker-track and M5 workstreams never edit the same file.
*/
import type { Express } from 'express';
import { execFile } from 'node:child_process';
import { join, dirname } from 'node:path';
import { fileURLToPath } from 'node:url';
import { query } from '../../db/pool.ts';
const __dirname = dirname(fileURLToPath(import.meta.url));
const ROOT = join(__dirname, '..', '..');
// whitelist: job name -> script (never run user-supplied paths)
const RUN_JOBS: Record<string, string> = {
zillow: 'src/ingest/zillow_research.ts',
redfin: 'src/ingest/redfin_tracker.ts',
fhfa: 'src/ingest/fhfa_hpi.ts',
acs: 'src/ingest/census_acs.ts',
derive: 'src/ingest/derive_metrics.ts',
score: 'src/score/opportunity.ts',
alerts: 'src/jobs/alerts_check.ts',
brokers: 'src/ingest/brokers/engine.ts',
};
let jobRunning: string | null = null; // one at a time — these are heavyweight
export function mountWatchlistAdmin(app: Express): void {
app.get('/api/watchlist', async (_req, res) => {
try {
const r = await query(
`SELECT w.id, w.label, w.thresholds, w.created_at,
rg.canonical_key AS key, rg.name, rg.state_code AS state, rg.region_type,
rs.opportunity_score,
(SELECT value FROM metric_series WHERE region_id = rg.id AND metric='zhvi' ORDER BY period DESC LIMIT 1) AS zhvi,
(SELECT value FROM metric_series WHERE region_id = rg.id AND metric='zhvi_yoy' ORDER BY period DESC LIMIT 1) AS zhvi_yoy,
(SELECT value FROM metric_series WHERE region_id = rg.id AND metric='dom' ORDER BY period DESC LIMIT 1) AS dom
FROM watchlist w
JOIN region rg ON rg.id = w.region_id
LEFT JOIN region_scores rs ON rs.region_id = rg.id
AND rs.score_date = (SELECT MAX(score_date) FROM region_scores)
ORDER BY w.created_at DESC`,
);
res.json({ items: r.rows });
} catch (e: any) { res.status(500).json({ error: String(e.message || e) }); }
});
app.post('/api/watchlist', async (req, res) => {
try {
const { key, thresholds, label } = req.body || {};
const rg = await query<{ id: number; name: string }>(
`SELECT id, name FROM region WHERE canonical_key = $1`, [String(key || '')],
);
if (!rg.rows.length) return res.status(400).json({ error: 'unknown region key' });
const r = await query(
`INSERT INTO watchlist (region_id, label, thresholds) VALUES ($1, $2, $3) RETURNING id`,
[rg.rows[0].id, label || rg.rows[0].name, thresholds ? JSON.stringify(thresholds) : null],
);
res.json({ ok: true, id: r.rows[0].id });
} catch (e: any) { res.status(500).json({ error: String(e.message || e) }); }
});
app.delete('/api/watchlist/:id', async (req, res) => {
try {
const id = Number(req.params.id);
if (!Number.isInteger(id)) return res.status(400).json({ error: 'bad id' });
await query(`DELETE FROM alerts_log WHERE watchlist_id = $1`, [id]);
await query(`DELETE FROM watchlist WHERE id = $1`, [id]);
res.json({ ok: true });
} catch (e: any) { res.status(500).json({ error: String(e.message || e) }); }
});
app.get('/api/alerts', async (_req, res) => {
try {
const r = await query(
`SELECT a.id, a.metric, a.observed, a.fired_at, a.detail,
rg.name, rg.state_code AS state, rg.canonical_key AS key
FROM alerts_log a
LEFT JOIN watchlist w ON w.id = a.watchlist_id
LEFT JOIN region rg ON rg.id = w.region_id
ORDER BY a.fired_at DESC LIMIT 50`,
);
res.json({ alerts: r.rows });
} catch (e: any) { res.status(500).json({ error: String(e.message || e) }); }
});
app.get('/api/admin/runs', async (_req, res) => {
try {
const r = await query(
`SELECT id, source, url, file_hash, started_at, finished_at, status,
rows_upserted, rows_skipped, notes
FROM ingest_runs ORDER BY id DESC LIMIT 60`,
);
res.json({ runs: r.rows });
} catch (e: any) { res.status(500).json({ error: String(e.message || e) }); }
});
app.post('/api/admin/run/:job', (req, res) => {
const job = String(req.params.job);
const script = RUN_JOBS[job];
if (!script) return res.status(400).json({ error: 'unknown job' });
if (jobRunning) return res.status(409).json({ error: `job already running: ${jobRunning}` });
jobRunning = job;
execFile('npx', ['tsx', script], { cwd: ROOT, timeout: 45 * 60_000, maxBuffer: 16 * 1024 * 1024 },
(err, stdout, stderr) => {
jobRunning = null;
const tail = (stdout || '').trim().split('\n').slice(-2).join(' | ');
if (err) return res.status(500).json({ error: String(err.message), tail: (stderr || '').slice(-500) });
res.json({ ok: true, summary: tail });
});
});
}