← back to AbramsOS
amazon-orders: 30-min Gmail poller (George bridge) → purchase + reorder_item, word-order-safe dedup; wired into scheduler
f450f59a61cdad5634c7e811f2ec9bdf4d8a3fd3 · 2026-08-18 15:32:54 -0700 · Steve
Files touched
A lib/amazon-orders.jsM lib/scheduler.js
Diff
commit f450f59a61cdad5634c7e811f2ec9bdf4d8a3fd3
Author: Steve <steve@designerwallcoverings.com>
Date: Tue Aug 18 15:32:54 2026 -0700
amazon-orders: 30-min Gmail poller (George bridge) → purchase + reorder_item, word-order-safe dedup; wired into scheduler
---
lib/amazon-orders.js | 96 ++++++++++++++++++++++++++++++++++++++++++++++++++++
lib/scheduler.js | 18 ++++++++++
2 files changed, 114 insertions(+)
diff --git a/lib/amazon-orders.js b/lib/amazon-orders.js
new file mode 100644
index 0000000..204ec51
--- /dev/null
+++ b/lib/amazon-orders.js
@@ -0,0 +1,96 @@
+// lib/amazon-orders.js — poll Gmail (via the George bridge) for Amazon order
+// confirmations and wire each into `purchase` + `reorder_item`. Runs every 30 min
+// from the in-process scheduler. Dedupes on the Gmail message id, so re-runs are safe.
+//
+// Gmail access = George HTTP bridge (GEORGE_URL + GEORGE_BASIC_AUTH). No-ops cleanly
+// if George isn't configured (same credential that gates the claims deadline email).
+// Honest: total_amount stays NULL unless a price is present; nothing is bought.
+
+const db = require('./db');
+const { id } = require('./ids');
+
+const GEORGE_URL = (process.env.GEORGE_URL || '').replace(/\/$/, '');
+const GEORGE_BASIC_AUTH = process.env.GEORGE_BASIC_AUTH || '';
+const USER = process.env.ABRAMSOS_USER_ID || 'user_steve';
+const ACCOUNTS = (process.env.AMAZON_ORDER_ACCOUNTS || 'info,steve-personal').split(',').map(s => s.trim());
+const QUERY = 'from:auto-confirm@amazon.com subject:Ordered newer_than:190d';
+
+// "Ordered: \"Nespresso Capsules Vertuo,...\" and 4 more items" -> "Nespresso Vertuo Capsules"
+function productFromSubject(subj) {
+ let s = String(subj || '').replace(/^Ordered:\s*/i, '').replace(/\s*and\s+[\d]+\s+more items?\.?$/i, '');
+ s = s.replace(/["“”]/g, '').replace(/,?\.\.\.$/,'').replace(/[,\s]+$/,'').trim();
+ if (/^\d+\s+\w+\s+item/i.test(s) || !s) return null; // generic "1 Office item"
+ return s;
+}
+// crude category for the reorder tracker
+function categoryFor(name) {
+ const n = name.toLowerCase();
+ if (/nespresso|coffee|capsule|pod/.test(n)) return 'coffee';
+ if (/label|shipping|poly bag|box|tape|office|printer|ink|toner/.test(n)) return 'office/shipping';
+ if (/headset|plantronics|ssd|adapter|usb|cable|electronic/.test(n)) return 'electronics';
+ return 'general';
+}
+
+async function georgeSearch(account) {
+ const url = `${GEORGE_URL}/api/messages?account=${encodeURIComponent(account)}&maxResults=50&q=${encodeURIComponent(QUERY)}`;
+ const res = await fetch(url, { headers: { Authorization: 'Basic ' + GEORGE_BASIC_AUTH } });
+ if (!res.ok) throw new Error(`George /api/messages ${res.status}`);
+ const j = await res.json();
+ return j.messages || j || [];
+}
+
+// Insert one order as a purchase (dedupe on source_message_id) + upsert a reorder_item.
+async function wireOrder(o) {
+ const gmailId = o.id;
+ const subj = o.subject || '';
+ const when = o.internalDate ? new Date(Number(o.internalDate)) : (o.date ? new Date(o.date) : new Date());
+ const product = productFromSubject(subj);
+
+ const ins = 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,'Amazon','amazon.com',NULL,$4,NULL,'USD',0.6,$5)
+ ON CONFLICT (source_message_id) DO NOTHING RETURNING id`,
+ [id('purchase'), USER, gmailId, when.toISOString().slice(0,10), JSON.stringify({ subject: subj, product, gmail_id: gmailId })]
+ ).catch(() => ({ rows: [] }));
+ const newPurchase = ins.rows.length > 0;
+
+ let reorderTouched = false;
+ if (product) {
+ // upsert reorder_item — match word-order-insensitively so "Nespresso Capsules Vertuo"
+ // and "Nespresso Vertuo Capsules" are the same item (no dup).
+ const norm = (s) => String(s||'').toLowerCase().replace(/[^a-z0-9 ]/g,'').split(/\s+/).filter(Boolean).sort().join(' ');
+ const all = await db.query(`SELECT id, name, reorder_cadence_days FROM reorder_item WHERE user_id=$1`, [USER]).catch(() => ({ rows: [] }));
+ const match = all.rows.find(r => norm(r.name) === norm(product));
+ const ex = { rows: match ? [match] : [] };
+ if (ex.rows.length) {
+ const cad = ex.rows[0].reorder_cadence_days || 90;
+ await db.query(`UPDATE reorder_item SET last_ordered_at=GREATEST(coalesce(last_ordered_at,'1970-01-01'),$2::timestamptz), next_due_date=$2::date + ($3||' days')::interval, updated_at=now() WHERE id=$1`, [ex.rows[0].id, when.toISOString(), cad]).catch(()=>{});
+ } else {
+ await db.query(
+ `INSERT INTO reorder_item (id,user_id,name,merchant,category,currency,last_ordered_at,status,notes,created_at,updated_at)
+ VALUES ($1,$2,$3,'Amazon',$4,'USD',$5,'active',$6,now(),now())`,
+ [id('reorder'), USER, product, categoryFor(product), when.toISOString(), 'auto-imported from Amazon order ' + when.toISOString().slice(0,10)]
+ ).catch(()=>{});
+ }
+ reorderTouched = true;
+ }
+ return { newPurchase, reorderTouched, product };
+}
+
+async function pollAndImport() {
+ if (!GEORGE_URL || !GEORGE_BASIC_AUTH) return { skipped: 'george-not-configured', imported: 0, reorders: 0 };
+ let imported = 0, reorders = 0, scanned = 0;
+ for (const acct of ACCOUNTS) {
+ let msgs = [];
+ try { msgs = await georgeSearch(acct); } catch (e) { continue; }
+ for (const m of msgs) {
+ scanned++;
+ const r = await wireOrder(m);
+ if (r.newPurchase) imported++;
+ if (r.reorderTouched) reorders++;
+ }
+ }
+ return { imported, reorders, scanned };
+}
+
+module.exports = { pollAndImport, wireOrder, productFromSubject, categoryFor };
diff --git a/lib/scheduler.js b/lib/scheduler.js
index a17caf7..3d1d73a 100644
--- a/lib/scheduler.js
+++ b/lib/scheduler.js
@@ -11,6 +11,7 @@ const matcher = require('./recall-matcher');
const digest = require('./digest');
const autopilot = require('./claims-autopilot');
const claimsAlert = require('./claims-alert');
+const amazonOrders = require('./amazon-orders');
const { id } = require('./ids');
const LOG_FILE = path.join(__dirname, '..', 'logs', 'scheduler.log');
@@ -100,6 +101,15 @@ async function runClaimsAutopilot() {
}
}
+async function runAmazonOrders() {
+ try {
+ const r = await amazonOrders.pollAndImport();
+ logLine(`amazon-orders: ${r.skipped ? 'skipped — ' + r.skipped : `imported=${r.imported} reorders=${r.reorders} scanned=${r.scanned}`}`);
+ } catch (err) {
+ logLine(`amazon-orders FAILED: ${err.message}`);
+ }
+}
+
let started = false;
function start() {
if (started) return;
@@ -125,6 +135,14 @@ function start() {
setTimeout(() => { runClaimsAutopilot(); }, 15000); // boot pass, after DB warms
}
+ // Amazon order poller — every 30 min, import new order confirmations into
+ // purchase + reorder_item (dedup by Gmail id). No-ops without George config.
+ const ordersOn = process.env.AMAZON_ORDERS_DISABLED !== '1';
+ if (ordersOn) {
+ cron.schedule('*/30 * * * *', runAmazonOrders);
+ setTimeout(() => { runAmazonOrders(); }, 25000);
+ }
+
started = true;
logLine(`scheduler started · reminders every 4h · cpsc daily 03:17 · digest ${digestOn ? 'daily 08:00' : 'DISABLED'} · claims-autopilot ${claimsOn ? 'every 6h' : 'DISABLED'}`);
}
← b2fd5b1 deploy: abramsos vhost with fleet Basic Auth (admin/DW2024!)
·
back to AbramsOS
·
settlements: 3 views (grid/list/table sortable) + one-click 103c476 →