Adjust the default maximum value for message processing to improve performance and accommodate larger datasets.
243 lines
7.0 KiB
JavaScript
243 lines
7.0 KiB
JavaScript
#!/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 <q> Gmail search (default: in:inbox newer_than:30d)
|
|
* --ids <a,b> Explicit message IDs (skips list search)
|
|
* --user <email> Mailbox to use (default: first user with tokens)
|
|
* --max <n> 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]) || 500);
|
|
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)) || 500);
|
|
}
|
|
}
|
|
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);
|
|
});
|