← back to AbramsOS
routes/import.js
80 lines
// /import — gated by requireStepUp (mounted at /import in server.js).
// Renders per-connector status cards and an "Import all" button.
// The button calls /import/run which fans out to every connected connector's sync endpoint.
const express = require('express');
const db = require('../lib/db');
const audit = require('../lib/audit');
const { CATALOG } = require('../lib/connectors-catalog');
const router = express.Router();
const DEV_USER_ID = 'user_steve';
async function getConnectorStatus() {
const r = await db.query(
`SELECT id, provider, external_subject_id, last_sync_at,
(SELECT count(*) FROM source_message sm WHERE sm.connector_id = ca.id)::int AS message_count
FROM connector_account ca
WHERE ca.user_id = $1 AND (SELECT revoked_at FROM consent_grant cg WHERE cg.id = ca.consent_grant_id) IS NULL
ORDER BY ca.created_at DESC`,
[DEV_USER_ID]
);
return r.rows;
}
router.get('/', async (_req, res) => {
const accounts = await getConnectorStatus();
// Build connector-card view: catalog × actual accounts
const cards = CATALOG.map(c => {
const matched = accounts.filter(a => a.provider === c.provider);
return {
...c,
accounts: matched,
connected: matched.length > 0,
};
});
res.render('import', { cards, plaidEnv: process.env.PLAID_ENV || 'sandbox' });
});
// Fan-out runner: trigger every connected connector's sync endpoint in parallel.
router.post('/run', async (req, res) => {
await audit.log({
actorType: 'user',
actorId: DEV_USER_ID,
objectType: 'user_account',
objectId: DEV_USER_ID,
eventType: 'import_initiated',
metadata: { source: 'import_all_button', ip: req.ip },
});
const accounts = await getConnectorStatus();
if (!accounts.length) return res.json({ ok: true, message: 'no connectors to sync', results: [] });
// Use absolute URL so we hit the same Express via HTTP (preserves middleware + audit).
// Cookie-forward for auth + step-up.
const base = `http://127.0.0.1:${process.env.PORT || 9931}`;
const cookie = req.headers.cookie || '';
const results = await Promise.all(
accounts.map(async (a) => {
const path = a.provider === 'plaid'
? `/api/connectors/plaid/${a.id}/sync`
: `/api/connectors/${a.id}/sync`;
try {
const r = await fetch(`${base}${path}`, {
method: 'POST',
headers: { cookie, 'Content-Type': 'application/json' },
});
const body = await r.json().catch(() => ({}));
return { connector_id: a.id, provider: a.provider, status: r.status, body };
} catch (err) {
return { connector_id: a.id, provider: a.provider, status: 0, error: err.message };
}
})
);
res.json({ ok: true, results });
});
module.exports = router;