← 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;