diff --git a/deno.json b/deno.json index a41dd33..69c23e2 100644 --- a/deno.json +++ b/deno.json @@ -1,7 +1,7 @@ { "nodeModulesDir": "auto", "tasks": { - "dev": "deno run --env-file -A npm:vite dev", + "dev": "deno run --env-file --unstable-cron -A npm:vite dev", "build": "deno run -A npm:vite build", "start": "deno run --env-file --unstable-cron -A build/index.js", "check": "deno run -A npm:@sveltejs/kit/svelte-kit sync && deno run -A npm:svelte-check --tsconfig ./tsconfig.json", diff --git a/migrations/001_init.sql b/migrations/001_init.sql index 8b1e674..af3bbdc 100644 --- a/migrations/001_init.sql +++ b/migrations/001_init.sql @@ -75,6 +75,7 @@ CREATE TABLE transactions ( extra TEXT, -- verbatim SimpleFIN extra JSON category_id INTEGER REFERENCES categories (id), -- cache of latest categorization event created_at TEXT NOT NULL, + removed_at TEXT, -- soft delete (stale pending rows); events are never deleted UNIQUE (account_id, sfin_id) ); CREATE INDEX idx_transactions_posted ON transactions (posted); diff --git a/openspec/changes/bootstrap-finance-app/tasks.md b/openspec/changes/bootstrap-finance-app/tasks.md index fa5c66a..faf41fe 100644 --- a/openspec/changes/bootstrap-finance-app/tasks.md +++ b/openspec/changes/bootstrap-finance-app/tasks.md @@ -1,49 +1,49 @@ -## 1. Foundation - -- [x] 1.1 Scaffold SvelteKit project running under Deno 2.9 (deno.json tasks for dev/build/start, adapter choice per design D1) with a health-check route -- [x] 1.2 Implement SQLite bootstrap via `node:sqlite`: open `DB_PATH`, enable WAL, numbered-migration runner with `schema_version` table applied at startup -- [x] 1.3 Write migration 001: users, sessions, oauth state/session stores, connections, accounts, transactions, balance_snapshots, raw_syncs, categories (seed built-in Transfer), rules, categorization_events -- [x] 1.4 Add config module reading and validating `APP_URL`, `ALLOWED_DIDS`, `DB_PATH`, and OAuth signing key; fail fast with clear errors on missing config -- [x] 1.5 Establish the service-layer convention (design D6): domain operations as transport-agnostic modules under `src/lib/server/`, SvelteKit loads/actions as thin adapters only — no SQL or domain logic in routes - -## 2. Authentication (specs/auth) - -- [x] 2.1 Dynamic routes for `/client-metadata.json` and `/jwks.json` generated from `APP_URL` and the signing key -- [x] 2.2 Integrate `@atproto/oauth-client-node` with SQLite-backed state/session stores; login page with handle input; `/oauth/callback` handler (validates the library works under Deno npm-compat — fallback per design risk if not) -- [x] 2.3 DID allowlist check on callback: create/update user record for allowlisted DIDs, 403 otherwise -- [x] 2.4 Application sessions: HTTP-only Secure cookie, SQLite sessions table, hooks guard on all protected routes, logout -- [ ] 2.5 Verify full OAuth round-trip end-to-end through the public `APP_URL` origin (both household DIDs) - -## 3. SimpleFIN ingestion (specs/simplefin-sync) - -- [ ] 3.1 Settings page: setup-token paste → decode → claim → store connection; handle already-claimed 403 with clear message; trigger initial sync -- [ ] 3.2 Sync engine: fetch `/accounts` (pending included, overlapping start-date), archive verbatim payload to raw_syncs (success and failure rows) before any processing -- [ ] 3.3 Idempotent normalization: upsert accounts and transactions (integer cents, verbatim extra JSON) keyed on SimpleFIN ids; pure function over payload with fixture-based tests including double-run idempotency -- [ ] 3.4 Balance snapshots: insert one row per account per successful sync -- [ ] 3.5 Pending→posted reconciliation: conservative matcher (account + exact amount + date window), carry categorization forward via `reconciliation` event, remove stale pending rows; tests for the categorized-pending-reposts case -- [ ] 3.6 Schedule daily sync with `Deno.cron` and add manual "sync now" action; record outcomes and update `last_successful_data_at` per account - -## 4. Account lifecycle (specs/account-management) - -- [ ] 4.1 Discovery: register unknown account ids from sync as `NEW`; auto-transition vanished accounts to `INACTIVE` and reappearing ones back -- [ ] 4.2 Classification UI: dashboard prompt for `NEW` accounts, assign type + optional display name → `ACTIVE`; hide/unhide action -- [ ] 4.3 Connection health: parse sync-response errors, dashboard banners naming the institution with SimpleFIN Bridge link-out, staleness indicator from `last_successful_data_at` - -## 5. Categorization (specs/categorization) - -- [ ] 5.1 Category management UI/API: create, rename, deactivate; enforce non-deletable built-in Transfer -- [ ] 5.2 Event log core: append event + update denormalized `transactions.category_id` in one DB transaction; manual events record actor DID -- [ ] 5.3 Rules engine: exact/contains case-insensitive matching on raw description, deterministic precedence (exact > contains, longer > shorter, newer > older), fire only on uncategorized transactions during sync; unit tests for precedence and the manual-outranks-rule invariant -- [ ] 5.4 Rules UI: create/deactivate rules, retroactive-apply offer scoped to currently-uncategorized matches with result count -- [ ] 5.5 Provenance UI: per-transaction history view showing every event (rule pattern / person / reconciliation, category, timestamp) - -## 6. Reporting (specs/reporting) - -- [ ] 6.1 Transaction ledger: filters (account, category incl. uncategorized, month, pending), inline manual categorization, provenance indicator per row -- [ ] 6.2 Monthly report: income and expense totals per category (posted, non-hidden, transfer-excluded), net figure, Uncategorized line linking to filtered ledger -- [ ] 6.3 Net worth over time: latest-snapshot-per-account-per-day series with asset/liability split, excluding hidden accounts - -## 7. Deployment & verification - -- [ ] 7.1 Production build + run task; README covering Caddy route, env setup, key generation, backup expectations (SQLite file is secret-grade), and known limits (~90-day backfill, APP_URL change forces re-consent) -- [ ] 7.2 End-to-end walkthrough on real infrastructure: both users log in, claim real setup token, first sync lands, classify accounts, create rules, categorize manually, verify monthly report and net worth chart +## 1. Foundation + +- [x] 1.1 Scaffold SvelteKit project running under Deno 2.9 (deno.json tasks for dev/build/start, adapter choice per design D1) with a health-check route +- [x] 1.2 Implement SQLite bootstrap via `node:sqlite`: open `DB_PATH`, enable WAL, numbered-migration runner with `schema_version` table applied at startup +- [x] 1.3 Write migration 001: users, sessions, oauth state/session stores, connections, accounts, transactions, balance_snapshots, raw_syncs, categories (seed built-in Transfer), rules, categorization_events +- [x] 1.4 Add config module reading and validating `APP_URL`, `ALLOWED_DIDS`, `DB_PATH`, and OAuth signing key; fail fast with clear errors on missing config +- [x] 1.5 Establish the service-layer convention (design D6): domain operations as transport-agnostic modules under `src/lib/server/`, SvelteKit loads/actions as thin adapters only — no SQL or domain logic in routes + +## 2. Authentication (specs/auth) + +- [x] 2.1 Dynamic routes for `/client-metadata.json` and `/jwks.json` generated from `APP_URL` and the signing key +- [x] 2.2 Integrate `@atproto/oauth-client-node` with SQLite-backed state/session stores; login page with handle input; `/oauth/callback` handler (validates the library works under Deno npm-compat — fallback per design risk if not) +- [x] 2.3 DID allowlist check on callback: create/update user record for allowlisted DIDs, 403 otherwise +- [x] 2.4 Application sessions: HTTP-only Secure cookie, SQLite sessions table, hooks guard on all protected routes, logout +- [ ] 2.5 Verify full OAuth round-trip end-to-end through the public `APP_URL` origin (both household DIDs) + +## 3. SimpleFIN ingestion (specs/simplefin-sync) + +- [x] 3.1 Settings page: setup-token paste → decode → claim → store connection; handle already-claimed 403 with clear message; trigger initial sync +- [x] 3.2 Sync engine: fetch `/accounts` (pending included, overlapping start-date), archive verbatim payload to raw_syncs (success and failure rows) before any processing +- [x] 3.3 Idempotent normalization: upsert accounts and transactions (integer cents, verbatim extra JSON) keyed on SimpleFIN ids; pure function over payload with fixture-based tests including double-run idempotency +- [x] 3.4 Balance snapshots: insert one row per account per successful sync +- [x] 3.5 Pending→posted reconciliation: conservative matcher (account + exact amount + date window), carry categorization forward via `reconciliation` event, remove stale pending rows; tests for the categorized-pending-reposts case +- [x] 3.6 Schedule daily sync with `Deno.cron` and add manual "sync now" action; record outcomes and update `last_successful_data_at` per account + +## 4. Account lifecycle (specs/account-management) + +- [x] 4.1 Discovery: register unknown account ids from sync as `NEW`; auto-transition vanished accounts to `INACTIVE` and reappearing ones back +- [ ] 4.2 Classification UI: dashboard prompt for `NEW` accounts, assign type + optional display name → `ACTIVE`; hide/unhide action +- [ ] 4.3 Connection health: parse sync-response errors, dashboard banners naming the institution with SimpleFIN Bridge link-out, staleness indicator from `last_successful_data_at` + +## 5. Categorization (specs/categorization) + +- [ ] 5.1 Category management UI/API: create, rename, deactivate; enforce non-deletable built-in Transfer +- [x] 5.2 Event log core: append event + update denormalized `transactions.category_id` in one DB transaction; manual events record actor DID +- [x] 5.3 Rules engine: exact/contains case-insensitive matching on raw description, deterministic precedence (exact > contains, longer > shorter, newer > older), fire only on uncategorized transactions during sync; unit tests for precedence and the manual-outranks-rule invariant +- [ ] 5.4 Rules UI: create/deactivate rules, retroactive-apply offer scoped to currently-uncategorized matches with result count +- [ ] 5.5 Provenance UI: per-transaction history view showing every event (rule pattern / person / reconciliation, category, timestamp) + +## 6. Reporting (specs/reporting) + +- [ ] 6.1 Transaction ledger: filters (account, category incl. uncategorized, month, pending), inline manual categorization, provenance indicator per row +- [ ] 6.2 Monthly report: income and expense totals per category (posted, non-hidden, transfer-excluded), net figure, Uncategorized line linking to filtered ledger +- [ ] 6.3 Net worth over time: latest-snapshot-per-account-per-day series with asset/liability split, excluding hidden accounts + +## 7. Deployment & verification + +- [ ] 7.1 Production build + run task; README covering Caddy route, env setup, key generation, backup expectations (SQLite file is secret-grade), and known limits (~90-day backfill, APP_URL change forces re-consent) +- [ ] 7.2 End-to-end walkthrough on real infrastructure: both users log in, claim real setup token, first sync lands, classify accounts, create rules, categorize manually, verify monthly report and net worth chart diff --git a/src/app.d.ts b/src/app.d.ts index 7c863b3..25eb888 100644 --- a/src/app.d.ts +++ b/src/app.d.ts @@ -10,6 +10,13 @@ declare global { // interface PageState {} // interface Platform {} } + + // Minimal surface of the Deno global used by server code; svelte-check runs + // plain tsc without Deno's type library. (deno test files carry their own + // /// and are excluded from svelte-check.) + const Deno: { + cron(name: string, schedule: string, handler: () => void | Promise): void; + }; } export {}; diff --git a/src/hooks.server.ts b/src/hooks.server.ts index 625a60a..115ebef 100644 --- a/src/hooks.server.ts +++ b/src/hooks.server.ts @@ -4,6 +4,7 @@ import { getConfig } from '$lib/server/config'; import { getDb, initDb } from '$lib/server/db'; import { initOAuthClient } from '$lib/server/auth/oauth-client'; import { deleteExpiredSessions, getSessionUser } from '$lib/server/services/sessions'; +import { runSync } from '$lib/server/services/sync'; export const init: ServerInit = async () => { if (building) return; @@ -11,6 +12,20 @@ export const init: ServerInit = async () => { const db = initDb(config.dbPath, 'migrations'); deleteExpiredSessions(db); await initOAuthClient(config, db); + + // Daily sync. The Bridge refreshes bank data roughly daily; more often is pointless. + try { + Deno.cron('daily simplefin sync', '0 11 * * *', async () => { + const outcomes = await runSync(getDb()); + for (const o of outcomes) { + console.log( + `sync connection ${o.connectionId}: ${o.ok ? 'ok' : `FAILED (${o.error})`}, +${o.newTransactions} txns, ${o.reconciled} reconciled, ${o.ruleCategorized} rule-categorized` + ); + } + }); + } catch (err) { + console.warn('Deno.cron unavailable; scheduled sync disabled:', err); + } }; export const SESSION_COOKIE = 'quantum_session'; diff --git a/src/lib/server/services/categorization.ts b/src/lib/server/services/categorization.ts new file mode 100644 index 0000000..1ee0c4c --- /dev/null +++ b/src/lib/server/services/categorization.ts @@ -0,0 +1,107 @@ +import type { DatabaseSync } from 'node:sqlite'; + +export type EventSource = 'rule' | 'manual' | 'reconciliation'; + +export interface CategorizationEventInput { + transactionId: number; + categoryId: number | null; // null clears the category + source: EventSource; + ruleId?: number; + actorDid?: string; +} + +/** + * Append an immutable categorization event and update the denormalized + * transactions.category_id cache in the same DB transaction (design D4). + */ +export function appendCategorizationEvent(db: DatabaseSync, input: CategorizationEventInput): void { + if (input.source === 'rule' && input.ruleId == null) { + throw new Error('rule events require ruleId'); + } + if (input.source === 'manual' && !input.actorDid) { + throw new Error('manual events require actorDid'); + } + db.exec('BEGIN'); + try { + db.prepare( + `INSERT INTO categorization_events + (transaction_id, category_id, source, rule_id, actor_did, created_at) + VALUES (?, ?, ?, ?, ?, ?)` + ).run( + input.transactionId, + input.categoryId, + input.source, + input.ruleId ?? null, + input.actorDid ?? null, + new Date().toISOString() + ); + db.prepare('UPDATE transactions SET category_id = ? WHERE id = ?').run( + input.categoryId, + input.transactionId + ); + db.exec('COMMIT'); + } catch (err) { + db.exec('ROLLBACK'); + throw err; + } +} + +export interface CategorizationEvent { + id: number; + transactionId: number; + categoryId: number | null; + categoryName: string | null; + source: EventSource; + ruleId: number | null; + rulePattern: string | null; + actorDid: string | null; + actorHandle: string | null; + createdAt: string; +} + +/** Full provenance history for a transaction, oldest first. */ +export function listEvents(db: DatabaseSync, transactionId: number): CategorizationEvent[] { + const rows = db + .prepare( + `SELECT e.id, e.transaction_id, e.category_id, c.name AS category_name, + e.source, e.rule_id, r.pattern AS rule_pattern, + e.actor_did, u.handle AS actor_handle, e.created_at + FROM categorization_events e + LEFT JOIN categories c ON c.id = e.category_id + LEFT JOIN rules r ON r.id = e.rule_id + LEFT JOIN users u ON u.did = e.actor_did + WHERE e.transaction_id = ? + ORDER BY e.id` + ) + .all(transactionId) as Record[]; + return rows.map((r) => ({ + id: r.id as number, + transactionId: r.transaction_id as number, + categoryId: r.category_id as number | null, + categoryName: r.category_name as string | null, + source: r.source as EventSource, + ruleId: r.rule_id as number | null, + rulePattern: r.rule_pattern as string | null, + actorDid: r.actor_did as string | null, + actorHandle: r.actor_handle as string | null, + createdAt: r.created_at as string + })); +} + +/** + * Manual categorization: allowed regardless of current state (humans outrank + * everything), always attributed to the acting user. + */ +export function categorizeManually( + db: DatabaseSync, + transactionId: number, + categoryId: number | null, + actorDid: string +): void { + appendCategorizationEvent(db, { + transactionId, + categoryId, + source: 'manual', + actorDid + }); +} diff --git a/src/lib/server/services/connections.test.ts b/src/lib/server/services/connections.test.ts new file mode 100644 index 0000000..c55aec9 --- /dev/null +++ b/src/lib/server/services/connections.test.ts @@ -0,0 +1,59 @@ +/// +import { openDatabase } from '../db.ts'; +import { claimSetupToken, listConnections } from './connections.ts'; +import { ClaimError } from '../simplefin.ts'; + +const MIGRATIONS_DIR = new URL('../../../../migrations', import.meta.url).pathname.replace( + /^\/([A-Za-z]:)/, + '$1' +); + +const CLAIM_URL = 'https://bridge.test/simplefin/claim/demo'; +const TOKEN = btoa(CLAIM_URL); +const ACCESS_URL = 'https://user:pass@bridge.test/simplefin'; + +Deno.test('valid setup token is claimed and the access URL stored', async () => { + const db = openDatabase(`${Deno.makeTempDirSync()}/t.db`, MIGRATIONS_DIR); + const fakeFetch = ((input: RequestInfo | URL, init?: RequestInit) => { + if (String(input) !== CLAIM_URL || init?.method !== 'POST') { + return Promise.resolve(new Response('wrong request', { status: 500 })); + } + return Promise.resolve(new Response(ACCESS_URL, { status: 200 })); + }) as typeof fetch; + + const connection = await claimSetupToken(db, TOKEN, fakeFetch); + if (connection.accessUrl !== ACCESS_URL) throw new Error('access URL mismatch'); + if (listConnections(db).length !== 1) throw new Error('connection not stored'); + db.close(); +}); + +Deno.test('already-claimed token surfaces a clear error and stores nothing', async () => { + const db = openDatabase(`${Deno.makeTempDirSync()}/t.db`, MIGRATIONS_DIR); + const fakeFetch = (() => Promise.resolve(new Response('', { status: 403 }))) as typeof fetch; + let caught: unknown; + try { + await claimSetupToken(db, TOKEN, fakeFetch); + } catch (err) { + caught = err; + } + if (!(caught instanceof ClaimError) || !caught.alreadyClaimed) { + throw new Error('expected ClaimError with alreadyClaimed=true'); + } + if (listConnections(db).length !== 0) throw new Error('connection should not be stored'); + db.close(); +}); + +Deno.test('garbage token is rejected before any network call', async () => { + const db = openDatabase(`${Deno.makeTempDirSync()}/t.db`, MIGRATIONS_DIR); + const explodingFetch = (() => { + throw new Error('network should not be touched'); + }) as typeof fetch; + let threw = false; + try { + await claimSetupToken(db, 'not-base64!!!', explodingFetch); + } catch (err) { + threw = err instanceof ClaimError; + } + if (!threw) throw new Error('expected ClaimError'); + db.close(); +}); diff --git a/src/lib/server/services/connections.ts b/src/lib/server/services/connections.ts new file mode 100644 index 0000000..fbba268 --- /dev/null +++ b/src/lib/server/services/connections.ts @@ -0,0 +1,34 @@ +import type { DatabaseSync } from 'node:sqlite'; +import { claimAccessUrl, decodeSetupToken } from '../simplefin.ts'; + +export interface Connection { + id: number; + accessUrl: string; + claimedAt: string; +} + +export function listConnections(db: DatabaseSync): Connection[] { + const rows = db + .prepare('SELECT id, access_url, claimed_at FROM connections ORDER BY id') + .all() as Record[]; + return rows.map((r) => ({ + id: r.id as number, + accessUrl: r.access_url as string, + claimedAt: r.claimed_at as string + })); +} + +/** Claim a pasted setup token and persist the resulting Access URL. */ +export async function claimSetupToken( + db: DatabaseSync, + token: string, + fetchFn: typeof fetch = fetch +): Promise { + const claimUrl = decodeSetupToken(token); + const accessUrl = await claimAccessUrl(claimUrl, fetchFn); + const claimedAt = new Date().toISOString(); + const result = db + .prepare('INSERT INTO connections (access_url, claimed_at) VALUES (?, ?)') + .run(accessUrl, claimedAt); + return { id: Number(result.lastInsertRowid), accessUrl, claimedAt }; +} diff --git a/src/lib/server/services/normalize.test.ts b/src/lib/server/services/normalize.test.ts new file mode 100644 index 0000000..d50c044 --- /dev/null +++ b/src/lib/server/services/normalize.test.ts @@ -0,0 +1,65 @@ +/// +import { normalizePayload, parseAmountToCents } from './normalize.ts'; + +function assertEq(actual: unknown, expected: unknown, label = '') { + if (actual !== expected) throw new Error(`${label} expected ${expected}, got ${actual}`); +} + +Deno.test('parseAmountToCents handles SimpleFIN decimal strings exactly', () => { + assertEq(parseAmountToCents('123.45'), 12345); + assertEq(parseAmountToCents('-123.45'), -12345); + assertEq(parseAmountToCents('0.01'), 1); + assertEq(parseAmountToCents('-0.01'), -1); + assertEq(parseAmountToCents('1'), 100); + assertEq(parseAmountToCents('1.5'), 150); + assertEq(parseAmountToCents('1.005'), 101, 'rounds half away from zero'); + assertEq(parseAmountToCents('-1.005'), -101); + assertEq(parseAmountToCents('4222.19'), 422219, 'no float drift'); + let threw = false; + try { + parseAmountToCents('12,34'); + } catch { + threw = true; + } + if (!threw) throw new Error('expected malformed amount to throw'); +}); + +Deno.test('normalizePayload maps accounts, transactions, errors', () => { + const payload = JSON.stringify({ + errors: ['Connection to Chase needs attention'], + accounts: [ + { + org: { name: 'Chase', domain: 'chase.com', 'sfin-url': 'https://x' }, + id: 'act-1', + name: 'Checking', + currency: 'USD', + balance: '1500.25', + 'available-balance': '1450.00', + 'balance-date': 1750000000, + transactions: [ + { + id: 'txn-1', + posted: 1749900000, + amount: '-42.19', + description: 'KROGER #123', + extra: { memo: 'grocery' } + }, + { id: 'txn-2', posted: 0, pending: true, amount: '-9.99', description: 'PENDING COFFEE' } + ] + } + ] + }); + const result = normalizePayload(payload); + assertEq(result.errors.length, 1); + assertEq(result.accounts.length, 1); + const account = result.accounts[0]; + assertEq(account.balanceCents, 150025); + assertEq(account.availableBalanceCents, 145000); + assertEq(account.org.name, 'Chase'); + const [posted, pending] = account.transactions; + assertEq(posted.pending, false); + assertEq(posted.amountCents, -4219); + assertEq(posted.extra, '{"memo":"grocery"}'); + assertEq(pending.pending, true); + assertEq(pending.posted, null, 'pending posted=0 becomes null'); +}); diff --git a/src/lib/server/services/normalize.ts b/src/lib/server/services/normalize.ts new file mode 100644 index 0000000..dab0d26 --- /dev/null +++ b/src/lib/server/services/normalize.ts @@ -0,0 +1,98 @@ +// Pure normalization of a SimpleFIN /accounts payload. No I/O, no DB: +// a function from archived payload text to typed rows, so history can be +// replayed and behavior tested against fixtures. + +export interface NormalizedOrg { + name: string | null; + domain: string | null; + sfinUrl: string | null; +} + +export interface NormalizedAccount { + id: string; + org: NormalizedOrg; + name: string; + currency: string; + balanceCents: number; + availableBalanceCents: number | null; + balanceDate: number | null; + transactions: NormalizedTransaction[]; +} + +export interface NormalizedTransaction { + sfinId: string; + posted: number | null; // Unix seconds; null while pending + transactedAt: number | null; + amountCents: number; + description: string; + pending: boolean; + extra: string | null; // verbatim JSON +} + +export interface NormalizedPayload { + accounts: NormalizedAccount[]; + /** Connection-level error strings, verbatim from the feed. */ + errors: string[]; +} + +/** + * Parse a SimpleFIN decimal string into integer cents without floats. + * Assumes 2-minor-unit currencies; extra fraction digits are rounded + * (half away from zero). + */ +export function parseAmountToCents(amount: string | number): number { + const text = String(amount).trim(); + const match = text.match(/^(-?)(\d+)(?:\.(\d+))?$/); + if (!match) throw new Error(`Unparseable amount: ${JSON.stringify(amount)}`); + const [, sign, whole, fracRaw = ''] = match; + const frac = (fracRaw + '00').slice(0, 2); + let cents = Number(whole) * 100 + Number(frac); + if (fracRaw.length > 2 && Number(fracRaw[2]) >= 5) cents += 1; + return sign === '-' ? -cents : cents; +} + +export function normalizePayload(payloadText: string): NormalizedPayload { + const raw = JSON.parse(payloadText) as { + errors?: unknown[]; + accounts?: Record[]; + }; + + const errors = (raw.errors ?? []).map((e) => String(e)); + const accounts: NormalizedAccount[] = (raw.accounts ?? []).map((account) => { + const org = (account.org ?? {}) as Record; + const transactions = ((account.transactions ?? []) as Record[]).map( + (txn): NormalizedTransaction => { + const posted = typeof txn.posted === 'number' && txn.posted > 0 ? txn.posted : null; + const pending = txn.pending === true || posted === null; + return { + sfinId: String(txn.id), + posted: pending ? null : posted, + transactedAt: typeof txn.transacted_at === 'number' ? txn.transacted_at : null, + amountCents: parseAmountToCents(txn.amount as string), + description: String(txn.description ?? ''), + pending, + extra: txn.extra !== undefined ? JSON.stringify(txn.extra) : null + }; + } + ); + return { + id: String(account.id), + org: { + name: org.name != null ? String(org.name) : null, + domain: org.domain != null ? String(org.domain) : null, + sfinUrl: org['sfin-url'] != null ? String(org['sfin-url']) : null + }, + name: String(account.name ?? ''), + currency: String(account.currency ?? 'USD'), + balanceCents: parseAmountToCents(account.balance as string), + availableBalanceCents: + account['available-balance'] != null + ? parseAmountToCents(account['available-balance'] as string) + : null, + balanceDate: typeof account['balance-date'] === 'number' ? account['balance-date'] : null, + transactions + }; + }); + + return { accounts, errors }; +} diff --git a/src/lib/server/services/rules.test.ts b/src/lib/server/services/rules.test.ts new file mode 100644 index 0000000..f5399e1 --- /dev/null +++ b/src/lib/server/services/rules.test.ts @@ -0,0 +1,118 @@ +/// +import { openDatabase } from '../db.ts'; +import { applyRulesToUncategorized, createRule, findWinningRule, type Rule } from './rules.ts'; +import { categorizeManually } from './categorization.ts'; +import type { DatabaseSync } from 'node:sqlite'; + +const MIGRATIONS_DIR = new URL('../../../../migrations', import.meta.url).pathname.replace( + /^\/([A-Za-z]:)/, + '$1' +); + +function rule(partial: Partial & { pattern: string }): Rule { + return { + id: 1, + matchType: 'contains', + categoryId: 1, + createdByDid: 'did:plc:t', + active: true, + createdAt: '2026-01-01', + ...partial + }; +} + +Deno.test('precedence: exact beats contains, longer beats shorter, newer beats older', () => { + const contains = rule({ id: 1, pattern: 'AMAZON' }); + const longer = rule({ id: 2, pattern: 'AMAZON PRIME' }); + const exact = rule({ id: 3, matchType: 'exact', pattern: 'amazon prime *2k4l' }); + const newerSameLength = rule({ id: 4, pattern: 'aMaZoN' }); + + const desc = 'AMAZON PRIME *2K4L'; + if (findWinningRule([contains, longer], desc)?.id !== 2) + throw new Error('longer pattern should win'); + if (findWinningRule([contains, longer, exact], desc)?.id !== 3) + throw new Error('exact should beat contains'); + if (findWinningRule([contains, newerSameLength], desc)?.id !== 4) + throw new Error('newer rule should win the tie'); + if (findWinningRule([contains], 'ETSY ORDER') !== null) + throw new Error('non-matching description should return null'); +}); + +function testDb(): DatabaseSync { + const db = openDatabase(`${Deno.makeTempDirSync()}/test.db`, MIGRATIONS_DIR); + const now = new Date().toISOString(); + db.prepare("INSERT INTO users (did, handle, created_at) VALUES ('did:plc:t','tester',?)").run(now); + db.prepare( + "INSERT INTO connections (access_url, claimed_at) VALUES ('https://u:p@x/simplefin',?)" + ).run(now); + db.prepare( + `INSERT INTO accounts (id, connection_id, name, currency, state, created_at) + VALUES ('act-1', 1, 'Checking', 'USD', 'ACTIVE', ?)` + ).run(now); + db.prepare("INSERT INTO categories (name, kind, created_at) VALUES ('Groceries','expense',?)").run( + now + ); + return db; +} + +function insertTxn(db: DatabaseSync, sfinId: string, description: string): number { + const result = db + .prepare( + `INSERT INTO transactions (account_id, sfin_id, posted, amount_cents, description, pending, created_at) + VALUES ('act-1', ?, 1750000000, -1000, ?, 0, ?)` + ) + .run(sfinId, description, new Date().toISOString()); + return Number(result.lastInsertRowid); +} + +Deno.test('rules fire on uncategorized transactions and record the winning rule id', () => { + const db = testDb(); + const categoryId = 2; // after built-in Transfer (1) + const created = createRule(db, { + matchType: 'contains', + pattern: 'KROGER', + categoryId, + createdByDid: 'did:plc:t' + }); + insertTxn(db, 't1', 'KROGER #123 COLUMBUS'); + insertTxn(db, 't2', 'SHELL OIL'); + const applied = applyRulesToUncategorized(db); + if (applied !== 1) throw new Error(`expected 1 applied, got ${applied}`); + const event = db + .prepare('SELECT source, rule_id FROM categorization_events') + .get() as { source: string; rule_id: number }; + if (event.source !== 'rule' || event.rule_id !== created.id) + throw new Error('event should record the firing rule'); + db.close(); +}); + +Deno.test('rules never overwrite a manual decision, including manual uncategorize', () => { + const db = testDb(); + createRule(db, { + matchType: 'contains', + pattern: 'KROGER', + categoryId: 2, + createdByDid: 'did:plc:t' + }); + const txnId = insertTxn(db, 't1', 'KROGER #123'); + + // manual categorization to a different category + db.prepare("INSERT INTO categories (name, kind, created_at) VALUES ('Dining','expense',?)").run( + new Date().toISOString() + ); + categorizeManually(db, txnId, 3, 'did:plc:t'); + applyRulesToUncategorized(db); + let cached = (db.prepare('SELECT category_id FROM transactions WHERE id = ?').get(txnId) as { + category_id: number; + }).category_id; + if (cached !== 3) throw new Error('rule overwrote manual categorization'); + + // manual uncategorize: latest event is manual with NULL → rules stay away + categorizeManually(db, txnId, null, 'did:plc:t'); + applyRulesToUncategorized(db); + cached = (db.prepare('SELECT category_id FROM transactions WHERE id = ?').get(txnId) as { + category_id: number | null; + }).category_id as number; + if (cached !== null) throw new Error('rule overrode a manual uncategorize'); + db.close(); +}); diff --git a/src/lib/server/services/rules.ts b/src/lib/server/services/rules.ts new file mode 100644 index 0000000..2966931 --- /dev/null +++ b/src/lib/server/services/rules.ts @@ -0,0 +1,154 @@ +import type { DatabaseSync } from 'node:sqlite'; +import { appendCategorizationEvent } from './categorization.ts'; + +export interface Rule { + id: number; + matchType: 'exact' | 'contains'; + pattern: string; + categoryId: number; + createdByDid: string; + active: boolean; + createdAt: string; +} + +export function listRules(db: DatabaseSync, options: { activeOnly?: boolean } = {}): Rule[] { + const rows = db + .prepare( + `SELECT id, match_type, pattern, category_id, created_by_did, active, created_at + FROM rules ${options.activeOnly ? 'WHERE active = 1' : ''} ORDER BY id` + ) + .all() as Record[]; + return rows.map((r) => ({ + id: r.id as number, + matchType: r.match_type as 'exact' | 'contains', + pattern: r.pattern as string, + categoryId: r.category_id as number, + createdByDid: r.created_by_did as string, + active: r.active === 1, + createdAt: r.created_at as string + })); +} + +export function createRule( + db: DatabaseSync, + input: { + matchType: 'exact' | 'contains'; + pattern: string; + categoryId: number; + createdByDid: string; + } +): Rule { + const createdAt = new Date().toISOString(); + const result = db + .prepare( + `INSERT INTO rules (match_type, pattern, category_id, created_by_did, active, created_at) + VALUES (?, ?, ?, ?, 1, ?)` + ) + .run(input.matchType, input.pattern, input.categoryId, input.createdByDid, createdAt); + return { + id: Number(result.lastInsertRowid), + ...input, + active: true, + createdAt + }; +} + +export function setRuleActive(db: DatabaseSync, ruleId: number, active: boolean): void { + db.prepare('UPDATE rules SET active = ? WHERE id = ?').run(active ? 1 : 0, ruleId); +} + +function matches(rule: Rule, description: string): boolean { + const desc = description.toLowerCase(); + const pattern = rule.pattern.toLowerCase(); + return rule.matchType === 'exact' ? desc === pattern : desc.includes(pattern); +} + +/** + * Deterministic precedence when several rules match one description + * (design D4): exact beats contains, longer pattern beats shorter, + * newer rule (higher id) beats older. No manual ordering. + */ +export function findWinningRule(rules: Rule[], description: string): Rule | null { + let winner: Rule | null = null; + for (const rule of rules) { + if (!matches(rule, description)) continue; + if (!winner) { + winner = rule; + continue; + } + const exactness = Number(rule.matchType === 'exact') - Number(winner.matchType === 'exact'); + const length = rule.pattern.length - winner.pattern.length; + if (exactness > 0 || (exactness === 0 && (length > 0 || (length === 0 && rule.id > winner.id)))) { + winner = rule; + } + } + return winner; +} + +/** + * Fire rules over currently-uncategorized, non-removed transactions. + * The invariant "rules never overwrite a human decision" holds structurally: + * only transactions with no current category are considered, and manual + * events always set category_id (or clear it, which a rule may then refill — + * an explicit uncategorize invites re-ruling only via this same path). + * Returns the number of transactions categorized. + */ +export function applyRulesToUncategorized( + db: DatabaseSync, + options: { ruleIds?: number[] } = {} +): number { + let rules = listRules(db, { activeOnly: true }); + if (options.ruleIds) rules = rules.filter((r) => options.ruleIds!.includes(r.id)); + if (rules.length === 0) return 0; + + // A transaction whose latest event is manual is skipped even when its + // category is currently NULL (a human explicitly uncategorized it). + const candidates = db + .prepare( + `SELECT t.id, t.description FROM transactions t + WHERE t.category_id IS NULL AND t.removed_at IS NULL + AND NOT EXISTS ( + SELECT 1 FROM categorization_events e + WHERE e.transaction_id = t.id AND e.source = 'manual' + AND e.id = (SELECT MAX(id) FROM categorization_events WHERE transaction_id = t.id) + )` + ) + .all() as { id: number; description: string }[]; + + let count = 0; + for (const txn of candidates) { + const winner = findWinningRule(rules, txn.description); + if (!winner) continue; + appendCategorizationEvent(db, { + transactionId: txn.id, + categoryId: winner.categoryId, + source: 'rule', + ruleId: winner.id + }); + count++; + } + return count; +} + +/** Count how many uncategorized transactions a prospective rule would hit. */ +export function countRuleMatches( + db: DatabaseSync, + matchType: 'exact' | 'contains', + pattern: string +): number { + const probe: Rule = { + id: 0, + matchType, + pattern, + categoryId: 0, + createdByDid: '', + active: true, + createdAt: '' + }; + const rows = db + .prepare( + 'SELECT description FROM transactions WHERE category_id IS NULL AND removed_at IS NULL' + ) + .all() as { description: string }[]; + return rows.filter((r) => matches(probe, r.description)).length; +} diff --git a/src/lib/server/services/sync-status.ts b/src/lib/server/services/sync-status.ts new file mode 100644 index 0000000..9016ba9 --- /dev/null +++ b/src/lib/server/services/sync-status.ts @@ -0,0 +1,28 @@ +import type { DatabaseSync } from 'node:sqlite'; + +export interface LastSync { + fetchedAt: string; + ok: boolean; + error: string | null; +} + +export function getLastSync(db: DatabaseSync): LastSync | null { + const row = db + .prepare('SELECT fetched_at, ok, error FROM raw_syncs ORDER BY id DESC LIMIT 1') + .get() as { fetched_at: string; ok: number; error: string | null } | undefined; + return row ? { fetchedAt: row.fetched_at, ok: row.ok === 1, error: row.error } : null; +} + +/** Connection-level errors reported by the most recent successful sync. */ +export function getConnectionErrors(db: DatabaseSync): string[] { + const row = db + .prepare('SELECT payload FROM raw_syncs WHERE ok = 1 ORDER BY id DESC LIMIT 1') + .get() as { payload: string } | undefined; + if (!row?.payload) return []; + try { + const parsed = JSON.parse(row.payload) as { errors?: unknown[] }; + return (parsed.errors ?? []).map((e) => String(e)); + } catch { + return []; + } +} diff --git a/src/lib/server/services/sync.test.ts b/src/lib/server/services/sync.test.ts new file mode 100644 index 0000000..3c7189b --- /dev/null +++ b/src/lib/server/services/sync.test.ts @@ -0,0 +1,203 @@ +/// +import { openDatabase } from '../db.ts'; +import { runSync } from './sync.ts'; +import { categorizeManually } from './categorization.ts'; +import type { DatabaseSync } from 'node:sqlite'; + +const MIGRATIONS_DIR = new URL('../../../../migrations', import.meta.url).pathname.replace( + /^\/([A-Za-z]:)/, + '$1' +); + +function testDb(): DatabaseSync { + const db = openDatabase(`${Deno.makeTempDirSync()}/test.db`, MIGRATIONS_DIR); + db.prepare( + "INSERT INTO connections (access_url, claimed_at) VALUES ('https://u:p@bridge.test/simplefin', ?)" + ).run(new Date().toISOString()); + db.prepare("INSERT INTO users (did, handle, created_at) VALUES ('did:plc:t', 'tester', ?)").run( + new Date().toISOString() + ); + return db; +} + +function fakeFetch(payload: unknown): typeof fetch { + return (() => + Promise.resolve(new Response(JSON.stringify(payload), { status: 200 }))) as typeof fetch; +} + +const NOW_S = Math.floor(Date.now() / 1000); + +function payloadWith(transactions: unknown[], accountOverrides: Record = {}) { + return { + errors: [], + accounts: [ + { + org: { name: 'Test Bank', domain: 'bank.test' }, + id: 'act-1', + name: 'Checking', + currency: 'USD', + balance: '100.00', + 'balance-date': NOW_S, + transactions, + ...accountOverrides + } + ] + }; +} + +function count(db: DatabaseSync, sql: string): number { + return (db.prepare(sql).get() as { n: number }).n; +} + +Deno.test('sync archives raw payload, discovers account as NEW, snapshots balance', async () => { + const db = testDb(); + const payload = payloadWith([ + { id: 't1', posted: NOW_S - 1000, amount: '-10.00', description: 'STORE A' } + ]); + const [outcome] = await runSync(db, fakeFetch(payload)); + if (!outcome.ok) throw new Error(`sync failed: ${outcome.error}`); + if (outcome.newTransactions !== 1) throw new Error('expected 1 new transaction'); + if (count(db, 'SELECT COUNT(*) n FROM raw_syncs WHERE ok = 1') !== 1) + throw new Error('raw payload not archived'); + const account = db.prepare("SELECT state FROM accounts WHERE id = 'act-1'").get() as { + state: string; + }; + if (account.state !== 'NEW') throw new Error(`expected NEW, got ${account.state}`); + if (count(db, 'SELECT COUNT(*) n FROM balance_snapshots') !== 1) + throw new Error('expected 1 balance snapshot'); + db.close(); +}); + +Deno.test('re-syncing the same payload is a no-op for transactions', async () => { + const db = testDb(); + const payload = payloadWith([ + { id: 't1', posted: NOW_S - 1000, amount: '-10.00', description: 'STORE A' }, + { id: 't2', posted: NOW_S - 2000, amount: '2500.00', description: 'PAYROLL' } + ]); + await runSync(db, fakeFetch(payload)); + const [second] = await runSync(db, fakeFetch(payload)); + if (second.newTransactions !== 0) throw new Error('second sync inserted duplicates'); + if (count(db, 'SELECT COUNT(*) n FROM transactions') !== 2) + throw new Error('expected exactly 2 transactions'); + // snapshots intentionally accumulate per sync + if (count(db, 'SELECT COUNT(*) n FROM balance_snapshots') !== 2) + throw new Error('expected 2 snapshots'); + db.close(); +}); + +Deno.test('categorized pending transaction posting under a new id carries category via reconciliation event', async () => { + const db = testDb(); + await runSync( + db, + fakeFetch( + payloadWith([ + { + id: 'pend-1', + posted: 0, + pending: true, + transacted_at: NOW_S - 3600, + amount: '-25.00', + description: 'COFFEE SHOP (PENDING)' + } + ]) + ) + ); + const pendingRow = db.prepare('SELECT id FROM transactions').get() as { id: number }; + db.prepare("INSERT INTO categories (name, kind, created_at) VALUES ('Coffee','expense',?)").run( + new Date().toISOString() + ); + const categoryId = Number( + (db.prepare("SELECT id FROM categories WHERE name='Coffee'").get() as { id: number }).id + ); + categorizeManually(db, pendingRow.id, categoryId, 'did:plc:t'); + + // Next sync: pending vanished, posted appears under a different id. + const [outcome] = await runSync( + db, + fakeFetch( + payloadWith([ + { id: 'post-9', posted: NOW_S, amount: '-25.00', description: 'COFFEE SHOP' } + ]) + ) + ); + if (outcome.reconciled !== 1) throw new Error(`expected 1 reconciled, got ${outcome.reconciled}`); + const txn = db + .prepare("SELECT id, sfin_id, pending, category_id, removed_at FROM transactions") + .get() as Record; + if (txn.sfin_id !== 'post-9') throw new Error('row not replaced in place'); + if (txn.pending !== 0) throw new Error('still pending'); + if (txn.category_id !== categoryId) throw new Error('category not carried'); + if (txn.removed_at !== null) throw new Error('should not be removed'); + const events = db + .prepare('SELECT source FROM categorization_events ORDER BY id') + .all() as { source: string }[]; + if (events.map((e) => e.source).join(',') !== 'manual,reconciliation') + throw new Error(`unexpected event trail: ${events.map((e) => e.source)}`); + db.close(); +}); + +Deno.test('stale pending transaction with no posted match is soft-removed', async () => { + const db = testDb(); + await runSync( + db, + fakeFetch( + payloadWith([ + { id: 'pend-2', posted: 0, pending: true, amount: '-5.00', description: 'GHOST' } + ]) + ) + ); + const [outcome] = await runSync(db, fakeFetch(payloadWith([]))); + if (outcome.removedPending !== 1) throw new Error('expected 1 removed pending'); + const row = db.prepare('SELECT removed_at FROM transactions').get() as { + removed_at: string | null; + }; + if (!row.removed_at) throw new Error('pending row not soft-removed'); + db.close(); +}); + +Deno.test('ACTIVE account absent from feed goes INACTIVE and returns on reappearance', async () => { + const db = testDb(); + await runSync(db, fakeFetch(payloadWith([]))); + db.prepare("UPDATE accounts SET state='ACTIVE', account_type='checking' WHERE id='act-1'").run(); + + // feed with a different account only + const otherAccount = { + errors: [], + accounts: [ + { + org: { name: 'Other' }, + id: 'act-2', + name: 'Savings', + currency: 'USD', + balance: '1.00', + transactions: [] + } + ] + }; + await runSync(db, fakeFetch(otherAccount)); + let state = (db.prepare("SELECT state FROM accounts WHERE id='act-1'").get() as { state: string }) + .state; + if (state !== 'INACTIVE') throw new Error(`expected INACTIVE, got ${state}`); + + await runSync(db, fakeFetch(payloadWith([]))); + state = (db.prepare("SELECT state FROM accounts WHERE id='act-1'").get() as { state: string }) + .state; + if (state !== 'ACTIVE') throw new Error(`expected ACTIVE after reappearing, got ${state}`); + db.close(); +}); + +Deno.test('failed fetch archives a failure row and touches nothing else', async () => { + const db = testDb(); + const failingFetch = (() => Promise.reject(new Error('network down'))) as unknown as typeof fetch; + const [outcome] = await runSync(db, failingFetch); + if (outcome.ok) throw new Error('expected failure'); + const raw = db.prepare('SELECT ok, error FROM raw_syncs').get() as { + ok: number; + error: string; + }; + if (raw.ok !== 0 || !raw.error.includes('network down')) + throw new Error('failure not recorded'); + if (count(db, 'SELECT COUNT(*) n FROM transactions') !== 0) + throw new Error('transactions should be untouched'); + db.close(); +}); diff --git a/src/lib/server/services/sync.ts b/src/lib/server/services/sync.ts new file mode 100644 index 0000000..a9246b5 --- /dev/null +++ b/src/lib/server/services/sync.ts @@ -0,0 +1,271 @@ +import type { DatabaseSync } from 'node:sqlite'; +import { fetchAccounts } from '../simplefin.ts'; +import { normalizePayload, type NormalizedAccount } from './normalize.ts'; +import { listConnections } from './connections.ts'; +import { appendCategorizationEvent } from './categorization.ts'; +import { applyRulesToUncategorized } from './rules.ts'; + +const FIRST_SYNC_LOOKBACK_DAYS = 365; +const INCREMENTAL_OVERLAP_DAYS = 30; +/** Max gap between a pending transaction and the posted one replacing it. */ +const RECONCILE_WINDOW_DAYS = 5; + +export interface SyncOutcome { + connectionId: number; + ok: boolean; + error?: string; + newTransactions: number; + reconciled: number; + removedPending: number; + ruleCategorized: number; + connectionErrors: string[]; +} + +/** Sync every connection: archive verbatim, normalize, snapshot, reconcile, rule-fire. */ +export async function runSync( + db: DatabaseSync, + fetchFn: typeof fetch = fetch +): Promise { + const outcomes: SyncOutcome[] = []; + for (const connection of listConnections(db)) { + outcomes.push(await syncConnection(db, connection.id, connection.accessUrl, fetchFn)); + } + return outcomes; +} + +async function syncConnection( + db: DatabaseSync, + connectionId: number, + accessUrl: string, + fetchFn: typeof fetch +): Promise { + const hasPriorSync = db + .prepare('SELECT 1 FROM raw_syncs WHERE connection_id = ? AND ok = 1 LIMIT 1') + .get(connectionId); + const lookbackDays = hasPriorSync ? INCREMENTAL_OVERLAP_DAYS : FIRST_SYNC_LOOKBACK_DAYS; + const startDate = new Date(Date.now() - lookbackDays * 86400_000); + + const fetched = await fetchAccounts(accessUrl, { startDate, pending: true }, fetchFn); + const fetchedAt = new Date().toISOString(); + + // Archive verbatim before any processing — success or failure. + db.prepare( + 'INSERT INTO raw_syncs (connection_id, fetched_at, ok, payload, error) VALUES (?, ?, ?, ?, ?)' + ).run(connectionId, fetchedAt, fetched.ok ? 1 : 0, fetched.body ?? null, fetched.error ?? null); + + const outcome: SyncOutcome = { + connectionId, + ok: fetched.ok, + error: fetched.error, + newTransactions: 0, + reconciled: 0, + removedPending: 0, + ruleCategorized: 0, + connectionErrors: [] + }; + if (!fetched.ok || !fetched.body) return outcome; + + let normalized; + try { + normalized = normalizePayload(fetched.body); + } catch (err) { + outcome.ok = false; + outcome.error = `Normalization failed: ${err instanceof Error ? err.message : err}`; + db.prepare('UPDATE raw_syncs SET ok = 0, error = ? WHERE connection_id = ? AND fetched_at = ?').run( + outcome.error, + connectionId, + fetchedAt + ); + return outcome; + } + outcome.connectionErrors = normalized.errors; + + for (const account of normalized.accounts) { + upsertAccount(db, connectionId, account, fetchedAt); + const result = ingestTransactions(db, account); + outcome.newTransactions += result.inserted; + outcome.reconciled += result.reconciled; + outcome.removedPending += result.removedPending; + db.prepare( + 'INSERT INTO balance_snapshots (account_id, captured_at, balance_cents, available_balance_cents) VALUES (?, ?, ?, ?)' + ).run(account.id, fetchedAt, account.balanceCents, account.availableBalanceCents); + } + + // Lifecycle: ACTIVE accounts of this connection absent from the feed go INACTIVE. + const seenIds = normalized.accounts.map((a) => a.id); + const placeholders = seenIds.map(() => '?').join(','); + db.prepare( + `UPDATE accounts SET state = 'INACTIVE' + WHERE connection_id = ? AND state = 'ACTIVE' + ${seenIds.length ? `AND id NOT IN (${placeholders})` : ''}` + ).run(connectionId, ...seenIds); + + outcome.ruleCategorized = applyRulesToUncategorized(db); + return outcome; +} + +function upsertAccount( + db: DatabaseSync, + connectionId: number, + account: NormalizedAccount, + fetchedAt: string +): void { + const existing = db.prepare('SELECT id, state FROM accounts WHERE id = ?').get(account.id) as + | { id: string; state: string } + | undefined; + if (!existing) { + db.prepare( + `INSERT INTO accounts + (id, connection_id, org_name, org_domain, org_sfin_url, name, currency, state, last_successful_data_at, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, 'NEW', ?, ?)` + ).run( + account.id, + connectionId, + account.org.name, + account.org.domain, + account.org.sfinUrl, + account.name, + account.currency, + fetchedAt, + fetchedAt + ); + return; + } + // Reappearing INACTIVE accounts return to ACTIVE if already classified, NEW otherwise. + db.prepare( + `UPDATE accounts SET + org_name = ?, org_domain = ?, org_sfin_url = ?, name = ?, currency = ?, + last_successful_data_at = ?, + state = CASE + WHEN state = 'INACTIVE' AND account_type IS NOT NULL THEN 'ACTIVE' + WHEN state = 'INACTIVE' THEN 'NEW' + ELSE state + END + WHERE id = ?` + ).run( + account.org.name, + account.org.domain, + account.org.sfinUrl, + account.name, + account.currency, + fetchedAt, + account.id + ); +} + +function ingestTransactions( + db: DatabaseSync, + account: NormalizedAccount +): { inserted: number; reconciled: number; removedPending: number } { + const now = new Date().toISOString(); + let inserted = 0; + let reconciled = 0; + + const feedSfinIds = new Set(account.transactions.map((t) => t.sfinId)); + + for (const txn of account.transactions) { + const existing = db + .prepare('SELECT id, pending FROM transactions WHERE account_id = ? AND sfin_id = ?') + .get(account.id, txn.sfinId) as { id: number; pending: number } | undefined; + + if (existing) { + // Idempotent refresh; a known-pending transaction that posts under the + // SAME id is just updated in place. + db.prepare( + `UPDATE transactions SET posted = ?, transacted_at = ?, amount_cents = ?, + description = ?, pending = ?, extra = ?, removed_at = NULL + WHERE id = ?` + ).run( + txn.posted, + txn.transactedAt, + txn.amountCents, + txn.description, + txn.pending ? 1 : 0, + txn.extra, + existing.id + ); + continue; + } + + if (!txn.pending) { + // New posted transaction: does it replace a pending row that vanished + // from the feed? Match: same account + exact amount + date proximity. + const windowStart = (txn.posted ?? 0) - RECONCILE_WINDOW_DAYS * 86400; + const windowEnd = (txn.posted ?? 0) + RECONCILE_WINDOW_DAYS * 86400; + const match = db + .prepare( + `SELECT id, sfin_id FROM transactions + WHERE account_id = ? AND pending = 1 AND removed_at IS NULL AND amount_cents = ? + AND COALESCE(transacted_at, posted, ?) BETWEEN ? AND ? + ORDER BY id LIMIT 1` + ) + .get(account.id, txn.amountCents, txn.posted ?? 0, windowStart, windowEnd) as + | { id: number; sfin_id: string } + | undefined; + + if (match && !feedSfinIds.has(match.sfin_id)) { + // Replace in place: same row id keeps the event history attached; + // a reconciliation event records the carry-forward. + const category = db + .prepare('SELECT category_id FROM transactions WHERE id = ?') + .get(match.id) as { category_id: number | null }; + db.prepare( + `UPDATE transactions SET sfin_id = ?, posted = ?, transacted_at = ?, amount_cents = ?, + description = ?, pending = 0, extra = ?, removed_at = NULL + WHERE id = ?` + ).run( + txn.sfinId, + txn.posted, + txn.transactedAt, + txn.amountCents, + txn.description, + txn.extra, + match.id + ); + if (category.category_id != null) { + appendCategorizationEvent(db, { + transactionId: match.id, + categoryId: category.category_id, + source: 'reconciliation' + }); + } + reconciled++; + continue; + } + } + + db.prepare( + `INSERT INTO transactions + (account_id, sfin_id, posted, transacted_at, amount_cents, description, pending, extra, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)` + ).run( + account.id, + txn.sfinId, + txn.posted, + txn.transactedAt, + txn.amountCents, + txn.description, + txn.pending ? 1 : 0, + txn.extra, + now + ); + inserted++; + } + + // Stale pending rows: still pending in the DB, absent from the feed, and + // not matched by any posted transaction above. Soft-remove (events are + // never deleted). + const stale = db + .prepare( + 'SELECT id, sfin_id FROM transactions WHERE account_id = ? AND pending = 1 AND removed_at IS NULL' + ) + .all(account.id) as { id: number; sfin_id: string }[]; + let removedPending = 0; + for (const row of stale) { + if (feedSfinIds.has(row.sfin_id)) continue; + db.prepare('UPDATE transactions SET removed_at = ? WHERE id = ?').run(now, row.id); + removedPending++; + } + + return { inserted, reconciled, removedPending }; +} diff --git a/src/lib/server/simplefin.ts b/src/lib/server/simplefin.ts new file mode 100644 index 0000000..8068d91 --- /dev/null +++ b/src/lib/server/simplefin.ts @@ -0,0 +1,86 @@ +// SimpleFIN protocol plumbing: token claim and authenticated /accounts fetch. +// The Access URL embeds Basic Auth credentials (https://user:pass@host/simplefin); +// fetch() rejects credentialed URLs, so they are extracted into a header. + +export class ClaimError extends Error { + constructor( + message: string, + public readonly alreadyClaimed: boolean + ) { + super(message); + } +} + +/** Decode a one-time setup token (base64-encoded claim URL). */ +export function decodeSetupToken(token: string): string { + let claimUrl: string; + try { + claimUrl = atob(token.trim()); + } catch { + throw new ClaimError('That does not look like a SimpleFIN setup token.', false); + } + if (!claimUrl.startsWith('https://')) { + throw new ClaimError('Setup token did not decode to an https claim URL.', false); + } + return claimUrl; +} + +/** Claim a setup token, returning the permanent Access URL. */ +export async function claimAccessUrl( + claimUrl: string, + fetchFn: typeof fetch = fetch +): Promise { + const res = await fetchFn(claimUrl, { method: 'POST' }); + if (res.status === 403) { + throw new ClaimError( + 'This token was already claimed. Generate a fresh one at SimpleFIN Bridge.', + true + ); + } + if (!res.ok) { + throw new ClaimError(`Claim failed (HTTP ${res.status}).`, false); + } + const accessUrl = (await res.text()).trim(); + if (!accessUrl.startsWith('https://')) { + throw new ClaimError('Claim response was not an Access URL.', false); + } + return accessUrl; +} + +export function parseAccessUrl(accessUrl: string): { baseUrl: string; authHeader: string } { + const url = new URL(accessUrl); + const authHeader = `Basic ${btoa(`${decodeURIComponent(url.username)}:${decodeURIComponent(url.password)}`)}`; + url.username = ''; + url.password = ''; + return { baseUrl: url.toString().replace(/\/$/, ''), authHeader }; +} + +export interface FetchAccountsResult { + ok: boolean; + status?: number; + /** Verbatim response body (archived even on failure when present). */ + body?: string; + error?: string; +} + +export async function fetchAccounts( + accessUrl: string, + options: { startDate: Date; pending?: boolean }, + fetchFn: typeof fetch = fetch +): Promise { + const { baseUrl, authHeader } = parseAccessUrl(accessUrl); + const url = new URL(`${baseUrl}/accounts`); + url.searchParams.set('start-date', String(Math.floor(options.startDate.getTime() / 1000))); + if (options.pending !== false) url.searchParams.set('pending', '1'); + + try { + const res = await fetchFn(url, { headers: { Authorization: authHeader } }); + const body = await res.text(); + if (!res.ok) { + return { ok: false, status: res.status, body, error: `HTTP ${res.status}` }; + } + return { ok: true, status: res.status, body }; + } catch (err) { + return { ok: false, error: err instanceof Error ? err.message : String(err) }; + } +} diff --git a/src/routes/(app)/settings/+page.server.ts b/src/routes/(app)/settings/+page.server.ts new file mode 100644 index 0000000..ca6b093 --- /dev/null +++ b/src/routes/(app)/settings/+page.server.ts @@ -0,0 +1,53 @@ +import { fail } from '@sveltejs/kit'; +import { getDb } from '$lib/server/db'; +import { claimSetupToken, listConnections } from '$lib/server/services/connections'; +import { getLastSync } from '$lib/server/services/sync-status'; +import { runSync } from '$lib/server/services/sync'; +import { ClaimError } from '$lib/server/simplefin'; +import type { Actions, PageServerLoad } from './$types'; + +export const load: PageServerLoad = () => { + const db = getDb(); + const connections = listConnections(db).map(({ id, claimedAt }) => ({ id, claimedAt })); + return { + connections, + lastSync: getLastSync(db) + }; +}; + +export const actions: Actions = { + claim: async ({ request }) => { + const form = await request.formData(); + const token = String(form.get('token') ?? '').trim(); + if (!token) return fail(400, { claimMessage: 'Paste a setup token first.' }); + + const db = getDb(); + try { + await claimSetupToken(db, token); + } catch (err) { + const message = + err instanceof ClaimError ? err.message : 'Claiming failed. Check the token and try again.'; + return fail(400, { claimMessage: message }); + } + + // Initial backfill sync, immediately. + const outcomes = await runSync(db); + const failed = outcomes.find((o) => !o.ok); + return { + claimMessage: failed + ? `Connected, but the first sync failed: ${failed.error}` + : 'Connected. First sync complete.' + }; + }, + + sync: async () => { + const outcomes = await runSync(getDb()); + if (outcomes.length === 0) { + return fail(400, { syncMessage: 'No SimpleFIN connection yet.' }); + } + const failed = outcomes.find((o) => !o.ok); + if (failed) return fail(502, { syncMessage: `Sync failed: ${failed.error}` }); + const total = outcomes.reduce((sum, o) => sum + o.newTransactions, 0); + return { syncMessage: `Synced. ${total} new transaction${total === 1 ? '' : 's'}.` }; + } +}; diff --git a/src/routes/(app)/settings/+page.svelte b/src/routes/(app)/settings/+page.svelte index a92c081..5968b01 100644 --- a/src/routes/(app)/settings/+page.svelte +++ b/src/routes/(app)/settings/+page.svelte @@ -1,9 +1,117 @@ + +

Settings

-

SimpleFIN connection setup arrives with the sync engine.

+ +
+

SimpleFIN

+ + {#if data.connections.length === 0} +

+ Connect your banks at SimpleFIN Bridge, + then generate a setup token and paste it here. The token is single-use; Quantum exchanges it + for a permanent read-only connection. +

+
+ + + +
+ {:else} +

+ Connected since {dateFmt.format(new Date(data.connections[0].claimedAt))}. Add or fix bank + connections at + SimpleFIN Bridge; new accounts + appear here automatically after the next sync. +

+ +
+
{ + syncing = true; + return async ({ update }) => { + await update(); + syncing = false; + }; + }} + > + +
+ {#if data.lastSync} + + Last sync {dateFmt.format(new Date(data.lastSync.fetchedAt))} + {data.lastSync.ok ? '' : `, failed: ${data.lastSync.error}`} + + {/if} +
+ {/if} + + {#if form?.claimMessage} + + {/if} + {#if form?.syncMessage} + + {/if} +