#!/usr/bin/env node /** * Manually process Gmail messages through match → parse → PocketBase save. * * Use when forwarding old receipts: Gmail Date may be "today", but trx_date * comes from the receipt body (Tanggal/Jam, Picked up on, …). Existing * invoice_id rows are skipped (unique constraint). * * Also expands .eml / message/rfc822 attachments (forward-as-attachment). * * Usage: * npm run process-messages -w watcher -- --query 'newer_than:2d subject:"Pembayaran Berhasil!"' * npm run process-messages -w watcher -- --ids 1a099bd61efd234d,1a099b8b9d239e5e * npm run process-messages -w watcher -- --dry-run --max 20 * * Options: * --query Gmail search (default: in:inbox newer_than:30d) * --ids Explicit message IDs (skips list search) * --user Mailbox to use (default: first user with tokens) * --max Max messages to process (default 500) * --dry-run Parse and log only; do not write PocketBase */ import { GOOGLE_CLIENT_ID, GOOGLE_CLIENT_SECRET, POCKETBASE_URL, } from "../src/config.js"; import { getUserByEmail, listUsersWithTokens } from "../src/db.js"; import { gmailClientForUser } from "../src/google.js"; import { fetchAndLogMessage } from "../src/handlers/fetch-and-log-message.js"; import { matchRule } from "../src/handlers/match-rule.js"; import { parseMatchedMessage } from "../src/handlers/parse-matched-message.js"; import { saveSpending } from "../src/handlers/save-spending.js"; import { runPipeline } from "../src/pipeline.js"; import { processFetchedMessage } from "../src/process-fetched.js"; import { pocketbaseConfigured } from "../src/pocketbase/client.js"; const DEFAULT_QUERY = "in:inbox newer_than:30d"; function parseArgs(argv) { const out = { query: null, ids: [], user: null, max: 500, dryRun: false, help: false, }; for (let i = 0; i < argv.length; i++) { const a = argv[i]; if (a === "--help" || a === "-h") out.help = true; else if (a === "--dry-run") out.dryRun = true; else if (a === "--query") out.query = argv[++i]; else if (a === "--ids") { out.ids = String(argv[++i] || "") .split(",") .map((s) => s.trim()) .filter(Boolean); } else if (a === "--user") out.user = argv[++i]; else if (a === "--max") out.max = Math.max(1, Number(argv[++i]) || 50); else if (a.startsWith("--query=")) out.query = a.slice("--query=".length); else if (a.startsWith("--ids=")) { out.ids = a .slice("--ids=".length) .split(",") .map((s) => s.trim()) .filter(Boolean); } else if (a.startsWith("--user=")) out.user = a.slice("--user=".length); else if (a.startsWith("--max=")) { out.max = Math.max(1, Number(a.slice("--max=".length)) || 50); } } return out; } function makeLog() { const line = (level, obj, msg) => { const payload = typeof obj === "string" ? { msg: obj } : { ...(obj || {}), msg }; const { msg: message, ...rest } = payload; console.log( JSON.stringify({ level, time: new Date().toISOString(), ...rest, msg: message || msg, }), ); }; return { info: (obj, msg) => line("info", obj, msg), warn: (obj, msg) => line("warn", obj, msg), error: (obj, msg) => line("error", obj, msg), }; } async function listMessageIds(gmail, query, max) { const ids = []; let pageToken; while (ids.length < max) { const { data } = await gmail.users.messages.list({ userId: "me", q: query, maxResults: Math.min(50, max - ids.length), pageToken, }); for (const m of data.messages || []) { if (m.id) ids.push(m.id); if (ids.length >= max) break; } pageToken = data.nextPageToken; if (!pageToken) break; } return ids; } async function main() { const args = parseArgs(process.argv.slice(2)); if (args.help) { console.log(`Usage: process-messages [--query q] [--ids id1,id2] [--user email] [--max n] [--dry-run] Default query: ${DEFAULT_QUERY}`); process.exit(0); } if (!GOOGLE_CLIENT_ID || !GOOGLE_CLIENT_SECRET) { console.error("Missing GOOGLE_CLIENT_ID / GOOGLE_CLIENT_SECRET in .env"); process.exit(1); } const users = listUsersWithTokens(); if (!users.length) { console.error("No signed-in users with tokens. Sign in via the web UI first."); process.exit(1); } let user = users[0]; if (args.user) { const byEmail = getUserByEmail(args.user); const byId = users.find((u) => String(u.id) === String(args.user)); user = byEmail || byId || null; if (!user) { console.error(`User not found: ${args.user}`); process.exit(1); } } const log = makeLog(); const gmail = await gmailClientForUser(user.id); let messageIds = args.ids; if (!messageIds.length) { const query = args.query || DEFAULT_QUERY; log.info({ query, max: args.max, user: user.email }, "process-messages: listing"); messageIds = await listMessageIds(gmail, query, args.max); } else { messageIds = messageIds.slice(0, args.max); } log.info( { count: messageIds.length, user: user.email, dryRun: args.dryRun, pocketbase: pocketbaseConfigured() ? POCKETBASE_URL : null, }, "process-messages: start", ); const fetchHandlers = [fetchAndLogMessage]; const processHandlers = args.dryRun ? [matchRule, parseMatchedMessage] : [matchRule, parseMatchedMessage, saveSpending]; const summary = { processed: 0, matched: 0, saved: 0, duplicate: 0, skipped: 0, errors: 0, }; for (const messageId of messageIds) { try { const fetched = await runPipeline(fetchHandlers, { notification: { emailAddress: user.email, historyId: null }, userId: user.id, messageId, gmail, log, }); if (!fetched.message) { summary.skipped++; continue; } const results = await processFetchedMessage(fetched, processHandlers); for (const result of results) { summary.processed++; if (!result.rule) { summary.skipped++; continue; } summary.matched++; if (args.dryRun) { log.info( { messageId: result.messageId, kind: result.parsed?.kind, invoice_id: result.parsed?.reference, amount: result.parsed?.amount, occurredAt: result.parsed?.occurredAt, }, "process-messages: dry-run parsed", ); continue; } if (result.saved) summary.saved++; else if (result.duplicate) summary.duplicate++; else if (result.error) summary.errors++; else summary.skipped++; } } catch (err) { summary.errors++; log.error( { messageId, err: err.message || String(err) }, "process-messages: failed", ); } } log.info(summary, "process-messages: done"); } main().catch((err) => { console.error(err); process.exit(1); });