From c97c029b3ceb530db6b4e9b7cab44135cfd389be Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 21 Sep 2026 07:56:40 +0000 Subject: [PATCH] Mirror transactions into a structured D1 store Transaction data is currently only kept as one uploaded file per transaction in the vector store, which makes date filters and totals approximate and leaves no way to correct or deduplicate a record. Write every transaction into a D1 table in parallel with the existing upload. The vector store stays the read path, so this change only accumulates structured data; no query behaviour changes. Schema notes: * Amounts are stored as integers in minor units, with the merchant-side amount kept separately when a bank settles a foreign charge in another currency. * The exchange rate is resolved and frozen when a transaction is recorded, so historical reports stop moving as rates change. A transaction is never blocked on a rate lookup; a daily job backfills any row recorded while no rate was available, using the rate for that row's own date. * An FTS5 index covers the free-text columns with diacritics folded, so untoned Vietnamese queries match. Duplicates are handled in two layers. A unique constraint drops exact repeats from delivery retries. A merchant receipt and the matching bank debit, which arrive minutes apart under different senders, are linked through duplicate_of instead of dropped, so both stay queryable while totals count the payment once. Direction separates that case from a transfer between accounts, which produces a debit and a credit and must remain two records. Writes are best-effort: failures are logged and never interrupt notification or the vector store upload, and the write is skipped when no binding exists. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01SwuhAs5n7j6dGB3fcaqCDp --- README.md | 26 ++++ handlers/assistant.ts | 6 +- handlers/transactions.ts | 31 +++- index.ts | 5 + migrations/0001_create_transactions.sql | 68 +++++++++ services/fx.ts | 137 +++++++++++++++++ services/telegram.ts | 24 +-- services/transactions-store.ts | 162 ++++++++++++++++++++ tests/transactions-store.test.ts | 190 ++++++++++++++++++++++++ types/index.ts | 44 ++++++ utils/money.ts | 61 ++++++++ wrangler.jsonc | 14 +- 12 files changed, 737 insertions(+), 31 deletions(-) create mode 100644 migrations/0001_create_transactions.sql create mode 100644 services/fx.ts create mode 100644 services/transactions-store.ts create mode 100644 tests/transactions-store.test.ts create mode 100644 utils/money.ts diff --git a/README.md b/README.md index 90ece8c..3616348 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,9 @@ Transaction emails will be forwarded to a "virtual" email address managed by [Cl 1. Trigger a notification (currently set to send alerts via Telegram). 2. Upload the processed text to the [vector database store](https://platform.openai.com/storage/vector_stores) on the OpenAI platform. +3. Mirror the same transaction into a [Cloudflare D1](https://developers.cloudflare.com/d1/) table as structured rows. + +The vector store remains the read path: questions asked through Telegram are still answered with file search. The D1 table is written in parallel so the structured data accumulates and can be verified before any query is switched over to it. Since the data is stored in a personal vector database, you can make queries by sending a message to your Telegram bot. The bot will then call the Cloudflare worker using a Telegram webhook. These "on-demand" requests will be processed by the [OpenAI Responses API](https://platform.openai.com/docs/api-reference/responses) with file search over the configured vector store. @@ -69,10 +72,33 @@ The application requires the following environment variables: | `OPENAI_ASSISTANT_MODEL` | The Responses API model used for transaction questions and report runs. | No | `gpt-5.6-luna` | | `OPENAI_ASSISTANT_ROUTER_MODEL` | The model used to route Telegram messages to assistant functions. | No | `gpt-5.6-luna` | | `OPENAI_ASSISTANT_VECTORSTORE_ID`| The vector store identifier for storing processed data in OpenAI and answering questions with file search. | Yes | - | +| `TRANSACTION_DUPLICATE_WINDOW_MINUTES` | Minutes either side of a transaction to look for the same payment reported by a second sender. | No | `15` | + +## Structured transaction store + +Transactions are mirrored into a D1 database alongside the vector store upload. Create the database and apply the migration: + +```bash +bunx wrangler d1 create personalaccountant +# copy the returned database_id into the d1_databases block in wrangler.jsonc +bunx wrangler d1 migrations apply personalaccountant --remote +``` + +The schema keeps amounts as integers in minor units, records the merchant-side amount separately when a bank settles a foreign charge in another currency, and freezes the exchange rate used at the time a transaction is recorded so historical reports do not drift. A `transactions_fts` FTS5 index covers the free-text columns with diacritics folded, so untoned Vietnamese queries still match. + +Two kinds of duplicate are handled differently: + +* **Delivery retries** — the same email or webhook arriving twice is dropped by a unique constraint. +* **One payment, two senders** — a merchant receipt and the matching bank debit arrive minutes apart. The second row is kept but linked to the first through `duplicate_of`, so both remain queryable while totals count the payment once. A debit paired with a credit is treated as a transfer between accounts rather than a duplicate. + +Writes to D1 are best-effort: a failure is logged and does not interrupt notification or the vector store upload. If no `DB` binding is configured, the structured write is skipped entirely. + +A daily job refreshes exchange rates and fills in any transaction recorded while no rate was available, converting each row with the rate for its own date. A transaction is never blocked on an exchange rate lookup. ## TODO - [ ] Whitelist email addresses. - [ ] Notify to channel, group chat instead +- [ ] Switch the query path from vector store file search to the structured store ## Additional information diff --git a/handlers/assistant.ts b/handlers/assistant.ts index 20ed452..dde653e 100644 --- a/handlers/assistant.ts +++ b/handlers/assistant.ts @@ -1,7 +1,7 @@ import { Buffer } from 'node:buffer'; import { formatDate } from '../utils/date'; import { createOpenAIClient } from '../services/openai'; -import { processTransaction, storeTransaction, notifyServices } from './transactions'; +import { processTransaction, persistTransaction, notifyServices } from './transactions'; import { buildMessageWithReplyContext, sendTelegramMessage } from '../services/telegram'; import type { Environment } from '../types'; @@ -100,7 +100,7 @@ const assistantOcr = async (message, c) => { const transactionDetails = await processTransaction(transaction, c.env); if (!transactionDetails) return 'Not okay'; - await Promise.all([storeTransaction(transactionDetails, c.env), notifyServices(transactionDetails, c.env)]); + await Promise.all([persistTransaction(transactionDetails, c.env, 'ocr'), notifyServices(transactionDetails, c.env)]); return '📬 Email processed successfully'; }; @@ -119,7 +119,7 @@ const assistantManualTransaction = async (transaction, env: Environment) => { const transactionDetails = await processTransaction(buildManualTransactionInput(transaction), env, 'manual'); if (!transactionDetails) return 'Not okay'; - await Promise.all([storeTransaction(transactionDetails, env), notifyServices(transactionDetails, env, '✅ *Đã thêm giao dịch thủ công*')]); + await Promise.all([persistTransaction(transactionDetails, env, 'manual'), notifyServices(transactionDetails, env, '✅ *Đã thêm giao dịch thủ công*')]); return '📬 Email processed successfully'; }; diff --git a/handlers/transactions.ts b/handlers/transactions.ts index 03f0ac8..d26de18 100644 --- a/handlers/transactions.ts +++ b/handlers/transactions.ts @@ -2,7 +2,8 @@ import PostalMime from 'postal-mime'; import { Buffer } from 'node:buffer'; import { createOpenAIClient } from '../services/openai'; import { formatTransactionDetails, sendTelegramMessage } from '../services/telegram'; -import type { Environment } from '../types'; +import { saveTransaction } from '../services/transactions-store'; +import type { Environment, TransactionDetails } from '../types'; export const processTransaction = async (emailData: string, env: Environment, source: 'email' | 'manual' | 'ocr' = 'email') => { console.log(`🤖 Processing ${source} content: ${emailData}`); @@ -96,6 +97,32 @@ export const storeTransaction = async (details, env: Environment) => { console.info(`🤖 Add ${fileName} to Vector store successfully`); }; +/** + * Writes the transaction to D1 alongside the vector store upload. + * + * The vector store remains the read path, so this is additive: a D1 failure is + * logged and swallowed rather than allowed to break notification or the + * existing upload. + */ +export const storeTransactionRecord = async (details: TransactionDetails, env: Environment, source: 'email' | 'manual' | 'ocr') => { + if (!env.DB) { + console.warn('🗄️ No D1 binding configured, skipping structured write'); + return null; + } + + try { + return await saveTransaction(env, details, source); + } catch (error) { + console.error('🗄️ Failed to store transaction in D1', error); + return null; + } +}; + +/** Uploads to the vector store and mirrors the transaction into D1. */ +export const persistTransaction = async (details: TransactionDetails, env: Environment, source: 'email' | 'manual' | 'ocr') => { + await Promise.all([storeTransaction(details, env), storeTransactionRecord(details, env, source)]); +}; + export const notifyServices = async (details: any, env: Environment, headline?: string) => { const message = formatTransactionDetails(details, headline); await sendTelegramMessage(env, message); @@ -115,6 +142,6 @@ export const email = async (message, env: Environment) => { if (!transactionDetails) return "Not okay"; - await Promise.all([storeTransaction(transactionDetails, env), notifyServices(transactionDetails, env)]); + await Promise.all([persistTransaction(transactionDetails, env, 'email'), notifyServices(transactionDetails, env)]); return "📬 Email processed successfully"; }; diff --git a/index.ts b/index.ts index 94f825a..76aabd1 100644 --- a/index.ts +++ b/index.ts @@ -5,6 +5,7 @@ import { requestId } from 'hono/request-id'; import { secureHeaders } from 'hono/secure-headers'; import { dailyReport, handleAssistantRequest, monthlyReport, weeklyReport } from './handlers/assistant'; import { email as processEmail } from './handlers/transactions'; +import { backfillExchangeRates } from './services/fx'; import type { Environment } from './types'; export { buildMessageWithReplyContext, convertCurrencyAmountsToVnd, formatCurrencyAmounts, formatTransactionDetails, normalize, stripTelegramMarkdown } from './services/telegram'; @@ -45,6 +46,10 @@ export default { console.info("⏰ Monthly scheduler triggered"); await monthlyReport(env); break; + case "0 18 * * *": + console.info("💱 Exchange rate scheduler triggered"); + if (env.DB) await backfillExchangeRates(env); + break; } }, diff --git a/migrations/0001_create_transactions.sql b/migrations/0001_create_transactions.sql new file mode 100644 index 0000000..ae227d6 --- /dev/null +++ b/migrations/0001_create_transactions.sql @@ -0,0 +1,68 @@ +-- Structured transaction store. +-- Written in parallel with the existing vector store upload; queries still run +-- against the vector store until the read path is switched over. + +CREATE TABLE IF NOT EXISTS transactions ( + id TEXT PRIMARY KEY, + occurred_at TEXT NOT NULL, -- ISO8601 UTC, sortable + amount_minor INTEGER NOT NULL, -- amount that hit the account, in minor units + currency TEXT NOT NULL, + original_amount_minor INTEGER, -- merchant-side amount when the bank converted it + original_currency TEXT, + amount_vnd_minor INTEGER, -- frozen at ingest, NULL until the rate is known + fx_rate REAL, + fx_rate_as_of TEXT, + bank_name TEXT NOT NULL DEFAULT '', -- '' not NULL: NULLs defeat the UNIQUE constraint + category TEXT, + direction TEXT NOT NULL DEFAULT 'debit' CHECK (direction IN ('debit', 'credit')), + source TEXT NOT NULL CHECK (source IN ('email', 'manual', 'ocr')), + source_kind TEXT NOT NULL DEFAULT 'bank' CHECK (source_kind IN ('bank', 'merchant')), + message TEXT NOT NULL, + plain_data TEXT NOT NULL, + duplicate_of TEXT REFERENCES transactions (id), -- NULL = counts toward totals + needs_review INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL, + UNIQUE (occurred_at, amount_minor, currency, bank_name) +); + +CREATE INDEX IF NOT EXISTS idx_transactions_occurred_at ON transactions (occurred_at); +CREATE INDEX IF NOT EXISTS idx_transactions_category ON transactions (category, occurred_at); +CREATE INDEX IF NOT EXISTS idx_transactions_duplicate_of ON transactions (duplicate_of); +CREATE INDEX IF NOT EXISTS idx_transactions_pending_fx ON transactions (currency, amount_vnd_minor); + +-- Full-text index over the free-text columns. External content table: rows live +-- in `transactions`, this holds only the index. +CREATE VIRTUAL TABLE IF NOT EXISTS transactions_fts USING fts5 ( + message, + plain_data, + content='transactions', + content_rowid='rowid', + tokenize='unicode61 remove_diacritics 2' +); + +CREATE TRIGGER IF NOT EXISTS transactions_fts_insert AFTER INSERT ON transactions BEGIN + INSERT INTO transactions_fts (rowid, message, plain_data) + VALUES (new.rowid, new.message, new.plain_data); +END; + +CREATE TRIGGER IF NOT EXISTS transactions_fts_delete AFTER DELETE ON transactions BEGIN + INSERT INTO transactions_fts (transactions_fts, rowid, message, plain_data) + VALUES ('delete', old.rowid, old.message, old.plain_data); +END; + +CREATE TRIGGER IF NOT EXISTS transactions_fts_update AFTER UPDATE ON transactions BEGIN + INSERT INTO transactions_fts (transactions_fts, rowid, message, plain_data) + VALUES ('delete', old.rowid, old.message, old.plain_data); + INSERT INTO transactions_fts (rowid, message, plain_data) + VALUES (new.rowid, new.message, new.plain_data); +END; + +-- Daily exchange rates, keyed by the source's own date so historical rates stay +-- auditable and backfill can use the rate for a transaction's own day. +CREATE TABLE IF NOT EXISTS fx_rates ( + currency TEXT NOT NULL, + as_of TEXT NOT NULL, + rate_to_vnd REAL NOT NULL, + fetched_at TEXT NOT NULL, + PRIMARY KEY (currency, as_of) +); diff --git a/services/fx.ts b/services/fx.ts new file mode 100644 index 0000000..4a0e4ba --- /dev/null +++ b/services/fx.ts @@ -0,0 +1,137 @@ +import { toVndMinorUnits } from '../utils/money'; +import type { Environment } from '../types'; + +export type FxRate = { + currency: string; + as_of: string; + rate_to_vnd: number; +}; + +const RATE_SOURCES = [ + (currency: string) => `https://cdn.jsdelivr.net/npm/@fawazahmed0/currency-api@latest/v1/currencies/${currency}.min.json`, + (currency: string) => `https://latest.currency-api.pages.dev/v1/currencies/${currency}.min.json`, +]; + +/** + * Fetches the current VND rate for a currency. The upstream dataset is updated + * once a day and carries its own `date`, which we keep as `as_of` so stored + * rates stay auditable. + */ +export const fetchRate = async (currency: string): Promise => { + const normalizedCurrency = currency.toLowerCase(); + + for (const buildUrl of RATE_SOURCES) { + try { + const response = await fetch(buildUrl(normalizedCurrency)); + if (!response.ok) continue; + + const data = await response.json() as Record; + const rate = (data[normalizedCurrency] as Record | undefined)?.vnd; + if (!Number.isFinite(rate)) continue; + + return { + currency: currency.toUpperCase(), + as_of: typeof data.date === 'string' ? data.date : new Date().toISOString().slice(0, 10), + rate_to_vnd: rate as number, + }; + } catch (error) { + console.warn(`⚠️ Failed to fetch ${currency.toUpperCase()} rate`, error); + } + } + + return null; +}; + +export const saveRate = async (env: Environment, rate: FxRate) => { + await env.DB.prepare( + `INSERT INTO fx_rates (currency, as_of, rate_to_vnd, fetched_at) + VALUES (?, ?, ?, ?) + ON CONFLICT (currency, as_of) DO UPDATE SET rate_to_vnd = excluded.rate_to_vnd, fetched_at = excluded.fetched_at`, + ).bind(rate.currency, rate.as_of, rate.rate_to_vnd, new Date().toISOString()).run(); +}; + +/** + * Returns the stored rate closest to (and not after) `onDate`, falling back to + * the most recent rate we hold. Returns null when the currency is unknown. + */ +export const lookupRate = async (env: Environment, currency: string, onDate?: string): Promise => { + const normalizedCurrency = currency.toUpperCase(); + const targetDate = onDate || new Date().toISOString().slice(0, 10); + + const onOrBefore = await env.DB.prepare( + `SELECT currency, as_of, rate_to_vnd FROM fx_rates + WHERE currency = ? AND as_of <= ? + ORDER BY as_of DESC LIMIT 1`, + ).bind(normalizedCurrency, targetDate).first(); + if (onOrBefore) return onOrBefore; + + return await env.DB.prepare( + `SELECT currency, as_of, rate_to_vnd FROM fx_rates + WHERE currency = ? + ORDER BY as_of ASC LIMIT 1`, + ).bind(normalizedCurrency).first(); +}; + +/** + * Resolves a rate from the local table, fetching and caching it only when we + * have nothing stored. Never throws: a missing rate leaves the transaction + * unconverted for the backfill job rather than failing ingest. + */ +export const resolveRate = async (env: Environment, currency: string, onDate?: string): Promise => { + try { + const stored = await lookupRate(env, currency, onDate); + if (stored) return stored; + + const fetched = await fetchRate(currency); + if (!fetched) return null; + + await saveRate(env, fetched); + return fetched; + } catch (error) { + console.warn(`⚠️ Could not resolve ${currency} rate`, error); + return null; + } +}; + +/** + * Refreshes stored rates for every foreign currency seen in the ledger, then + * fills in transactions that were recorded while no rate was available. Each + * row is converted with the rate for its own date, not today's. + */ +export const backfillExchangeRates = async (env: Environment) => { + const { results: currencies } = await env.DB.prepare( + `SELECT DISTINCT currency FROM transactions WHERE currency != 'VND' + UNION + SELECT DISTINCT original_currency FROM transactions WHERE original_currency IS NOT NULL AND original_currency != 'VND'`, + ).all<{ currency: string }>(); + + for (const row of currencies) { + if (!row.currency) continue; + const rate = await fetchRate(row.currency); + if (rate) await saveRate(env, rate); + } + + const { results: pending } = await env.DB.prepare( + `SELECT id, occurred_at, amount_minor, currency FROM transactions + WHERE amount_vnd_minor IS NULL AND currency != 'VND'`, + ).all<{ id: string; occurred_at: string; amount_minor: number; currency: string }>(); + + let converted = 0; + for (const transaction of pending) { + const rate = await lookupRate(env, transaction.currency, transaction.occurred_at.slice(0, 10)); + if (!rate) continue; + + await env.DB.prepare( + `UPDATE transactions SET amount_vnd_minor = ?, fx_rate = ?, fx_rate_as_of = ? WHERE id = ?`, + ).bind( + toVndMinorUnits(transaction.amount_minor, transaction.currency, rate.rate_to_vnd), + rate.rate_to_vnd, + rate.as_of, + transaction.id, + ).run(); + converted += 1; + } + + console.info(`💱 Refreshed ${currencies.length} rate(s), backfilled ${converted} transaction(s)`); + return { currencies: currencies.length, converted }; +}; diff --git a/services/telegram.ts b/services/telegram.ts index b16c873..b1550f2 100644 --- a/services/telegram.ts +++ b/services/telegram.ts @@ -1,4 +1,5 @@ import { Telegraf } from 'telegraf'; +import { parseCurrencyAmount } from '../utils/money'; import type { Environment } from '../types'; const formatVietnameseNumber = (value: string) => { @@ -18,29 +19,6 @@ const formatVietnameseNumber = (value: string) => { }).format(parsedValue); }; -const parseCurrencyAmount = (value: string) => { - const normalizedValue = value.trim().replace(/\s/g, ''); - const lastComma = normalizedValue.lastIndexOf(','); - const lastDot = normalizedValue.lastIndexOf('.'); - - if (lastComma > -1 && lastDot > -1) { - const decimalSeparator = lastComma > lastDot ? ',' : '.'; - const thousandsSeparator = decimalSeparator === ',' ? '.' : ','; - return Number(normalizedValue.replace(new RegExp(`\\${thousandsSeparator}`, 'g'), '').replace(decimalSeparator, '.')); - } - - const separator = lastComma > -1 ? ',' : lastDot > -1 ? '.' : ''; - if (!separator) return Number(normalizedValue); - - const separatorIndex = normalizedValue.lastIndexOf(separator); - const digitsAfterSeparator = normalizedValue.length - separatorIndex - 1; - const isDecimalSeparator = digitsAfterSeparator > 0 && digitsAfterSeparator <= 2; - - return Number(isDecimalSeparator - ? normalizedValue.replace(separator, '.') - : normalizedValue.replace(new RegExp(`\\${separator}`, 'g'), '')); -}; - const formatDong = (value: number) => `${new Intl.NumberFormat('vi-VN', { maximumFractionDigits: 0 }).format(Math.round(value))}đ`; const getExchangeRateToVnd = async (currency: string) => { diff --git a/services/transactions-store.ts b/services/transactions-store.ts new file mode 100644 index 0000000..a7590b6 --- /dev/null +++ b/services/transactions-store.ts @@ -0,0 +1,162 @@ +import { resolveRate } from './fx'; +import { normalizeCurrency, toMinorUnits, toVndMinorUnits } from '../utils/money'; +import type { Environment, TransactionDetails, TransactionRow } from '../types'; + +const DEFAULT_DUPLICATE_WINDOW_MINUTES = 15; + +/** + * Parses the extractor's `dd/MM/yyyy hh:mm:ss` datetime (local Asia/Bangkok + * time) into a sortable ISO8601 UTC string. Falls back to now when the value is + * missing or unparseable, so a bad date never blocks recording a transaction. + */ +export const parseOccurredAt = (value?: string | null, now: Date = new Date()) => { + const raw = (value || '').trim(); + if (!raw) return now.toISOString(); + + const local = raw.match(/^(\d{1,2})\/(\d{1,2})\/(\d{4})(?:[\sT]+(\d{1,2}):(\d{2})(?::(\d{2}))?)?/); + if (local) { + const [, day, month, year, hour = '0', minute = '0', second = '0'] = local; + // Asia/Bangkok is UTC+7 year-round, so a fixed offset is exact here. + const utc = Date.UTC(Number(year), Number(month) - 1, Number(day), Number(hour) - 7, Number(minute), Number(second)); + if (Number.isFinite(utc)) return new Date(utc).toISOString(); + } + + const parsed = new Date(raw); + return Number.isNaN(parsed.getTime()) ? now.toISOString() : parsed.toISOString(); +}; + +const normalizeDirection = (value?: string | null) => (String(value || '').toLowerCase() === 'credit' ? 'credit' : 'debit'); + +const normalizeSourceKind = (value?: string | null) => (String(value || '').toLowerCase() === 'merchant' ? 'merchant' : 'bank'); + +/** Every (currency, minor amount) pair a transaction can be recognised by. */ +const amountKeys = (row: Pick) => { + const keys = new Set([`${row.currency}:${row.amount_minor}`]); + if (row.original_currency && row.original_amount_minor !== null && row.original_amount_minor !== undefined) { + keys.add(`${row.original_currency}:${row.original_amount_minor}`); + } + return keys; +}; + +/** Maps extractor output onto a table row. Returns null when there is no usable amount. */ +export const buildTransactionRow = (details: TransactionDetails, source: 'email' | 'manual' | 'ocr'): TransactionRow | null => { + const currency = normalizeCurrency(details.currency) || 'VND'; + const amountMinor = toMinorUnits(details.amount, currency); + if (amountMinor === null) return null; + + const originalCurrency = normalizeCurrency(details.original_currency); + const originalAmountMinor = originalCurrency ? toMinorUnits(details.original_amount, originalCurrency) : null; + + return { + id: crypto.randomUUID(), + occurred_at: parseOccurredAt(details.datetime), + amount_minor: amountMinor, + currency, + original_amount_minor: originalAmountMinor, + original_currency: originalAmountMinor === null ? null : originalCurrency, + amount_vnd_minor: currency === 'VND' ? amountMinor : null, + fx_rate: null, + fx_rate_as_of: null, + bank_name: (details.bank_name || '').trim(), + category: details.category?.trim() || null, + direction: normalizeDirection(details.direction), + source, + source_kind: normalizeSourceKind(details.source_kind), + message: details.message || '', + plain_data: details.plain_data || details.message || '', + duplicate_of: null, + needs_review: 0, + created_at: new Date().toISOString(), + }; +}; + +/** + * Looks for an existing transaction that is the same real-world payment seen + * from the other side — a merchant receipt and the matching bank debit, which + * arrive minutes apart under different senders. + * + * A debit paired with a credit is a transfer between the user's own accounts, + * not a duplicate, so direction must match. + */ +export const findCrossSourceDuplicate = async (env: Environment, row: TransactionRow, windowMinutes: number) => { + const occurredAt = new Date(row.occurred_at).getTime(); + const from = new Date(occurredAt - windowMinutes * 60_000).toISOString(); + const to = new Date(occurredAt + windowMinutes * 60_000).toISOString(); + + const { results } = await env.DB.prepare( + `SELECT id, currency, amount_minor, original_currency, original_amount_minor + FROM transactions + WHERE occurred_at BETWEEN ? AND ? + AND direction = ? + AND source_kind != ? + AND duplicate_of IS NULL + ORDER BY occurred_at ASC`, + ).bind(from, to, row.direction, row.source_kind).all(); + + const incoming = amountKeys(row); + const settledKey = `${row.currency}:${row.amount_minor}`; + + for (const candidate of results) { + const candidateKeys = amountKeys(candidate); + const overlap = [...incoming].filter((key) => candidateKeys.has(key)); + if (overlap.length === 0) continue; + + // Matching only on the pre-conversion amount is the weaker signal, so + // flag it for review rather than silently folding the rows together. + const matchedOnSettledAmount = overlap.includes(settledKey) && candidate.currency === row.currency; + return { id: candidate.id, needsReview: matchedOnSettledAmount ? 0 : 1 }; + } + + return null; +}; + +/** + * Writes a transaction to D1. Exact repeats (webhook or delivery retries) are + * dropped by the UNIQUE constraint; the same payment arriving from a second + * sender is kept but linked, so totals count it once without losing detail. + */ +export const saveTransaction = async (env: Environment, details: TransactionDetails, source: 'email' | 'manual' | 'ocr' = 'email') => { + const row = buildTransactionRow(details, source); + if (!row) { + console.warn('🗄️ Skipping D1 write: no usable amount in extracted transaction'); + return null; + } + + if (row.currency !== 'VND') { + const rate = await resolveRate(env, row.currency, row.occurred_at.slice(0, 10)); + if (rate) { + row.amount_vnd_minor = toVndMinorUnits(row.amount_minor, row.currency, rate.rate_to_vnd); + row.fx_rate = rate.rate_to_vnd; + row.fx_rate_as_of = rate.as_of; + } + } + + const windowMinutes = Number(env.TRANSACTION_DUPLICATE_WINDOW_MINUTES) || DEFAULT_DUPLICATE_WINDOW_MINUTES; + const duplicate = await findCrossSourceDuplicate(env, row, windowMinutes); + if (duplicate) { + row.duplicate_of = duplicate.id; + row.needs_review = duplicate.needsReview; + console.info(`🗄️ Linked transaction to existing ${duplicate.id} (needs_review=${duplicate.needsReview})`); + } + + const result = await env.DB.prepare( + `INSERT INTO transactions ( + id, occurred_at, amount_minor, currency, original_amount_minor, original_currency, + amount_vnd_minor, fx_rate, fx_rate_as_of, bank_name, category, direction, + source, source_kind, message, plain_data, duplicate_of, needs_review, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (occurred_at, amount_minor, currency, bank_name) DO NOTHING`, + ).bind( + row.id, row.occurred_at, row.amount_minor, row.currency, row.original_amount_minor, row.original_currency, + row.amount_vnd_minor, row.fx_rate, row.fx_rate_as_of, row.bank_name, row.category, row.direction, + row.source, row.source_kind, row.message, row.plain_data, row.duplicate_of, row.needs_review, row.created_at, + ).run(); + + if (!result.meta?.changes) { + console.info('🗄️ Transaction already recorded, skipping duplicate insert'); + return null; + } + + console.info(`🗄️ Stored transaction ${row.id} in D1`); + return row; +}; diff --git a/tests/transactions-store.test.ts b/tests/transactions-store.test.ts new file mode 100644 index 0000000..a41d017 --- /dev/null +++ b/tests/transactions-store.test.ts @@ -0,0 +1,190 @@ +import { describe, expect, it } from 'bun:test'; + +const { currencyDecimals, fromMinorUnits, normalizeCurrency, toMinorUnits, toVndMinorUnits } = await import('../utils/money'); +const { buildTransactionRow, findCrossSourceDuplicate, parseOccurredAt, saveTransaction } = await import('../services/transactions-store'); +const { storeTransactionRecord } = await import('../handlers/transactions'); + +type Stub = { sql: RegExp; rows?: any[]; first?: any; changes?: number }; + +/** Minimal D1 stand-in: matches queued stubs against the SQL text. */ +const makeDb = (stubs: Stub[]) => { + const statements: { sql: string; args: unknown[] }[] = []; + + const db = { + statements, + prepare(sql: string) { + return { + bind: (...args: unknown[]) => { + statements.push({ sql, args }); + const stub = stubs.find((candidate) => candidate.sql.test(sql)); + return { + all: async () => ({ results: stub?.rows ?? [] }), + first: async () => stub?.first ?? null, + run: async () => ({ meta: { changes: stub?.changes ?? 1 } }), + }; + }, + }; + }, + }; + + return db as any; +}; + +const details = { + bank_name: 'VCB', + datetime: '21/09/2026 14:30:00', + amount: '120.000', + currency: 'VNĐ', + message: 'Ăn trưa 120.000 VNĐ', + plain_data: 'Thanh toán ăn trưa tại quán gần văn phòng', +}; + +describe('money helpers', () => { + it('treats VND as a zero-decimal currency', () => { + expect(currencyDecimals('VND')).toBe(0); + expect(currencyDecimals('VNĐ')).toBe(0); + expect(currencyDecimals('USD')).toBe(2); + }); + + it('normalizes the Vietnamese currency spelling', () => { + expect(normalizeCurrency('vnđ')).toBe('VND'); + expect(normalizeCurrency(' usd ')).toBe('USD'); + expect(normalizeCurrency(null)).toBe(''); + }); + + it('converts human-formatted amounts into minor units', () => { + expect(toMinorUnits('4.320.000', 'VND')).toBe(4320000); + expect(toMinorUnits('8,99', 'USD')).toBe(899); + expect(toMinorUnits('1,234.56', 'USD')).toBe(123456); + expect(toMinorUnits('', 'USD')).toBeNull(); + expect(toMinorUnits('abc', 'USD')).toBeNull(); + }); + + it('round-trips minor units and converts to đồng', () => { + expect(fromMinorUnits(899, 'USD')).toBeCloseTo(8.99); + expect(toVndMinorUnits(10000, 'USD', 25400)).toBe(2540000); + }); +}); + +describe('parseOccurredAt', () => { + it('reads dd/MM/yyyy hh:mm:ss as Asia/Bangkok time', () => { + expect(parseOccurredAt('21/09/2026 14:30:00')).toBe('2026-09-21T07:30:00.000Z'); + }); + + it('accepts a date without a time', () => { + expect(parseOccurredAt('01/02/2026')).toBe('2026-01-31T17:00:00.000Z'); + }); + + it('falls back to now when the value is missing or unparseable', () => { + const now = new Date('2026-09-21T00:00:00.000Z'); + expect(parseOccurredAt('', now)).toBe(now.toISOString()); + expect(parseOccurredAt('not a date', now)).toBe(now.toISOString()); + }); +}); + +describe('buildTransactionRow', () => { + it('maps extracted details onto a row', () => { + const row = buildTransactionRow(details, 'email')!; + expect(row.currency).toBe('VND'); + expect(row.amount_minor).toBe(120000); + expect(row.amount_vnd_minor).toBe(120000); + expect(row.direction).toBe('debit'); + expect(row.source_kind).toBe('bank'); + expect(row.bank_name).toBe('VCB'); + expect(row.duplicate_of).toBeNull(); + }); + + it('keeps the merchant-side amount when the bank settled in another currency', () => { + const row = buildTransactionRow({ + ...details, + amount: '2.540.000', + currency: 'VNĐ', + original_amount: '100.00', + original_currency: 'USD', + }, 'email')!; + + expect(row.amount_minor).toBe(2540000); + expect(row.original_amount_minor).toBe(10000); + expect(row.original_currency).toBe('USD'); + }); + + it('returns null when there is no usable amount', () => { + expect(buildTransactionRow({ ...details, amount: undefined }, 'email')).toBeNull(); + }); + + it('defaults an empty bank name to the empty string so dedup still applies', () => { + const row = buildTransactionRow({ ...details, bank_name: undefined }, 'manual')!; + expect(row.bank_name).toBe(''); + }); +}); + +describe('findCrossSourceDuplicate', () => { + const row = buildTransactionRow({ ...details, source_kind: 'merchant' }, 'email')!; + + it('links a merchant receipt to the matching bank debit', async () => { + const db = makeDb([{ sql: /FROM transactions/, rows: [{ id: 'bank-row', currency: 'VND', amount_minor: 120000, original_currency: null, original_amount_minor: null }] }]); + expect(await findCrossSourceDuplicate({ DB: db } as any, row, 15)).toEqual({ id: 'bank-row', needsReview: 0 }); + }); + + it('flags a match made only on the pre-conversion amount', async () => { + const converted = buildTransactionRow({ + ...details, + amount: '2.540.000', + currency: 'VNĐ', + original_amount: '100.00', + original_currency: 'USD', + source_kind: 'merchant', + }, 'email')!; + const db = makeDb([{ sql: /FROM transactions/, rows: [{ id: 'bank-row', currency: 'USD', amount_minor: 10000, original_currency: null, original_amount_minor: null }] }]); + + expect(await findCrossSourceDuplicate({ DB: db } as any, converted, 15)).toEqual({ id: 'bank-row', needsReview: 1 }); + }); + + it('returns null when no candidate amount matches', async () => { + const db = makeDb([{ sql: /FROM transactions/, rows: [{ id: 'other', currency: 'VND', amount_minor: 999, original_currency: null, original_amount_minor: null }] }]); + expect(await findCrossSourceDuplicate({ DB: db } as any, row, 15)).toBeNull(); + }); + + it('only considers the opposite source kind and the same direction', async () => { + const db = makeDb([{ sql: /FROM transactions/, rows: [] }]); + await findCrossSourceDuplicate({ DB: db } as any, row, 15); + + const [query] = db.statements; + expect(query.sql).toContain('source_kind !='); + expect(query.sql).toContain('direction ='); + expect(query.args).toContain('debit'); + expect(query.args).toContain('merchant'); + }); +}); + +describe('saveTransaction', () => { + it('stores a transaction and reports the inserted row', async () => { + const db = makeDb([{ sql: /SELECT id, currency/, rows: [] }, { sql: /INSERT INTO transactions/, changes: 1 }]); + const saved = await saveTransaction({ DB: db } as any, details, 'email'); + + expect(saved?.amount_minor).toBe(120000); + expect(db.statements.some(({ sql }) => sql.includes('ON CONFLICT'))).toBe(true); + }); + + it('reports nothing when the unique constraint drops an exact repeat', async () => { + const db = makeDb([{ sql: /SELECT id, currency/, rows: [] }, { sql: /INSERT INTO transactions/, changes: 0 }]); + expect(await saveTransaction({ DB: db } as any, details, 'email')).toBeNull(); + }); + + it('skips the write when the amount could not be extracted', async () => { + const db = makeDb([]); + expect(await saveTransaction({ DB: db } as any, { ...details, amount: undefined }, 'email')).toBeNull(); + expect(db.statements).toHaveLength(0); + }); +}); + +describe('storeTransactionRecord', () => { + it('skips quietly when no D1 binding is configured', async () => { + expect(await storeTransactionRecord(details, {} as any, 'email')).toBeNull(); + }); + + it('swallows D1 failures so the existing flow is unaffected', async () => { + const db = { prepare: () => { throw new Error('D1 unavailable'); } }; + expect(await storeTransactionRecord(details, { DB: db } as any, 'email')).toBeNull(); + }); +}); diff --git a/types/index.ts b/types/index.ts index 143fb20..40445d3 100644 --- a/types/index.ts +++ b/types/index.ts @@ -9,4 +9,48 @@ export type Environment = Env & { readonly OPENAI_API_KEY: string; readonly OPENAI_ASSISTANT_VECTORSTORE_ID: string; + + /** Minutes either side of a transaction to look for the same payment from another sender. */ + readonly TRANSACTION_DUPLICATE_WINDOW_MINUTES?: string; + + readonly DB: D1Database; +}; + +/** Shape returned by the extraction prompt. */ +export type TransactionDetails = { + bank_name?: string; + datetime?: string; + amount?: string; + currency?: string; + original_amount?: string; + original_currency?: string; + category?: string; + direction?: string; + source_kind?: string; + message?: string; + plain_data?: string; + result?: string; + error?: string; +}; + +export type TransactionRow = { + id: string; + occurred_at: string; + amount_minor: number; + currency: string; + original_amount_minor: number | null; + original_currency: string | null; + amount_vnd_minor: number | null; + fx_rate: number | null; + fx_rate_as_of: string | null; + bank_name: string; + category: string | null; + direction: 'debit' | 'credit'; + source: 'email' | 'manual' | 'ocr'; + source_kind: 'bank' | 'merchant'; + message: string; + plain_data: string; + duplicate_of: string | null; + needs_review: number; + created_at: string; }; diff --git a/utils/money.ts b/utils/money.ts new file mode 100644 index 0000000..b957921 --- /dev/null +++ b/utils/money.ts @@ -0,0 +1,61 @@ +/** + * Currencies whose smallest unit equals the major unit (no decimal subunit). + * Everything else is assumed to use two decimal places. + */ +const ZERO_DECIMAL_CURRENCIES = new Set(['VND', 'VNĐ', 'JPY', 'KRW', 'IDR', 'CLP', 'ISK', 'PYG', 'RWF', 'UGX', 'VUV', 'XAF', 'XOF', 'XPF']); + +export const normalizeCurrency = (currency?: string | null) => { + const normalized = (currency || '').trim().toUpperCase(); + if (!normalized) return ''; + return normalized === 'VNĐ' ? 'VND' : normalized; +}; + +export const currencyDecimals = (currency?: string | null) => + ZERO_DECIMAL_CURRENCIES.has(normalizeCurrency(currency)) ? 0 : 2; + +/** + * Parses a human-formatted amount into a Number, handling both Vietnamese + * (1.234.567,89) and Anglo (1,234,567.89) separator conventions. + */ +export const parseCurrencyAmount = (value: string) => { + const normalizedValue = value.trim().replace(/\s/g, ''); + const lastComma = normalizedValue.lastIndexOf(','); + const lastDot = normalizedValue.lastIndexOf('.'); + + if (lastComma > -1 && lastDot > -1) { + const decimalSeparator = lastComma > lastDot ? ',' : '.'; + const thousandsSeparator = decimalSeparator === ',' ? '.' : ','; + return Number(normalizedValue.replace(new RegExp(`\\${thousandsSeparator}`, 'g'), '').replace(decimalSeparator, '.')); + } + + const separator = lastComma > -1 ? ',' : lastDot > -1 ? '.' : ''; + if (!separator) return Number(normalizedValue); + + const separatorIndex = normalizedValue.lastIndexOf(separator); + const digitsAfterSeparator = normalizedValue.length - separatorIndex - 1; + const isDecimalSeparator = digitsAfterSeparator > 0 && digitsAfterSeparator <= 2; + + return Number(isDecimalSeparator + ? normalizedValue.replace(separator, '.') + : normalizedValue.replace(new RegExp(`\\${separator}`, 'g'), '')); +}; + +/** Converts a human-formatted amount into integer minor units for the currency. */ +export const toMinorUnits = (value: string | number | null | undefined, currency?: string | null) => { + if (value === null || value === undefined || value === '') return null; + + const parsed = typeof value === 'number' ? value : parseCurrencyAmount(String(value)); + if (!Number.isFinite(parsed)) return null; + + return Math.round(parsed * 10 ** currencyDecimals(currency)); +}; + +/** Inverse of `toMinorUnits`, returning the major-unit Number. */ +export const fromMinorUnits = (minor: number, currency?: string | null) => minor / 10 ** currencyDecimals(currency); + +/** + * Converts an amount in minor units to VND minor units (VND has no subunit, so + * minor units are whole đồng) using a major-unit exchange rate. + */ +export const toVndMinorUnits = (minor: number, currency: string, rateToVnd: number) => + Math.round(fromMinorUnits(minor, currency) * rateToVnd); diff --git a/wrangler.jsonc b/wrangler.jsonc index da8a5fe..32fc490 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -5,8 +5,16 @@ "compatibility_date": "2026-06-16", "compatibility_flags": ["nodejs_compat"], "triggers": { - "crons": ["0 15 * * *", "58 16 * * 1", "0 15 1 * *"] + "crons": ["0 15 * * *", "58 16 * * 1", "0 15 1 * *", "0 18 * * *"] }, + "d1_databases": [ + { + "binding": "DB", + "database_name": "personalaccountant", + "database_id": "", + "migrations_dir": "migrations" + } + ], "observability": { "enabled": true, "head_sampling_rate": 1 @@ -17,8 +25,8 @@ "OPENAI_ASSISTANT_MODEL": "gpt-5.6-luna", "OPENAI_ASSISTANT_ROUTER_MODEL": "gpt-5.6-luna", "OPENAI_ASSISTANT_RESPONSE_FORMAT_INSTRUCTIONS": "Format answers for Telegram MarkdownV2. Use *single asterisks* for bold labels/headings; do not use double-asterisk Markdown because Telegram MarkdownV2 bold uses single asterisks. When listing multiple transactions, use short bullet points. Format all money amounts in Vietnamese number style: use dots for thousands and commas for decimal fractions; remove insignificant trailing decimal zeros for non-VND currencies (for example, write 16 AUD instead of 16.000 AUD, and 5,24 AUD instead of 5.240 AUD). When showing transaction times, include only hour and minute (HH:mm); do not include seconds. When the user does not specify a limit or count, assume they want all matching transactions. If all matching transactions would be too large or token-expensive to list, return the 20 transactions closest to the requested time period and clearly say that the list was limited to 20. When the answer includes multiple dates, split the response into separate bold date sections. When there are 3 or more transactions, include a bold total summary at the bottom. When there are 10 or more transactions, do not use Markdown tables because Telegram MarkdownV2 does not support them reliably; use grouped date sections with bullet points instead.", - "OPENAI_PROCESS_EMAIL_SYSTEM_PROMPT": "You act like an personal accountant, extracts email transaction details and summary them in JSON plain format, not markdown, not HTML, just plain JSON like below. No yapping.\n\n{\n\"bank_name\": \"...\",\n\"datetime\": \"...\",\n\"amount\": \"...\",\n\"currency\": \"...\",\n\"message: \"...\",\n\"plain_data\": \"...\"\n}\n\nIf it's not a transaction, just return { \"result\": \"failed\" }.", - "OPENAI_PROCESS_EMAIL_USER_PROMPT": "No yapping. Extract the transaction details such as:\n\nbank name:\namount: human-readable, example: 1.000, 4.320.000\ndatetime: dd/MM/yyyy hh:mm:ss\ncurrency: try to stick with three letter code. If VND, must change to VNĐ\nmessage: summary, classify the transaction detail in one paragraph (30-50 words or lesser, non-formal). Focus on purchase order, product name or person who I sent the money to; Grab/Uber group orders; sometimes I sent money between my bank account, you can know by lookup my name (DUONG ANH TUAN, or Đường Anh Tuấn), and no need to repeat my name in the message; or they are credit card transaction. In Vietnamese only. DO NOT REMOVE amount, currency in message.\nplain_data: collect, classify the transaction detail in around 200-400 words. Ensure you collect and store enough information of orders (example if food order: pay method (momo, credit card, COD); food store name; number, list of food I ordered; total discount) so we can query later. Just plain information, do not include any advertisement, do not include any offers advertisement (like group order), sales in the email.\n\nFrom the following email content:", + "OPENAI_PROCESS_EMAIL_SYSTEM_PROMPT": "You act like an personal accountant, extracts email transaction details and summary them in JSON plain format, not markdown, not HTML, just plain JSON like below. No yapping.\n\n{\n\"bank_name\": \"...\",\n\"datetime\": \"...\",\n\"amount\": \"...\",\n\"currency\": \"...\",\n\"message: \"...\",\n\"plain_data\": \"...\",\n\"category\": \"...\",\n\"direction\": \"debit|credit\",\n\"source_kind\": \"bank|merchant\",\n\"original_amount\": \"...\",\n\"original_currency\": \"...\"\n}\n\nIf it's not a transaction, just return { \"result\": \"failed\" }.", + "OPENAI_PROCESS_EMAIL_USER_PROMPT": "No yapping. Extract the transaction details such as:\n\nbank name:\namount: human-readable, example: 1.000, 4.320.000\ndatetime: dd/MM/yyyy hh:mm:ss\ncurrency: try to stick with three letter code. If VND, must change to VNĐ\nmessage: summary, classify the transaction detail in one paragraph (30-50 words or lesser, non-formal). Focus on purchase order, product name or person who I sent the money to; Grab/Uber group orders; sometimes I sent money between my bank account, you can know by lookup my name (DUONG ANH TUAN, or Đường Anh Tuấn), and no need to repeat my name in the message; or they are credit card transaction. In Vietnamese only. DO NOT REMOVE amount, currency in message.\ncategory: one lowercase slug describing what was bought, optionally with a sub-category separated by a slash, for example food/delivery, transport/ride, entertainment/cinema, shopping/online, bills/utilities, cloud/hosting. Reuse the same slug for similar purchases so they can be grouped later.\ndirection: debit when money left my account, credit when money arrived. Default to debit.\nsource_kind: bank when the email comes from a bank or card issuer, merchant when it comes from a shop or service such as a ride hailing app, marketplace or cloud provider.\noriginal_amount, original_currency: only when the merchant charged in one currency and the bank settled in another. Put the bank settled amount in amount/currency, and the merchant original amount here. Leave both empty otherwise.\nplain_data: collect, classify the transaction detail in around 200-400 words. Ensure you collect and store enough information of orders (example if food order: pay method (momo, credit card, COD); food store name; number, list of food I ordered; total discount) so we can query later. Just plain information, do not include any advertisement, do not include any offers advertisement (like group order), sales in the email.\n\nFrom the following email content:", "OPENAI_ASSISTANT_SCHEDULED_PROMPT": "Tổng hợp ngắn, phân loại các giao dịch được tạo trong ngày hôm nay %DATETIME%" } }