/** * Sources Routes * Manage data sources */ const express = require('express'); const multer = require('multer'); const { parse } = require('csv-parse/sync'); const { lit, arr } = require('../lib/sql'); const simplefin = require('../lib/simplefin'); const { inferFields } = require('../lib/fields'); const upload = multer({ storage: multer.memoryStorage() }); module.exports = (pool) => { const router = express.Router(); // Global import log (all sources) router.get('/import-log', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM get_all_import_logs()`); res.json(result.rows); } catch (err) { next(err); } }); // SimpleFIN helpers. Declared before /:name so they aren't shadowed by it. // List the accounts behind a bridge — used to find the account_id for a source router.get('/simplefin-accounts', async (req, res, next) => { try { res.json(await simplefin.listAccounts(req.query.access_url_env)); } catch (err) { if (err instanceof simplefin.SimpleFinError) return res.status(err.status || 502).json({ error: err.message }); next(err); } }); // Sample an account's real transactions and infer its field list — the API // equivalent of uploading a CSV to /suggest. Whatever the account actually // returns is what gets offered, so investment or loan accounts describe // themselves rather than being forced into a checking-account shape. router.get('/simplefin-sample', async (req, res, next) => { try { const { account_id, access_url_env } = req.query; if (!account_id) return res.status(400).json({ error: 'account_id is required' }); const { fetched, records, errors } = await simplefin.fetchTransactions({ accountId: account_id, accessUrlEnv: access_url_env, // The bridge advises staying under 45 days and warns at exactly // 45, so 44 is the largest quiet sample window days: req.query.days !== undefined ? parseInt(req.query.days) : 44, includePending: true, // widen the sample; pending rows can carry extra keys }); const { fields, sampleRows } = inferFields(records); res.json({ fields, sampleRows, fetched, errors }); } catch (err) { if (err instanceof simplefin.SimpleFinError) return res.status(err.status || 502).json({ error: err.message }); next(err); } }); // Exchange a one-shot setup token for the permanent access URL to put in .env router.post('/simplefin-claim', async (req, res, next) => { try { const { setup_token } = req.body || {}; if (!setup_token) return res.status(400).json({ error: 'setup_token is required' }); res.json({ access_url: await simplefin.claimSetupToken(setup_token) }); } catch (err) { if (err instanceof simplefin.SimpleFinError) return res.status(err.status || 502).json({ error: err.message }); next(err); } }); // List all sources router.get('/', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM list_sources()`); res.json(result.rows); } catch (err) { next(err); } }); // Get single source router.get('/:name', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM get_source(${lit(req.params.name)})`); if (result.rows.length === 0) return res.status(404).json({ error: 'Source not found' }); res.json(result.rows[0]); } catch (err) { next(err); } }); // Suggest source definition from CSV router.post('/suggest', upload.single('file'), async (req, res, next) => { try { if (!req.file) return res.status(400).json({ error: 'No file uploaded' }); const records = parse(req.file.buffer, { columns: true, skip_empty_lines: true, trim: true }); if (records.length === 0) return res.status(400).json({ error: 'CSV file is empty' }); const { fields, sampleRows } = inferFields(records); res.json({ name: '', constraint_fields: [], fields, sampleRows }); } catch (err) { next(err); } }); // Create source router.post('/', async (req, res, next) => { try { const { name, constraint_fields, config, global_picklist } = req.body; if (!name || !constraint_fields || !Array.isArray(constraint_fields)) { return res.status(400).json({ error: 'Missing required fields: name, constraint_fields (array)' }); } const result = await pool.query( `SELECT * FROM create_source(${lit(name)}, ${arr(constraint_fields)}, ${lit(config || {})}, ${lit(global_picklist !== false)})` ); res.status(201).json(result.rows[0]); } catch (err) { if (err.code === '23505') return res.status(409).json({ error: 'Source already exists' }); next(err); } }); // Update source router.put('/:name', async (req, res, next) => { try { const { constraint_fields, config, global_picklist } = req.body; const gpVal = global_picklist !== undefined ? lit(global_picklist) : 'NULL'; const result = await pool.query( `SELECT * FROM update_source(${lit(req.params.name)}, ${constraint_fields ? arr(constraint_fields) : 'NULL'}, ${config ? lit(config) : 'NULL'}, ${gpVal})` ); if (result.rows.length === 0) return res.status(404).json({ error: 'Source not found' }); res.json(result.rows[0]); } catch (err) { next(err); } }); // Delete source router.delete('/:name', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM delete_source(${lit(req.params.name)})`); if (result.rows.length === 0) return res.status(404).json({ error: 'Source not found' }); res.json({ success: true, deleted: result.rows[0].delete_source }); } catch (err) { next(err); } }); // Import CSV data and apply transformations to new records router.post('/:name/import', upload.single('file'), async (req, res, next) => { try { if (!req.file) return res.status(400).json({ error: 'No file uploaded' }); const records = parse(req.file.buffer, { columns: true, skip_empty_lines: true, trim: true }); const importResult = await pool.query( `SELECT import_records(${lit(req.params.name)}, ${lit(records)}) as result` ); const importData = importResult.rows[0].result; if (!importData.success) return res.json(importData); const transformResult = await pool.query( `SELECT apply_transformations(${lit(req.params.name)}) as result` ); const transformData = transformResult.rows[0].result; res.json({ ...importData, transform: transformData }); } catch (err) { next(err); } }); // Pull transactions from SimpleFIN and import them, same as a CSV upload. // Safe to re-run: overlapping transactions are skipped by constraint key. router.post('/:name/sync', async (req, res, next) => { try { const sourceResult = await pool.query(`SELECT * FROM get_source(${lit(req.params.name)})`); const source = sourceResult.rows[0]; if (!source || !source.name) return res.status(404).json({ error: 'Source not found' }); const cfg = (source.config || {}).simplefin; if (!cfg || !cfg.account_id) { return res.status(400).json({ error: `Source "${req.params.name}" has no simplefin.account_id in its config` }); } const opts = { ...req.query, ...req.body }; const { fetched, errors, records } = await simplefin.fetchTransactions({ accountId: cfg.account_id, accessUrlEnv: cfg.access_url_env, days: opts.days !== undefined ? parseInt(opts.days) : cfg.days, includePending: opts.include_pending === true || opts.include_pending === 'true', }); if (records.length === 0) { return res.json({ success: true, fetched, errors, imported: 0, duplicates: 0 }); } const importResult = await pool.query( `SELECT import_records(${lit(req.params.name)}, ${lit(records)}) as result` ); const importData = importResult.rows[0].result; if (!importData.success) return res.json({ ...importData, fetched, errors }); const transformResult = await pool.query( `SELECT apply_transformations(${lit(req.params.name)}) as result` ); res.json({ ...importData, fetched, errors, transform: transformResult.rows[0].result }); } catch (err) { if (err instanceof simplefin.SimpleFinError) return res.status(err.status || 502).json({ error: err.message }); next(err); } }); // Get import log router.get('/:name/import-log', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM get_import_log(${lit(req.params.name)})`); res.json(result.rows); } catch (err) { next(err); } }); // Delete an import (removes all records from that batch and the log entry) router.delete('/:name/import-log/:id', async (req, res, next) => { try { const result = await pool.query( `SELECT delete_import(${lit(parseInt(req.params.id))}) as result` ); const data = result.rows[0].result; if (!data.success) return res.status(404).json(data); res.json(data); } catch (err) { next(err); } }); // Apply transformations router.post('/:name/transform', async (req, res, next) => { try { const result = await pool.query(`SELECT apply_transformations(${lit(req.params.name)}) as result`); res.json(result.rows[0].result); } catch (err) { next(err); } }); // Get all known field names for a source router.get('/:name/fields', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM get_source_fields(${lit(req.params.name)})`); res.json(result.rows); } catch (err) { next(err); } }); // Generate output view router.post('/:name/view', async (req, res, next) => { try { const result = await pool.query(`SELECT generate_source_view(${lit(req.params.name)}) as result`); const data = result.rows[0].result; if (data && data.success) { await pool.query(`UPDATE dataflow.sources SET view_generated_at = NOW() WHERE name = ${lit(req.params.name)}`); } res.json(data); } catch (err) { next(err); } }); // Reprocess all records router.post('/:name/reprocess', async (req, res, next) => { try { const result = await pool.query(`SELECT reprocess_records(${lit(req.params.name)}) as result`); res.json(result.rows[0].result); } catch (err) { next(err); } }); // Get statistics router.get('/:name/stats', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM get_source_stats(${lit(req.params.name)})`); res.json(result.rows[0]); } catch (err) { next(err); } }); // Get view data (paginated, sortable) router.get('/:name/view-data', async (req, res, next) => { try { const { limit = 100, offset = 0, sort_col, sort_dir, filters } = req.query; let parsedFilters = null; if (filters) { try { parsedFilters = JSON.parse(filters); } catch { /* ignore bad JSON */ } } const result = await pool.query( `SELECT get_view_data(${lit(req.params.name)}, ${lit(parseInt(limit))}, ${lit(parseInt(offset))}, ${lit(sort_col || null)}, ${lit(sort_dir || 'asc')}, ${parsedFilters ? lit(parsedFilters) : 'NULL'}) as result` ); res.json(result.rows[0].result); } catch (err) { next(err); } }); // Override keys — distinct field names used in overrides across all records for this source router.get('/:name/override-keys', async (req, res, next) => { try { const result = await pool.query( `SELECT DISTINCT jsonb_object_keys(overrides) AS key FROM dataflow.records WHERE source_name = ${lit(req.params.name)} AND overrides IS NOT NULL ORDER BY key` ); res.json(result.rows.map(r => r.key)); } catch (err) { next(err); } }); // Pivot layouts router.get('/:name/layouts', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM list_pivot_layouts(${lit(req.params.name)})`); res.json(result.rows); } catch (err) { next(err); } }); router.post('/:name/layouts', async (req, res, next) => { try { const { layout_name, config } = req.body; if (!layout_name || !config) return res.status(400).json({ error: 'layout_name and config required' }); const result = await pool.query( `SELECT * FROM save_pivot_layout(${lit(req.params.name)}, ${lit(layout_name)}, ${lit(config)})` ); res.json(result.rows[0]); } catch (err) { next(err); } }); router.delete('/:name/layouts/:id', async (req, res, next) => { try { const result = await pool.query(`SELECT * FROM delete_pivot_layout(${lit(parseInt(req.params.id))})`); if (result.rows.length === 0) return res.status(404).json({ error: 'Layout not found' }); res.json({ success: true }); } catch (err) { next(err); } }); return router; };