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