Files
nsrbandCursor f270f531ea feat(watcher): add agent Gmail CLI with audit and defaults
Provide list/read/summarize plus reason-gated archive, label, draft, and forward-as-draft, with SQLite audit logging and default-mailbox resolution.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-22 10:09:25 +07:00

233 lines
6.7 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 (default: resolveUser — --user / env / is_default / sole)
* --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 { gmailClientForUser } from "../src/google.js";
import { resolveUser } from "../src/cli/resolve-user.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 resolved = resolveUser(args.user);
if (!resolved.ok) {
console.error(resolved.error);
process.exit(1);
}
const user = resolved.user;
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);
});