Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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

Expand Down
6 changes: 3 additions & 3 deletions handlers/assistant.ts
Original file line number Diff line number Diff line change
@@ -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';

Expand Down Expand Up @@ -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';
};

Expand All @@ -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';
};

Expand Down
31 changes: 29 additions & 2 deletions handlers/transactions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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}`);
Expand Down Expand Up @@ -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);
Expand All @@ -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";
};
5 changes: 5 additions & 0 deletions index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
}
},

Expand Down
68 changes: 68 additions & 0 deletions migrations/0001_create_transactions.sql
Original file line number Diff line number Diff line change
@@ -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)
);
137 changes: 137 additions & 0 deletions services/fx.ts
Original file line number Diff line number Diff line change
@@ -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<FxRate | null> => {
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<string, unknown>;
const rate = (data[normalizedCurrency] as Record<string, number> | 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<FxRate | null> => {
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<FxRate>();
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<FxRate>();
};

/**
* 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<FxRate | null> => {
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 };
};
24 changes: 1 addition & 23 deletions services/telegram.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { Telegraf } from 'telegraf';
import { parseCurrencyAmount } from '../utils/money';
import type { Environment } from '../types';

const formatVietnameseNumber = (value: string) => {
Expand All @@ -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) => {
Expand Down
Loading
Loading