← back to AbramsOS

routes/connectors.js

319 lines

const express = require('express');
const crypto = require('crypto');
const db = require('../lib/db');
const audit = require('../lib/audit');
const { id } = require('../lib/ids');
const { decrypt } = require('../lib/crypto');
const fs = require('fs');
const path = require('path');
const fetcher = require('../lib/gmail-fetcher');
const drive = require('../lib/drive-fetcher');
const extractor = require('../lib/receipt-extractor');
const commitmentExtractor = require('../lib/commitment-extractor');
const pdf = require('../lib/pdf-parser');

const UPLOADS_DIR = path.join(__dirname, '..', 'uploads');
fs.mkdirSync(UPLOADS_DIR, { recursive: true });

async function syncDrive(refreshToken, connectorId) {
  const files = await drive.listReceiptFiles(refreshToken, { max: 50 });
  let added = 0;
  for (const f of files) {
    // Skip if already ingested (look up by hash if Drive supplied md5, else by file id)
    const existing = await db.query(
      `SELECT id FROM document WHERE user_id = $1 AND (object_path = $2 OR hash = $3)`,
      [DEV_USER_ID, `drive:${f.id}`, f.md5Checksum || `drive:${f.id}`]
    );
    if (existing.rows.length) continue;

    // Stream the file to uploads/drive-<id>.<ext>
    let ext = '.bin';
    if (f.mimeType === 'application/pdf') ext = '.pdf';
    else if (f.mimeType.startsWith('image/jpeg')) ext = '.jpg';
    else if (f.mimeType.startsWith('image/png')) ext = '.png';

    const localPath = path.join(UPLOADS_DIR, `drive-${f.id}${ext}`);
    try {
      const stream = await drive.downloadFile(refreshToken, f.id);
      await new Promise((resolve, reject) => {
        const out = fs.createWriteStream(localPath);
        stream.on('error', reject);
        out.on('error', reject);
        out.on('finish', resolve);
        stream.pipe(out);
      });
    } catch (err) {
      console.warn('[drive sync] download failed for', f.id, err.message);
      continue;
    }

    const docId = id('document');
    await db.query(
      `INSERT INTO document (id, source_message_id, user_id, kind, mime, hash, object_path, parsed_status)
       VALUES ($1, NULL, $2, $3, $4, $5, $6, 'pending')`,
      [docId, DEV_USER_ID, 'attachment', f.mimeType, f.md5Checksum || `drive:${f.id}`, `drive:${f.id}`]
    );
    added += 1;
    await audit.log({
      actorType: 'system',
      objectType: 'document',
      objectId: docId,
      eventType: 'document_persisted',
      metadata: { source: 'drive', drive_id: f.id, name: f.name, mime: f.mimeType },
    });

    // Auto-parse PDFs (images deferred until OCR ships)
    if (f.mimeType === 'application/pdf') {
      try {
        const parsed = await pdf.parseFile(localPath);
        await db.query(`UPDATE document SET parsed_status = 'parsed' WHERE id = $1`, [docId]);

        const synthSummary = {
          headers: { subject: f.name || `drive:${f.id}`, from: '', date: new Date(f.modifiedTime || Date.now()).toISOString() },
          body: { text: parsed.text, html: '' },
        };
        let extracted = extractor.extract(synthSummary);
        if (extracted && extracted.confidence < extractor.LLM_THRESHOLD) {
          try {
            extracted = await extractor.extractWithFallback(synthSummary, { llm: true, timeoutMs: 25_000 });
          } catch (_) { /* keep heuristic */ }
        }

        if (extracted) {
          const purchaseId = id('purchase');
          await db.query(
            `INSERT INTO purchase
               (id, user_id, source_message_id, merchant_name, merchant_domain, order_number, purchase_date, total_amount, currency, confidence, raw_extract)
             VALUES ($1, $2, NULL, $3, $4, $5, $6, $7, $8, $9, $10)`,
            [purchaseId, DEV_USER_ID, extracted.merchant, extracted.merchantDomain, extracted.orderNumber, extracted.purchaseDate, extracted.total, extracted.currency, extracted.confidence, extracted]
          );
          await audit.log({
            actorType: 'system',
            objectType: 'purchase',
            objectId: purchaseId,
            eventType: 'purchase_extracted',
            metadata: { source: 'drive_pdf', document_id: docId, merchant: extracted.merchant, total: extracted.total, confidence: extracted.confidence, source_tier: extracted.source },
          });
        }
        await audit.log({
          actorType: 'system',
          objectType: 'document',
          objectId: docId,
          eventType: 'document_parsed',
          metadata: { pages: parsed.pages, chars: parsed.text.length },
        });
      } catch (err) {
        console.warn('[drive-pdf-parse]', f.id, err.message);
        await db.query(`UPDATE document SET parsed_status = 'failed' WHERE id = $1`, [docId]);
      }
    }
  }
  return { fetched: files.length, added };
}

const router = express.Router();
const DEV_USER_ID = 'user_steve';

router.get('/connectors', async (_req, res) => {
  const result = await db.query(
    `SELECT ca.id, ca.provider, ca.external_subject_id, ca.last_sync_at,
            cg.granted_at, cg.revoked_at,
            (SELECT count(*) FROM source_message sm WHERE sm.connector_id = ca.id)::int AS message_count
       FROM connector_account ca
       JOIN consent_grant cg ON cg.id = ca.consent_grant_id
      WHERE ca.user_id = $1
      ORDER BY cg.granted_at DESC`,
    [DEV_USER_ID]
  );
  res.render('connectors', { connectors: result.rows });
});

router.get('/api/connectors', async (_req, res) => {
  const result = await db.query(
    `SELECT id, provider, external_subject_id, last_sync_at FROM connector_account WHERE user_id = $1`,
    [DEV_USER_ID]
  );
  res.json(result.rows);
});

router.post('/api/connectors/:id/sync', async (req, res) => {
  const connectorId = req.params.id;
  const result = await db.query(
    `SELECT id, refresh_token_encrypted, refresh_token_iv, refresh_token_tag
       FROM connector_account WHERE id = $1 AND user_id = $2`,
    [connectorId, DEV_USER_ID]
  );
  if (!result.rows.length) return res.status(404).json({ error: 'connector not found' });
  const c = result.rows[0];

  let refreshToken;
  try {
    refreshToken = decrypt(c.refresh_token_encrypted, c.refresh_token_iv, c.refresh_token_tag);
  } catch (err) {
    return res.status(500).json({ error: 'failed to decrypt refresh token: ' + err.message });
  }

  await audit.log({
    actorType: 'user',
    actorId: DEV_USER_ID,
    objectType: 'connector_account',
    objectId: connectorId,
    eventType: 'sync_started',
  });

  let inserted = 0;
  let purchases = 0;
  try {
    const ids = await fetcher.listReceiptIds(refreshToken, { max: 50 });
    for (const mid of ids) {
      // Skip if already ingested
      const existing = await db.query(
        `SELECT id FROM source_message WHERE connector_id = $1 AND source_type = 'gmail' AND external_id = $2`,
        [connectorId, mid]
      );
      if (existing.rows.length) continue;

      const message = await fetcher.getMessage(refreshToken, mid);
      const summary = fetcher.summarize(message);
      const payloadHash = crypto.createHash('sha256').update(JSON.stringify(message)).digest('hex');

      const sourceId = id('source');
      await db.query(
        `INSERT INTO source_message
           (id, connector_id, source_type, external_id, thread_id, received_at, subject, sender, recipient, payload_jsonb, payload_hash)
         VALUES ($1, $2, 'gmail', $3, $4, $5, $6, $7, $8, $9, $10)
         ON CONFLICT (connector_id, source_type, external_id) DO NOTHING`,
        [
          sourceId,
          connectorId,
          mid,
          summary.threadId,
          summary.headers.date ? new Date(summary.headers.date) : new Date(),
          summary.headers.subject,
          summary.headers.from,
          summary.headers.to,
          message,
          payloadHash,
        ]
      );
      inserted += 1;

      await audit.log({
        actorType: 'system',
        objectType: 'source_message',
        objectId: sourceId,
        eventType: 'message_persisted',
        metadata: { gmail_id: mid, subject: summary.headers.subject?.slice(0, 120) },
      });

      // Tier 1 heuristic; if conf < threshold and Mac1 Ollama is reachable, enrich with LLM.
      let extracted = extractor.extract(summary);
      if (extracted && extracted.confidence < extractor.LLM_THRESHOLD) {
        try {
          extracted = await extractor.extractWithFallback(summary, { llm: true, timeoutMs: 25_000 });
        } catch (_) { /* fall through with heuristic-only result */ }
      }
      if (extracted) {
        const purchaseId = id('purchase');
        await db.query(
          `INSERT INTO purchase
             (id, user_id, source_message_id, merchant_name, merchant_domain, order_number, purchase_date, total_amount, currency, confidence, raw_extract)
           VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)`,
          [
            purchaseId,
            DEV_USER_ID,
            sourceId,
            extracted.merchant,
            extracted.merchantDomain,
            extracted.orderNumber,
            extracted.purchaseDate,
            extracted.total,
            extracted.currency,
            extracted.confidence,
            extracted,
          ]
        );
        purchases += 1;

        await audit.log({
          actorType: 'system',
          objectType: 'purchase',
          objectId: purchaseId,
          eventType: 'purchase_extracted',
          metadata: { merchant: extracted.merchant, total: extracted.total, source: extracted.source, confidence: extracted.confidence },
        });

        // Extract any service commitments from the same email body
        try {
          const commitments = commitmentExtractor.extract(summary);
          for (const c of commitments) {
            const commitmentId = id('purchase'); // reuse ulid prefix
            const refundEnd = commitmentExtractor.computeWindowEnd(extracted.purchaseDate, c.window_days);
            await db.query(
              `INSERT INTO service_commitment
                 (id, user_id, purchase_id, source_message_id, provider_name,
                  commitment_type, guarantee_text, promised_outcome, window_days,
                  refund_window_ends_at, conditions, source, confidence, raw_extract)
               VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14)`,
              [
                commitmentId, DEV_USER_ID, purchaseId, sourceId,
                extracted.merchant, c.commitment_type, c.guarantee_text,
                c.promised_outcome, c.window_days, refundEnd, c.conditions,
                c.source, c.confidence, JSON.stringify(c),
              ]
            );
            await audit.log({
              actorType: 'system',
              objectType: 'service_commitment',
              objectId: commitmentId,
              eventType: 'commitment_extracted',
              metadata: { type: c.commitment_type, window_days: c.window_days, source: c.source, confidence: c.confidence },
            });
          }
        } catch (err) {
          console.warn('[commitment-extract]', err.message);
        }
      }
    }

    // Walk Drive too (same Google account, drive.readonly scope)
    let driveStats = { fetched: 0, added: 0 };
    try {
      driveStats = await syncDrive(refreshToken, connectorId);
    } catch (err) {
      console.warn('[drive sync] failed:', err.message);
      await audit.log({
        actorType: 'system',
        objectType: 'connector_account',
        objectId: connectorId,
        eventType: 'drive_sync_failed',
        metadata: { error: err.message },
      });
    }

    await db.query(`UPDATE connector_account SET last_sync_at = now() WHERE id = $1`, [connectorId]);

    await audit.log({
      actorType: 'system',
      objectType: 'connector_account',
      objectId: connectorId,
      eventType: 'sync_completed',
      metadata: { fetched: ids.length, inserted, purchases, drive: driveStats },
    });

    res.json({ ok: true, fetched: ids.length, inserted, purchases, drive: driveStats });
  } catch (err) {
    console.error('[sync]', err);
    await audit.log({
      actorType: 'system',
      objectType: 'connector_account',
      objectId: connectorId,
      eventType: 'sync_failed',
      metadata: { error: err.message },
    });
    res.status(500).json({ error: err.message });
  }
});

module.exports = router;