phase-3: ingestion, from a QR in the camera to a card in the bandeja
The whole pipeline: storage behind one driver interface (local disk and S3), a portable job queue with a poller, QR and CDC parsing, OCR through the Anthropic API, dedupe, manual entry, and the bandeja that turns all of it into one decision per card. Scanning tries the trustworthy door first: a QR is parsed and prefilled from its CDC; a photo without one is queued for OCR when a key is configured and otherwise opens the manual form against the stored file. A QR that will not parse records an ingest error and falls through rather than losing the photo. Job claiming is the only dialect divergence, as SPEC allows: FOR UPDATE SKIP LOCKED on Postgres, a conditional UPDATE against SQLite's single writer. Retry backoff follows SPEC exactly and a job abandoned by a killed process returns to the queue once its lock goes stale, which is the phase 3 acceptance case. OCR uses structured outputs rather than parsing prose, so the model cannot return anything but the RULES.md schema, and every field is nullable because unreadable is a real answer. Web: scan with live QR decoding (BarcodeDetector, ZXing fallback, wasm served from our own origin), manual entry with the IVA split worked out from the total, the bandeja with swipe, buttons and keyboard all doing the same thing, and the documents list and detail with an editable classification. The seed now carries Maria's 34 purchases and 8 sales and Carlos's 6, all classified through the real rules, plus the two open ingest errors. Two defects found and fixed with tests: seeded documents could be dated in the future, which would corrupt any projection computed from them, and the category buttons announced their keyboard shortcut as part of their name. 255 vitest tests, 37 Playwright tests, rules coverage still 100%, typecheck and lint clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
0d7651b17c
commit
b074456b70
@@ -9,9 +9,13 @@
|
||||
"start": "node dist/index.js",
|
||||
"typecheck": "tsc --noEmit",
|
||||
"db:migrate": "tsx src/db/migrate.cli.ts",
|
||||
"db:seed": "tsx src/db/seed.cli.ts"
|
||||
"db:seed": "tsx src/db/seed.cli.ts",
|
||||
"fixtures": "tsx scripts/render-fixtures.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@anthropic-ai/sdk": "^0.123.0",
|
||||
"@aws-sdk/client-s3": "^3.1126.0",
|
||||
"@aws-sdk/s3-request-presigner": "^3.1126.0",
|
||||
"@hono/node-server": "^2.1.1",
|
||||
"@impuestos/contracts": "workspace:*",
|
||||
"@impuestos/i18n": "workspace:*",
|
||||
@@ -28,6 +32,10 @@
|
||||
"@types/better-sqlite3": "^9.6.0",
|
||||
"@types/node": "^26.4.1",
|
||||
"@types/pg": "^8.23.1",
|
||||
"@types/pngjs": "^6.0.5",
|
||||
"@types/qrcode": "^1.5.6",
|
||||
"pngjs": "^7.0.0",
|
||||
"qrcode": "^1.5.4",
|
||||
"tsup": "^8.5.1",
|
||||
"tsx": "^4.23.13",
|
||||
"typescript": "^5.9.3"
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
import { PNG } from 'pngjs';
|
||||
import QRCode from 'qrcode';
|
||||
import type { Fixture } from '../src/db/fixtures';
|
||||
|
||||
const WIDTH = 620;
|
||||
const HEIGHT = 840;
|
||||
const QR_MODULE_PX = 6;
|
||||
const QUIET_ZONE_MODULES = 4;
|
||||
|
||||
/**
|
||||
* Draws a fixture as a PNG "page" with a real, decodable QR where a KUDE prints one.
|
||||
*
|
||||
* SPEC-GAP: CONTRACTS.md section 4 describes rendering these through HTML or canvas. They
|
||||
* are drawn with raw pixels instead, because the only thing any test reads from them is
|
||||
* the QR, and a faithful visual replica would mean carrying a browser or an SVG
|
||||
* rasteriser purely to produce three images nothing looks at. The emitter, date and
|
||||
* totals live in fixtures.ts, which is where tests and the seed read them from.
|
||||
*/
|
||||
export function renderKude(fixture: Fixture): Buffer {
|
||||
const png = new PNG({ width: WIDTH, height: HEIGHT });
|
||||
fill(png, 0xff, 0xff, 0xff);
|
||||
|
||||
// A masthead bar and a rule, so the image reads as a document rather than a blank page.
|
||||
rect(png, 0, 0, WIDTH, 72, 0x0f, 0x4c, 0x4c);
|
||||
rect(png, 40, 120, WIDTH - 80, 2, 0xd0, 0xd0, 0xd0);
|
||||
rect(png, 40, 220, WIDTH - 80, 2, 0xd0, 0xd0, 0xd0);
|
||||
|
||||
if (fixture.qrUrl) drawQr(png, fixture.qrUrl);
|
||||
|
||||
return PNG.sync.write(png);
|
||||
}
|
||||
|
||||
function drawQr(png: PNG, payload: string): void {
|
||||
// Level M with an explicit quiet zone: a QR without margin is unreliable to decode.
|
||||
const qr = QRCode.create(payload, { errorCorrectionLevel: 'M' });
|
||||
const size = qr.modules.size;
|
||||
const data = qr.modules.data;
|
||||
|
||||
const side = (size + QUIET_ZONE_MODULES * 2) * QR_MODULE_PX;
|
||||
const originX = Math.floor((WIDTH - side) / 2);
|
||||
const originY = HEIGHT - side - 60;
|
||||
|
||||
rect(png, originX, originY, side, side, 0xff, 0xff, 0xff);
|
||||
|
||||
for (let row = 0; row < size; row++) {
|
||||
for (let column = 0; column < size; column++) {
|
||||
if (!data[row * size + column]) continue;
|
||||
rect(
|
||||
png,
|
||||
originX + (column + QUIET_ZONE_MODULES) * QR_MODULE_PX,
|
||||
originY + (row + QUIET_ZONE_MODULES) * QR_MODULE_PX,
|
||||
QR_MODULE_PX,
|
||||
QR_MODULE_PX,
|
||||
0x00,
|
||||
0x00,
|
||||
0x00,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function fill(png: PNG, r: number, g: number, b: number): void {
|
||||
rect(png, 0, 0, png.width, png.height, r, g, b);
|
||||
}
|
||||
|
||||
function rect(
|
||||
png: PNG,
|
||||
x: number,
|
||||
y: number,
|
||||
width: number,
|
||||
height: number,
|
||||
r: number,
|
||||
g: number,
|
||||
b: number,
|
||||
): void {
|
||||
for (let row = y; row < y + height && row < png.height; row++) {
|
||||
for (let column = x; column < x + width && column < png.width; column++) {
|
||||
const index = (png.width * row + column) << 2;
|
||||
png.data[index] = r;
|
||||
png.data[index + 1] = g;
|
||||
png.data[index + 2] = b;
|
||||
png.data[index + 3] = 0xff;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
/**
|
||||
* Renders the fixture comprobantes to PNG under /fixtures, so the e2e tests scan a real
|
||||
* image rather than a hand built payload. Run with `pnpm fixtures`.
|
||||
*/
|
||||
import { mkdirSync, writeFileSync } from 'node:fs';
|
||||
import { dirname, join } from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { FIXTURES } from '../src/db/fixtures';
|
||||
import { renderKude } from './kude-png';
|
||||
|
||||
const here = dirname(fileURLToPath(import.meta.url));
|
||||
const outputDir = join(here, '..', '..', '..', 'fixtures');
|
||||
|
||||
mkdirSync(outputDir, { recursive: true });
|
||||
|
||||
for (const fixture of FIXTURES) {
|
||||
const png = renderKude(fixture);
|
||||
const target = join(outputDir, `${fixture.name}.png`);
|
||||
writeFileSync(target, png);
|
||||
console.info(`[fixtures] ${target} (${png.byteLength} bytes)${fixture.cdc ? ' with QR' : ' no QR'}`);
|
||||
}
|
||||
|
||||
console.info(`[fixtures] ${FIXTURES.length} rendered`);
|
||||
@@ -0,0 +1,114 @@
|
||||
import { buildCdc, computeRucDv } from '@impuestos/rules';
|
||||
|
||||
/**
|
||||
* The synthetic comprobantes the fixtures script renders and the tests scan.
|
||||
*
|
||||
* Every CDC is assembled from the field table in RULES.md section 4 rather than written
|
||||
* out as a string, so a change to that table breaks these loudly instead of silently
|
||||
* producing a CDC that no longer parses.
|
||||
*/
|
||||
export interface Fixture {
|
||||
name: string;
|
||||
emitterName: string;
|
||||
emitterRucBase: string;
|
||||
issueDate: string; // YYYY-MM-DD
|
||||
total: number;
|
||||
iva10: number;
|
||||
/** Absent for the fixture that deliberately carries no QR. */
|
||||
cdc?: string;
|
||||
qrUrl?: string;
|
||||
}
|
||||
|
||||
const QR_HOST = 'https://ekuatia.set.gov.py/consultas/qr';
|
||||
|
||||
function fixture(args: {
|
||||
name: string;
|
||||
emitterName: string;
|
||||
emitterRucBase: string;
|
||||
issueDate: string;
|
||||
total: number;
|
||||
iva10: number;
|
||||
numeroDocumento: string;
|
||||
codigoSeguridad: string;
|
||||
withQr: boolean;
|
||||
}): Fixture {
|
||||
const base = args.emitterRucBase.padStart(8, '0');
|
||||
const dv = String(computeRucDv(args.emitterRucBase));
|
||||
|
||||
if (!args.withQr) {
|
||||
return {
|
||||
name: args.name,
|
||||
emitterName: args.emitterName,
|
||||
emitterRucBase: args.emitterRucBase,
|
||||
issueDate: args.issueDate,
|
||||
total: args.total,
|
||||
iva10: args.iva10,
|
||||
};
|
||||
}
|
||||
|
||||
const cdc = buildCdc({
|
||||
tipoDocumento: '01',
|
||||
rucEmisor: base,
|
||||
dvEmisor: dv,
|
||||
establecimiento: '001',
|
||||
puntoExpedicion: '001',
|
||||
numeroDocumento: args.numeroDocumento,
|
||||
tipoContribuyente: '1',
|
||||
fechaEmision: args.issueDate.replaceAll('-', ''),
|
||||
tipoEmision: '1',
|
||||
codigoSeguridad: args.codigoSeguridad,
|
||||
digitoVerificador: '4',
|
||||
});
|
||||
|
||||
const qrUrl =
|
||||
`${QR_HOST}?nVersion=150&Id=${cdc}&dFeEmiDE=${args.issueDate}` +
|
||||
`&dTotGralOpe=${args.total}&dTotIVA=${args.iva10}&cItems=3`;
|
||||
|
||||
return {
|
||||
name: args.name,
|
||||
emitterName: args.emitterName,
|
||||
emitterRucBase: args.emitterRucBase,
|
||||
issueDate: args.issueDate,
|
||||
total: args.total,
|
||||
iva10: args.iva10,
|
||||
cdc,
|
||||
qrUrl,
|
||||
};
|
||||
}
|
||||
|
||||
/** Two with a scannable QR, one without, per CONTRACTS.md section 4. */
|
||||
export const FIXTURES: readonly Fixture[] = [
|
||||
fixture({
|
||||
name: 'factura-qr-supermercado',
|
||||
emitterName: 'SUPERMERCADO REAL',
|
||||
emitterRucBase: '80011223',
|
||||
issueDate: '2026-08-14',
|
||||
total: 385_000,
|
||||
iva10: 35_000,
|
||||
numeroDocumento: '0004521',
|
||||
codigoSeguridad: '481920573',
|
||||
withQr: true,
|
||||
}),
|
||||
fixture({
|
||||
name: 'factura-qr-farmacia',
|
||||
emitterName: 'FARMACIA CATEDRAL',
|
||||
emitterRucBase: '80022114',
|
||||
issueDate: '2026-08-19',
|
||||
total: 176_000,
|
||||
iva10: 16_000,
|
||||
numeroDocumento: '0009087',
|
||||
codigoSeguridad: '736451028',
|
||||
withQr: true,
|
||||
}),
|
||||
fixture({
|
||||
name: 'factura-sin-qr-ferreteria',
|
||||
emitterName: 'FERRETERIA SAN MIGUEL',
|
||||
emitterRucBase: '80033005',
|
||||
issueDate: '2026-08-22',
|
||||
total: 540_000,
|
||||
iva10: 49_091,
|
||||
numeroDocumento: '0001777',
|
||||
codigoSeguridad: '905172634',
|
||||
withQr: false,
|
||||
}),
|
||||
];
|
||||
@@ -0,0 +1,32 @@
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { seedDate } from './seed-documents';
|
||||
|
||||
describe('seedDate', () => {
|
||||
it('anchors to the month it is asked for', () => {
|
||||
const now = new Date('2026-09-04T12:00:00Z');
|
||||
expect(seedDate(now, 1, 14)).toBe('2026-08-14');
|
||||
expect(seedDate(now, 3, 2)).toBe('2026-06-02');
|
||||
});
|
||||
|
||||
// A comprobante dated after today belongs to a period that has not happened.
|
||||
it('never produces a future date in the current month', () => {
|
||||
const now = new Date('2026-09-04T12:00:00Z');
|
||||
expect(seedDate(now, 0, 8)).toBe('2026-09-04');
|
||||
expect(seedDate(now, 0, 3)).toBe('2026-09-03');
|
||||
});
|
||||
|
||||
it('clamps to the length of a short month', () => {
|
||||
const now = new Date('2026-03-15T12:00:00Z');
|
||||
expect(seedDate(now, 1, 31)).toBe('2026-02-28');
|
||||
});
|
||||
|
||||
it('crosses the year boundary', () => {
|
||||
const now = new Date('2026-01-20T12:00:00Z');
|
||||
expect(seedDate(now, 1, 10)).toBe('2025-12-10');
|
||||
});
|
||||
|
||||
it('never emits a day before the first', () => {
|
||||
const now = new Date('2026-09-01T12:00:00Z');
|
||||
expect(seedDate(now, 0, 8)).toBe('2026-09-01');
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,252 @@
|
||||
import { computeRucDv } from '@impuestos/rules';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import { dedupeHash } from '../modules/documents/dedupe';
|
||||
import { classifyDocument } from '../modules/documents/service';
|
||||
import { recordIngestError } from '../modules/ingest/errors';
|
||||
import type { DbHandle } from './index';
|
||||
import type { DocumentsTable } from './schema';
|
||||
|
||||
/**
|
||||
* Maria's year of comprobantes, per CONTRACTS.md section 4. Deterministic: the amounts,
|
||||
* dates and emitters are written down rather than generated, so every assertion about
|
||||
* her dashboard and her declarations is a fixed number somebody can check by hand.
|
||||
*
|
||||
* Dates are anchored to the seed's "today" so the previous month is always complete.
|
||||
*/
|
||||
|
||||
const SEEDED_AT = '2026-01-15T12:00:00.000Z';
|
||||
|
||||
interface SeedDoc {
|
||||
emitter: string;
|
||||
rucBase: string;
|
||||
/** Day of the month. The month comes from the offset the caller passes. */
|
||||
day: number;
|
||||
monthsAgo: number;
|
||||
total: number;
|
||||
rate: 10 | 5 | 0;
|
||||
direction?: 'purchase' | 'sale';
|
||||
regime?: 'normal' | 'resimple' | 'unknown';
|
||||
status?: 'confirmed' | 'needs_review';
|
||||
source?: 'scan_qr' | 'scan_ocr' | 'manual';
|
||||
}
|
||||
|
||||
/**
|
||||
* 34 purchases and 8 sales for Maria, 29 of the purchases confirmed and 5 waiting in the
|
||||
* bandeja, exactly as CONTRACTS.md section 4 specifies.
|
||||
*/
|
||||
|
||||
/** IVA is included in the printed total, so the base is the total less the tax. */
|
||||
function split(total: number, rate: 10 | 5 | 0): {
|
||||
base10: number;
|
||||
base5: number;
|
||||
exenta: number;
|
||||
iva10: number;
|
||||
iva5: number;
|
||||
} {
|
||||
if (rate === 0) return { base10: 0, base5: 0, exenta: total, iva10: 0, iva5: 0 };
|
||||
if (rate === 5) {
|
||||
const iva5 = Math.round(total / 21);
|
||||
return { base10: 0, base5: total - iva5, exenta: 0, iva10: 0, iva5 };
|
||||
}
|
||||
const iva10 = Math.round(total / 11);
|
||||
return { base10: total - iva10, base5: 0, exenta: 0, iva10, iva5: 0 };
|
||||
}
|
||||
|
||||
const PURCHASES: SeedDoc[] = [
|
||||
// Groceries, six of them across the year.
|
||||
{ emitter: 'SUPERMERCADO REAL', rucBase: '80011223', day: 4, monthsAgo: 1, total: 412_000, rate: 10 },
|
||||
{ emitter: 'SUPERMERCADO REAL', rucBase: '80011223', day: 12, monthsAgo: 1, total: 268_000, rate: 10 },
|
||||
{ emitter: 'SUPERMERCADO REAL', rucBase: '80011223', day: 21, monthsAgo: 1, total: 735_000, rate: 10 },
|
||||
{ emitter: 'SUPERMERCADO REAL', rucBase: '80011223', day: 8, monthsAgo: 2, total: 189_000, rate: 10 },
|
||||
{ emitter: 'SUPERMERCADO REAL', rucBase: '80011223', day: 19, monthsAgo: 3, total: 903_000, rate: 10 },
|
||||
{ emitter: 'SUPERMERCADO REAL', rucBase: '80011223', day: 26, monthsAgo: 4, total: 544_000, rate: 10 },
|
||||
|
||||
{ emitter: 'FARMACIA CATEDRAL', rucBase: '80022114', day: 6, monthsAgo: 1, total: 176_000, rate: 10 },
|
||||
{ emitter: 'FARMACIA CATEDRAL', rucBase: '80022114', day: 15, monthsAgo: 2, total: 321_000, rate: 10 },
|
||||
{ emitter: 'FARMACIA CATEDRAL', rucBase: '80022114', day: 3, monthsAgo: 3, total: 88_000, rate: 10 },
|
||||
{ emitter: 'FARMACIA CATEDRAL', rucBase: '80022114', day: 22, monthsAgo: 5, total: 245_000, rate: 10 },
|
||||
|
||||
{ emitter: 'COLEGIO SAN JOSE', rucBase: '80033441', day: 5, monthsAgo: 1, total: 1_650_000, rate: 0 },
|
||||
{ emitter: 'COLEGIO SAN JOSE', rucBase: '80033441', day: 5, monthsAgo: 2, total: 1_650_000, rate: 0 },
|
||||
{ emitter: 'COLEGIO SAN JOSE', rucBase: '80033441', day: 5, monthsAgo: 3, total: 1_650_000, rate: 0 },
|
||||
|
||||
{ emitter: 'PETROBRAS ESTACION 12', rucBase: '80044552', day: 9, monthsAgo: 1, total: 350_000, rate: 10 },
|
||||
{ emitter: 'PETROBRAS ESTACION 12', rucBase: '80044552', day: 24, monthsAgo: 2, total: 420_000, rate: 10 },
|
||||
{ emitter: 'PETROBRAS ESTACION 12', rucBase: '80044552', day: 17, monthsAgo: 4, total: 305_000, rate: 10 },
|
||||
|
||||
// Rent, at the 5% rate.
|
||||
{ emitter: 'INMOBILIARIA DEL SOL', rucBase: '80055663', day: 2, monthsAgo: 1, total: 2_800_000, rate: 5 },
|
||||
{ emitter: 'INMOBILIARIA DEL SOL', rucBase: '80055663', day: 2, monthsAgo: 2, total: 2_800_000, rate: 5 },
|
||||
|
||||
{ emitter: 'BOUTIQUE ANDREA', rucBase: '80066774', day: 14, monthsAgo: 2, total: 480_000, rate: 10 },
|
||||
{ emitter: 'BOUTIQUE ANDREA', rucBase: '80066774', day: 28, monthsAgo: 5, total: 615_000, rate: 10 },
|
||||
|
||||
{ emitter: 'CINE ITAU', rucBase: '80077885', day: 16, monthsAgo: 1, total: 95_000, rate: 10 },
|
||||
{ emitter: 'CINE ITAU', rucBase: '80077885', day: 23, monthsAgo: 3, total: 120_000, rate: 10 },
|
||||
|
||||
// RESIMPLE supplier: deductible for IRP but capped, and no IVA credit at all.
|
||||
{ emitter: 'DESPENSA DON JUAN', rucBase: '20033445', day: 11, monthsAgo: 1, total: 145_000, rate: 10, regime: 'resimple' },
|
||||
{ emitter: 'DESPENSA DON JUAN', rucBase: '20033445', day: 27, monthsAgo: 2, total: 210_000, rate: 10, regime: 'resimple' },
|
||||
|
||||
// The rest: a mix, including some with no keyword match at all.
|
||||
{ emitter: 'FERRETERIA SAN MIGUEL', rucBase: '80033005', day: 22, monthsAgo: 1, total: 540_000, rate: 10 },
|
||||
{ emitter: 'CLINICA SANTA RITA', rucBase: '80010118', day: 13, monthsAgo: 3, total: 890_000, rate: 10 },
|
||||
{ emitter: 'SERVICIOS INTEGRALES SA', rucBase: '80011229', day: 20, monthsAgo: 4, total: 375_000, rate: 10 },
|
||||
{ emitter: 'ANDE', rucBase: '80000001', day: 10, monthsAgo: 1, total: 318_000, rate: 10 },
|
||||
{ emitter: 'ESSAP', rucBase: '80000002', day: 10, monthsAgo: 1, total: 96_000, rate: 10 },
|
||||
|
||||
// Waiting in the bandeja, two of them from a low confidence OCR read.
|
||||
{ emitter: 'PANADERIA LA UNION', rucBase: '80012330', day: 3, monthsAgo: 0, total: 62_000, rate: 10, status: 'needs_review' },
|
||||
{ emitter: 'TALLER MECANICO RUIZ', rucBase: '80012331', day: 5, monthsAgo: 0, total: 780_000, rate: 10, status: 'needs_review' },
|
||||
{ emitter: 'CONSULTORA NORTE', rucBase: '80012332', day: 6, monthsAgo: 0, total: 1_200_000, rate: 10, status: 'needs_review' },
|
||||
{ emitter: 'ALMACEN CENTRAL', rucBase: '80012333', day: 7, monthsAgo: 0, total: 143_000, rate: 10, status: 'needs_review', source: 'scan_ocr' },
|
||||
{ emitter: 'OPTICA VISION', rucBase: '80012334', day: 8, monthsAgo: 0, total: 455_000, rate: 10, status: 'needs_review', source: 'scan_ocr' },
|
||||
];
|
||||
|
||||
/** Services she invoices, all at 10%. */
|
||||
const SALES: SeedDoc[] = [
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 3, monthsAgo: 1, total: 6_600_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 17, monthsAgo: 1, total: 4_400_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 5, monthsAgo: 2, total: 9_900_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 20, monthsAgo: 2, total: 5_500_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 9, monthsAgo: 3, total: 7_700_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 22, monthsAgo: 3, total: 4_400_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 11, monthsAgo: 4, total: 8_800_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'MARIA GONZALEZ', rucBase: '4123456', day: 25, monthsAgo: 5, total: 6_050_000, rate: 10, direction: 'sale' },
|
||||
];
|
||||
|
||||
const CARLOS_DOCS: SeedDoc[] = [
|
||||
{ emitter: 'DISTRIBUIDORA DEL ESTE', rucBase: '80020001', day: 4, monthsAgo: 1, total: 4_400_000, rate: 10 },
|
||||
{ emitter: 'IMPRENTA MODERNA', rucBase: '80020002', day: 12, monthsAgo: 1, total: 1_320_000, rate: 10 },
|
||||
{ emitter: 'FERRETERIA SAN MIGUEL', rucBase: '80033005', day: 19, monthsAgo: 2, total: 880_000, rate: 10 },
|
||||
{ emitter: 'BENITEZ Y ASOCIADOS SRL', rucBase: '80012345', day: 6, monthsAgo: 1, total: 12_100_000, rate: 10, direction: 'sale' },
|
||||
{ emitter: 'CONSULTORA SUR', rucBase: '80020003', day: 8, monthsAgo: 0, total: 660_000, rate: 10, status: 'needs_review' },
|
||||
{ emitter: 'PAPELERIA CENTRAL', rucBase: '80020004', day: 9, monthsAgo: 0, total: 231_000, rate: 10, status: 'needs_review' },
|
||||
];
|
||||
|
||||
/**
|
||||
* Anchors every seeded date to the same "now", so the previous month is always complete.
|
||||
*
|
||||
* Never returns a future date: the current month is only partly elapsed, and a comprobante
|
||||
* dated after today would land in a period that has not happened, quietly corrupting every
|
||||
* projection and deadline computed from it.
|
||||
*/
|
||||
export function seedDate(now: Date, monthsAgo: number, day: number): string {
|
||||
const anchor = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth() - monthsAgo, 1));
|
||||
const lastOfMonth = new Date(
|
||||
Date.UTC(anchor.getUTCFullYear(), anchor.getUTCMonth() + 1, 0),
|
||||
).getUTCDate();
|
||||
|
||||
const latestAllowed = monthsAgo === 0 ? Math.min(lastOfMonth, now.getUTCDate()) : lastOfMonth;
|
||||
const safeDay = Math.max(1, Math.min(day, latestAllowed));
|
||||
|
||||
return `${anchor.getUTCFullYear()}-${String(anchor.getUTCMonth() + 1).padStart(2, '0')}-${String(safeDay).padStart(2, '0')}`;
|
||||
}
|
||||
|
||||
export async function seedDocuments(handle: DbHandle, now = new Date()): Promise<number> {
|
||||
const db = handle.db;
|
||||
|
||||
const already = await db.selectFrom('documents').select('id').limit(1).executeTakeFirst();
|
||||
if (already) return 0;
|
||||
|
||||
const maria = await userIdFor(handle, 'maria@demo.local');
|
||||
const carlos = await userIdFor(handle, 'carlos@demo.local');
|
||||
|
||||
let inserted = 0;
|
||||
for (const [userId, docs] of [
|
||||
[maria, [...PURCHASES, ...SALES]],
|
||||
[carlos, CARLOS_DOCS],
|
||||
] as const) {
|
||||
for (const doc of docs) {
|
||||
await insertDoc(handle, userId, doc, now);
|
||||
inserted += 1;
|
||||
}
|
||||
}
|
||||
|
||||
// Two open ingest errors against Maria, per CONTRACTS.md section 4.
|
||||
await recordIngestError(db, {
|
||||
userId: maria,
|
||||
stage: 'ocr',
|
||||
message: 'la foto salio movida y no se pudieron leer los importes',
|
||||
payload: { seeded: true },
|
||||
});
|
||||
await recordIngestError(db, {
|
||||
userId: maria,
|
||||
stage: 'qr_parse',
|
||||
message: 'el QR de la factura no tenia un CDC valido',
|
||||
payload: { seeded: true },
|
||||
});
|
||||
|
||||
return inserted;
|
||||
}
|
||||
|
||||
async function insertDoc(
|
||||
handle: DbHandle,
|
||||
userId: string,
|
||||
doc: SeedDoc,
|
||||
now: Date,
|
||||
): Promise<void> {
|
||||
const issueDate = seedDate(now, doc.monthsAgo, doc.day);
|
||||
const amounts = split(doc.total, doc.rate);
|
||||
const direction = doc.direction ?? 'purchase';
|
||||
const status = doc.status ?? 'confirmed';
|
||||
const id = uuidv7();
|
||||
|
||||
const row: DocumentsTable = {
|
||||
id,
|
||||
user_id: userId,
|
||||
source: doc.source ?? 'scan_qr',
|
||||
status,
|
||||
cdc: null,
|
||||
qr_url: null,
|
||||
doc_kind: 'factura',
|
||||
direction,
|
||||
emitter_ruc: doc.rucBase,
|
||||
emitter_dv: String(computeRucDv(doc.rucBase)),
|
||||
emitter_name: doc.emitter,
|
||||
receiver_doc: null,
|
||||
issue_date: issueDate,
|
||||
currency: 'PYG',
|
||||
total: doc.total,
|
||||
amount_iva10: amounts.base10,
|
||||
amount_iva5: amounts.base5,
|
||||
amount_exenta: amounts.exenta,
|
||||
iva10: amounts.iva10,
|
||||
iva5: amounts.iva5,
|
||||
supplier_regime_hint: doc.regime ?? 'normal',
|
||||
verified_dnit: 0,
|
||||
verification_status: 'unverified',
|
||||
dedupe_hash: dedupeHash({
|
||||
emitterRuc: doc.rucBase,
|
||||
issueDate,
|
||||
total: doc.total,
|
||||
docKind: 'factura',
|
||||
}),
|
||||
file_id: null,
|
||||
raw_extraction: null,
|
||||
created_at: SEEDED_AT,
|
||||
confirmed_at: status === 'confirmed' ? SEEDED_AT : null,
|
||||
};
|
||||
|
||||
await handle.db.insertInto('documents').values(row).execute();
|
||||
await classifyDocument({ db: handle.db, storage: nullStorage }, id);
|
||||
}
|
||||
|
||||
/** The seed never reads a file, so classification does not need a real driver. */
|
||||
const nullStorage = {
|
||||
kind: 'local' as const,
|
||||
put: async () => undefined,
|
||||
getStream: async () => new ReadableStream<Uint8Array>(),
|
||||
delete: async () => undefined,
|
||||
url: async () => '',
|
||||
check: async () => undefined,
|
||||
};
|
||||
|
||||
async function userIdFor(handle: DbHandle, email: string): Promise<string> {
|
||||
const row = await handle.db
|
||||
.selectFrom('user')
|
||||
.select('id')
|
||||
.where('email', '=', email)
|
||||
.executeTakeFirstOrThrow();
|
||||
return row.id;
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import type { Auth, Role } from '../auth/options';
|
||||
import { writeAudit } from '../modules/audit';
|
||||
import { CONSENT_TEXT_VERSION } from '../modules/pii';
|
||||
import type { DbHandle } from './index';
|
||||
import { seedDocuments } from './seed-documents';
|
||||
|
||||
export interface SeedAccount {
|
||||
email: string;
|
||||
@@ -71,6 +72,7 @@ export async function seed(handle: DbHandle, auth: Auth): Promise<SeedResult> {
|
||||
}
|
||||
|
||||
await seedProfiles(handle);
|
||||
await seedDocuments(handle);
|
||||
await seedAuditTrail(handle);
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -23,7 +23,7 @@ describe('health endpoints', () => {
|
||||
expect(response.status).toBe(200);
|
||||
expect(await response.json()).toEqual({
|
||||
ok: true,
|
||||
checks: { database: 'ok', migrations: 'ok' },
|
||||
checks: { database: 'ok', migrations: 'ok', storage: 'ok' },
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -3,6 +3,8 @@ import type { AppDeps, AppEnv } from './context';
|
||||
import { HttpError, toEnvelope } from './errors';
|
||||
import { liveness, readiness } from './health';
|
||||
import { localeMiddleware, sessionMiddleware } from './middleware';
|
||||
import { documentRoutes } from './routes/documents';
|
||||
import { fileRoutes } from './routes/files';
|
||||
import { lookupRoutes } from './routes/lookup';
|
||||
import { meRoutes } from './routes/me';
|
||||
|
||||
@@ -48,6 +50,8 @@ export function createApp(deps: AppDeps): AppHandle {
|
||||
api.use('*', localeMiddleware(deps));
|
||||
api.route('/lookup', lookupRoutes());
|
||||
api.route('/me', meRoutes(deps));
|
||||
api.route('/documents', documentRoutes(deps));
|
||||
api.route('/files', fileRoutes(deps));
|
||||
|
||||
app.route('/api', api);
|
||||
|
||||
|
||||
@@ -2,6 +2,8 @@ import type { Locale } from '@impuestos/i18n';
|
||||
import type { Auth } from '../auth/options';
|
||||
import type { DbHandle } from '../db/index';
|
||||
import type { Env } from '../lib/env';
|
||||
import type { IngestDeps } from '../modules/ingest/pipeline';
|
||||
import type { StorageDriver } from '../modules/storage';
|
||||
|
||||
export interface SessionUser {
|
||||
id: string;
|
||||
@@ -13,6 +15,9 @@ export interface AppDeps {
|
||||
env: Env;
|
||||
handle: DbHandle;
|
||||
auth: Auth;
|
||||
storage: StorageDriver;
|
||||
/** Bundles the db and storage the documents and ingest modules need. */
|
||||
ingest: IngestDeps;
|
||||
}
|
||||
|
||||
/** Hono context typing shared by every route and middleware. */
|
||||
|
||||
@@ -42,5 +42,12 @@ export async function readiness(
|
||||
checks['migrations'] = 'error';
|
||||
}
|
||||
|
||||
try {
|
||||
await deps.storage.check();
|
||||
checks['storage'] = 'ok';
|
||||
} catch {
|
||||
checks['storage'] = 'error';
|
||||
}
|
||||
|
||||
return { ok: Object.values(checks).every((state) => state === 'ok'), checks };
|
||||
}
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
import {
|
||||
ClassificationPatchInput,
|
||||
DocumentListQuery,
|
||||
DocumentPatchInput,
|
||||
ManualDocumentInput,
|
||||
RejectInput,
|
||||
} from '@impuestos/contracts';
|
||||
import { Hono } from 'hono';
|
||||
import type { z } from 'zod';
|
||||
import {
|
||||
confirmDocument,
|
||||
getDocument,
|
||||
listDocuments,
|
||||
patchClassification,
|
||||
patchDocument,
|
||||
rejectDocument,
|
||||
} from '../../modules/documents/service';
|
||||
import { listIngestErrorsForUser } from '../../modules/ingest/errors';
|
||||
import { ingestManual, ingestScan } from '../../modules/ingest/pipeline';
|
||||
import type { AppDeps, AppEnv } from '../context';
|
||||
import { HttpError } from '../errors';
|
||||
import { requireUser } from '../middleware';
|
||||
|
||||
/** Bigger than any phone photograph, small enough that a bad request cannot exhaust RAM. */
|
||||
const MAX_UPLOAD_BYTES = 12 * 1024 * 1024;
|
||||
|
||||
const ALLOWED_MIME = new Set([
|
||||
'image/jpeg',
|
||||
'image/png',
|
||||
'image/webp',
|
||||
'image/heic',
|
||||
'image/heif',
|
||||
'application/pdf',
|
||||
]);
|
||||
|
||||
export function documentRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
const routes = new Hono<AppEnv>();
|
||||
|
||||
routes.post('/scan', async (c) => {
|
||||
const user = requireUser(c);
|
||||
|
||||
const form = await c.req.formData().catch(() => null);
|
||||
const file = form?.get('file');
|
||||
if (!form || !(file instanceof File)) throw new HttpError('validation_error', { field: 'file' });
|
||||
|
||||
if (file.size > MAX_UPLOAD_BYTES) {
|
||||
throw new HttpError('validation_error', { field: 'file', detail: { maxBytes: MAX_UPLOAD_BYTES } });
|
||||
}
|
||||
const mime = file.type || 'image/jpeg';
|
||||
if (!ALLOWED_MIME.has(mime)) {
|
||||
throw new HttpError('validation_error', { field: 'file', detail: { mime } });
|
||||
}
|
||||
|
||||
const qrValue = form.get('qrPayload');
|
||||
const outcome = await ingestScan(deps.ingest, user.id, {
|
||||
file: { data: new Uint8Array(await file.arrayBuffer()), mime },
|
||||
qrPayload: typeof qrValue === 'string' ? qrValue : undefined,
|
||||
});
|
||||
|
||||
if (outcome.kind === 'needs_manual') return c.json(outcome.result);
|
||||
return c.json(
|
||||
outcome.merged ? { ...outcome.document, merged: true } : outcome.document,
|
||||
outcome.merged ? 200 : 201,
|
||||
);
|
||||
});
|
||||
|
||||
routes.post('/manual', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const input = parse(ManualDocumentInput, await body(c));
|
||||
const { document, merged } = await ingestManual(deps.ingest, user.id, input);
|
||||
return c.json(merged ? { ...document, merged: true } : document, merged ? 200 : 201);
|
||||
});
|
||||
|
||||
routes.get('/', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const query = parse(DocumentListQuery, stripUndefined(c.req.query()));
|
||||
return c.json(await listDocuments(deps.ingest, user.id, query));
|
||||
});
|
||||
|
||||
routes.get('/errors', async (c) => {
|
||||
const user = requireUser(c);
|
||||
return c.json(await listIngestErrorsForUser(deps.handle.db, user.id));
|
||||
});
|
||||
|
||||
routes.get('/:id', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const document = await getDocument(deps.ingest, user.id, c.req.param('id'));
|
||||
if (!document) throw new HttpError('not_found');
|
||||
return c.json(document);
|
||||
});
|
||||
|
||||
routes.patch('/:id', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const patch = parse(DocumentPatchInput, await body(c));
|
||||
const document = await patchDocument(deps.ingest, user.id, c.req.param('id'), patch);
|
||||
if (!document) throw new HttpError('not_found');
|
||||
return c.json(document);
|
||||
});
|
||||
|
||||
routes.post('/:id/confirm', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const document = await confirmDocument(deps.ingest, user.id, c.req.param('id'));
|
||||
if (!document) throw new HttpError('not_found');
|
||||
return c.json(document);
|
||||
});
|
||||
|
||||
routes.post('/:id/reject', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const { reason } = parse(RejectInput, await body(c));
|
||||
const document = await rejectDocument(deps.ingest, user.id, c.req.param('id'), reason);
|
||||
if (!document) throw new HttpError('not_found');
|
||||
return c.json(document);
|
||||
});
|
||||
|
||||
routes.patch('/:id/classification', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const patch = parse(ClassificationPatchInput, await body(c));
|
||||
const document = await patchClassification(deps.ingest, user.id, c.req.param('id'), patch);
|
||||
if (!document) throw new HttpError('not_found');
|
||||
return c.json(document);
|
||||
});
|
||||
|
||||
return routes;
|
||||
}
|
||||
|
||||
async function body(c: { req: { json: () => Promise<unknown> } }): Promise<unknown> {
|
||||
try {
|
||||
return await c.req.json();
|
||||
} catch {
|
||||
throw new HttpError('validation_error');
|
||||
}
|
||||
}
|
||||
|
||||
function stripUndefined(query: Record<string, string | undefined>): Record<string, string> {
|
||||
return Object.fromEntries(
|
||||
Object.entries(query).filter((entry): entry is [string, string] => entry[1] !== undefined),
|
||||
);
|
||||
}
|
||||
|
||||
function parse<T>(schema: z.ZodType<T>, value: unknown): T {
|
||||
const result = schema.safeParse(value);
|
||||
if (!result.success) {
|
||||
const issue = result.error.issues[0];
|
||||
throw new HttpError('validation_error', {
|
||||
...(issue?.path.length ? { field: issue.path.join('.') } : {}),
|
||||
detail: result.error.issues,
|
||||
});
|
||||
}
|
||||
return result.data;
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
import { Hono } from 'hono';
|
||||
import { ADMIN_ROLES } from '../../auth/options';
|
||||
import { writeAudit } from '../../modules/audit';
|
||||
import type { AppDeps, AppEnv } from '../context';
|
||||
import { HttpError } from '../errors';
|
||||
import { requireUser } from '../middleware';
|
||||
|
||||
/**
|
||||
* Files are private. Ownership is proved by joining through the document that references
|
||||
* the file, so a stolen file id is worthless without the session that owns it. Staff may
|
||||
* read a file for support, and every one of those reads is audited (SPEC.md section 7).
|
||||
*/
|
||||
export function fileRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
const routes = new Hono<AppEnv>();
|
||||
|
||||
routes.get('/:id', async (c) => {
|
||||
const user = requireUser(c);
|
||||
const id = c.req.param('id');
|
||||
|
||||
const file = await deps.handle.db
|
||||
.selectFrom('document_files')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.executeTakeFirst();
|
||||
if (!file) throw new HttpError('not_found');
|
||||
|
||||
const owner = await deps.handle.db
|
||||
.selectFrom('documents')
|
||||
.select('user_id')
|
||||
.where('file_id', '=', id)
|
||||
.executeTakeFirst();
|
||||
|
||||
const isOwner = owner?.user_id === user.id;
|
||||
const isStaff = ADMIN_ROLES.includes(user.role as (typeof ADMIN_ROLES)[number]);
|
||||
|
||||
// A file with no document yet belongs to whoever just uploaded it, which the storage
|
||||
// key records; anything else needs the ownership join or a staff role.
|
||||
const isUploader = owner === undefined && file.path.startsWith(`${user.id}/`);
|
||||
if (!isOwner && !isUploader && !isStaff) throw new HttpError('forbidden');
|
||||
|
||||
if (isStaff && !isOwner) {
|
||||
await writeAudit(deps.handle.db, {
|
||||
actorUserId: user.id,
|
||||
actorRole: user.role,
|
||||
subjectUserId: owner?.user_id ?? null,
|
||||
action: 'admin.file_access',
|
||||
resource: `document_files/${id}`,
|
||||
});
|
||||
}
|
||||
|
||||
const stream = await deps.storage.getStream(file.path);
|
||||
return new Response(stream, {
|
||||
headers: {
|
||||
'content-type': file.mime,
|
||||
'content-length': String(file.size),
|
||||
'cache-control': 'private, max-age=300',
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
return routes;
|
||||
}
|
||||
+28
-5
@@ -5,7 +5,10 @@ import { pendingMigrations } from './db/migrator';
|
||||
import { createApp } from './http/app';
|
||||
import type { AppDeps } from './http/context';
|
||||
import { loadEnv } from './lib/env';
|
||||
import { createHandlers, createPoller } from './modules/jobs';
|
||||
import { createOcrProvider } from './modules/documents/ocr';
|
||||
import { createOtpSender } from './modules/notifications/mailer';
|
||||
import { createStorage } from './modules/storage';
|
||||
|
||||
const DRAIN_TIMEOUT_MS = 25_000;
|
||||
|
||||
@@ -18,7 +21,11 @@ const auth = createAuth({
|
||||
sendOtp: createOtpSender(env),
|
||||
});
|
||||
|
||||
const deps: AppDeps = { env, handle, auth };
|
||||
const storage = createStorage(env);
|
||||
const ocr = createOcrProvider(env);
|
||||
const ingest = { db: handle.db, storage, ocrEnabled: ocr !== null };
|
||||
|
||||
const deps: AppDeps = { env, handle, auth, storage, ingest };
|
||||
const { app, startDraining, inFlight } = createApp(deps);
|
||||
|
||||
const pending = await pendingMigrations(handle).catch(() => ['<database unreachable>']);
|
||||
@@ -26,10 +33,25 @@ if (pending.length > 0) {
|
||||
console.warn(`[boot] pending migrations: ${pending.join(', ')}. Run pnpm db:migrate.`);
|
||||
}
|
||||
|
||||
if (env.ROLE === 'worker') {
|
||||
// The worker shares this image and this bootstrap. It serves only the health
|
||||
// endpoints; the job poller is wired in with the jobs module.
|
||||
console.info('[boot] role=worker');
|
||||
/**
|
||||
* One poller per process (SPEC.md section 10). A dedicated worker always runs one; a
|
||||
* server runs one only when JOBS_INLINE is on, which env validation already refuses to
|
||||
* turn off on SQLite.
|
||||
*/
|
||||
const wantsPoller = env.ROLE === 'worker' || env.JOBS_INLINE;
|
||||
const poller = wantsPoller
|
||||
? createPoller({
|
||||
db: handle.db,
|
||||
dialect: handle.dialect,
|
||||
handlers: createHandlers({ db: handle.db, storage, ocr }),
|
||||
intervalMs: env.JOBS_POLL_INTERVAL_MS,
|
||||
staleMinutes: env.JOBS_STALE_MINUTES,
|
||||
})
|
||||
: null;
|
||||
|
||||
if (poller) {
|
||||
poller.start();
|
||||
console.info(`[boot] job poller running (role=${env.ROLE}, ocr=${ocr ? 'on' : 'off'})`);
|
||||
}
|
||||
|
||||
const server = serve({ fetch: app.fetch, port: env.PORT, hostname: '0.0.0.0' }, (info) => {
|
||||
@@ -56,6 +78,7 @@ async function shutdown(signal: string): Promise<void> {
|
||||
if ('closeAllConnections' in server) server.closeAllConnections();
|
||||
}
|
||||
|
||||
await poller?.stop();
|
||||
await handle.close();
|
||||
console.info('[shutdown] done');
|
||||
process.exit(0);
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
import { createHash } from 'node:crypto';
|
||||
|
||||
/**
|
||||
* SPEC-GAP: SPEC.md section 8 requires a `dedupe_hash` but never defines it, and RULES.md
|
||||
* does not either, so this is an implementation choice rather than a tax rule.
|
||||
*
|
||||
* A CDC already identifies a comprobante uniquely and nationally, so when there is one it
|
||||
* is the whole key. Without a CDC (manual entry, OCR of a paper factura) the key is the
|
||||
* tuple that identifies the document to a human: who issued it, when, for how much, and
|
||||
* of what kind. Two photographs of the same factura collapse; two genuinely different
|
||||
* facturas from the same shop on the same day for the same amount would collapse too,
|
||||
* which is why a merge is presented to the user rather than applied silently.
|
||||
*/
|
||||
export function dedupeHash(input: {
|
||||
cdc?: string | null;
|
||||
emitterRuc: string;
|
||||
issueDate: string;
|
||||
total: number;
|
||||
docKind: string;
|
||||
}): string {
|
||||
const material = input.cdc
|
||||
? `cdc:${input.cdc}`
|
||||
: ['manual', input.emitterRuc.trim(), input.issueDate, String(input.total), input.docKind].join('|');
|
||||
|
||||
return createHash('sha256').update(material).digest('hex');
|
||||
}
|
||||
|
||||
export function sha256(data: Uint8Array): string {
|
||||
return createHash('sha256').update(data).digest('hex');
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
import type { ClassificationDto, DocumentDto } from '@impuestos/contracts';
|
||||
import type { ClassificationsTable, DocumentsTable } from '../../db/schema';
|
||||
|
||||
export type DocumentRow = DocumentsTable;
|
||||
export type ClassificationRow = ClassificationsTable;
|
||||
|
||||
export function toClassificationDto(row: ClassificationRow): ClassificationDto {
|
||||
return {
|
||||
ivaCreditEligible: row.iva_credit_eligible === 1,
|
||||
ivaCreditAmount: row.iva_credit_amount,
|
||||
irpCategory: row.irp_category as ClassificationDto['irpCategory'],
|
||||
irpDeductibleAmount: row.irp_deductible_amount,
|
||||
dependentId: row.dependent_id,
|
||||
confidence: row.confidence,
|
||||
decidedBy: row.decided_by,
|
||||
rulesVersion: row.rules_version,
|
||||
reasons: [],
|
||||
};
|
||||
}
|
||||
|
||||
export function toDocumentDto(
|
||||
row: DocumentRow,
|
||||
classification: ClassificationRow | null,
|
||||
fileUrl: string | null,
|
||||
): DocumentDto {
|
||||
const reasons = readReasons(row.raw_extraction);
|
||||
return {
|
||||
id: row.id,
|
||||
source: row.source,
|
||||
status: row.status,
|
||||
cdc: row.cdc,
|
||||
docKind: row.doc_kind,
|
||||
direction: row.direction,
|
||||
emitterRuc: row.emitter_ruc,
|
||||
emitterDv: row.emitter_dv,
|
||||
emitterName: row.emitter_name,
|
||||
receiverDoc: row.receiver_doc,
|
||||
issueDate: row.issue_date,
|
||||
currency: 'PYG',
|
||||
total: row.total,
|
||||
amountIva10: row.amount_iva10,
|
||||
amountIva5: row.amount_iva5,
|
||||
amountExenta: row.amount_exenta,
|
||||
iva10: row.iva10,
|
||||
iva5: row.iva5,
|
||||
supplierRegimeHint: row.supplier_regime_hint,
|
||||
verifiedDnit: row.verified_dnit === 1,
|
||||
verificationStatus: row.verification_status,
|
||||
fileUrl,
|
||||
classification: classification
|
||||
? { ...toClassificationDto(classification), reasons }
|
||||
: null,
|
||||
createdAt: row.created_at,
|
||||
confirmedAt: row.confirmed_at,
|
||||
};
|
||||
}
|
||||
|
||||
/** The classifier's reason codes ride along in raw_extraction rather than in a column. */
|
||||
function readReasons(rawExtraction: string | null): string[] {
|
||||
if (!rawExtraction) return [];
|
||||
try {
|
||||
const parsed = JSON.parse(rawExtraction) as { reasons?: unknown };
|
||||
return Array.isArray(parsed.reasons) ? parsed.reasons.filter((r): r is string => typeof r === 'string') : [];
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,126 @@
|
||||
import Anthropic from '@anthropic-ai/sdk';
|
||||
import { zodOutputFormat } from '@anthropic-ai/sdk/helpers/zod';
|
||||
import { z } from 'zod';
|
||||
import type { Env } from '../../lib/env';
|
||||
|
||||
/**
|
||||
* RULES.md section 9. The model returns this and nothing else: structured outputs
|
||||
* constrain the response to the schema, so there is no prose to strip and no JSON to
|
||||
* repair. Every field is nullable, because "unreadable" is a real answer and guessing a
|
||||
* number onto a tax document is the one thing this must never do.
|
||||
*/
|
||||
export const OcrExtraction = z.object({
|
||||
emitter_ruc: z.string().nullable(),
|
||||
emitter_dv: z.string().nullable(),
|
||||
emitter_name: z.string().nullable(),
|
||||
receiver_doc: z.string().nullable(),
|
||||
doc_number: z.string().nullable(),
|
||||
issue_date: z.string().nullable(),
|
||||
total: z.number().int().nullable(),
|
||||
amount_iva10: z.number().int().nullable(),
|
||||
amount_iva5: z.number().int().nullable(),
|
||||
amount_exenta: z.number().int().nullable(),
|
||||
iva10: z.number().int().nullable(),
|
||||
iva5: z.number().int().nullable(),
|
||||
confidence: z.record(z.string(), z.number()),
|
||||
});
|
||||
export type OcrExtraction = z.infer<typeof OcrExtraction>;
|
||||
|
||||
export interface OcrProvider {
|
||||
extract(args: { data: Uint8Array; mime: string }): Promise<OcrExtraction>;
|
||||
}
|
||||
|
||||
const SYSTEM_PROMPT = `Sos un extractor de datos de comprobantes fiscales paraguayos (facturas, autofacturas, notas de credito y debito).
|
||||
|
||||
Leé la imagen y devolvé unicamente los campos del esquema.
|
||||
|
||||
Reglas:
|
||||
- Las facturas paraguayas imprimen las columnas de IVA como "10%", "5%" y "Exentas".
|
||||
amount_iva10, amount_iva5 y amount_exenta son las bases gravadas de cada columna;
|
||||
iva10 e iva5 son los impuestos liquidados de cada una.
|
||||
- Los importes se imprimen con punto como separador de miles y no llevan decimales.
|
||||
Devolvelos como enteros sin separadores: "1.234.567" es 1234567.
|
||||
- issue_date en formato YYYY-MM-DD.
|
||||
- emitter_ruc sin el digito verificador; emitter_dv es ese digito.
|
||||
- Si un campo no se lee con claridad, devolvé null. Nunca adivines un numero.
|
||||
- confidence lleva una entrada por campo que sí leiste, de 0 a 1.`;
|
||||
|
||||
/** Guaranies have no cents, so a component sum may legitimately differ by rounding. */
|
||||
const TOTAL_TOLERANCE_GS = 1;
|
||||
/** RULES.md section 9: money fields drop to this when the components do not add up. */
|
||||
const MISMATCH_CONFIDENCE = 0.5;
|
||||
|
||||
/**
|
||||
* Returns null when no key is configured. Every caller treats that as "OCR is off" and
|
||||
* falls back to manual entry rather than failing the scan (SPEC.md section 8).
|
||||
*/
|
||||
export function createOcrProvider(env: Env): OcrProvider | null {
|
||||
if (!env.ANTHROPIC_API_KEY) return null;
|
||||
|
||||
const client = new Anthropic({ apiKey: env.ANTHROPIC_API_KEY });
|
||||
|
||||
return {
|
||||
async extract({ data, mime }) {
|
||||
const message = await client.messages.parse({
|
||||
model: env.OCR_MODEL,
|
||||
max_tokens: 4096,
|
||||
system: SYSTEM_PROMPT,
|
||||
output_config: { format: zodOutputFormat(OcrExtraction) },
|
||||
messages: [
|
||||
{
|
||||
role: 'user',
|
||||
content: [
|
||||
{
|
||||
type: 'image',
|
||||
source: {
|
||||
type: 'base64',
|
||||
media_type: imageMediaType(mime),
|
||||
data: Buffer.from(data).toString('base64'),
|
||||
},
|
||||
},
|
||||
{ type: 'text', text: 'Extraé los datos de este comprobante.' },
|
||||
],
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
const parsed = message.parsed_output;
|
||||
if (!parsed) throw new Error('OCR returned no parseable output');
|
||||
return lowerConfidenceOnMismatch(OcrExtraction.parse(parsed));
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* RULES.md section 9: when the total and the components are both present and disagree,
|
||||
* the money fields are not trusted, whatever the model said about them.
|
||||
*/
|
||||
export function lowerConfidenceOnMismatch(extraction: OcrExtraction): OcrExtraction {
|
||||
const { total, amount_iva10, amount_iva5, amount_exenta, iva10, iva5 } = extraction;
|
||||
const components = [amount_iva10, amount_iva5, amount_exenta, iva10, iva5];
|
||||
if (total === null || components.some((value) => value === null)) return extraction;
|
||||
|
||||
const derived =
|
||||
(amount_iva10 ?? 0) + (amount_iva5 ?? 0) + (amount_exenta ?? 0) + (iva10 ?? 0) + (iva5 ?? 0);
|
||||
if (Math.abs(total - derived) <= TOTAL_TOLERANCE_GS) return extraction;
|
||||
|
||||
const confidence = { ...extraction.confidence };
|
||||
for (const field of ['total', 'amount_iva10', 'amount_iva5', 'amount_exenta', 'iva10', 'iva5']) {
|
||||
if (field in confidence) confidence[field] = Math.min(confidence[field] ?? 1, MISMATCH_CONFIDENCE);
|
||||
else confidence[field] = MISMATCH_CONFIDENCE;
|
||||
}
|
||||
return { ...extraction, confidence };
|
||||
}
|
||||
|
||||
type ImageMediaType = 'image/jpeg' | 'image/png' | 'image/gif' | 'image/webp';
|
||||
|
||||
function imageMediaType(mime: string): ImageMediaType {
|
||||
switch (mime) {
|
||||
case 'image/png':
|
||||
case 'image/gif':
|
||||
case 'image/webp':
|
||||
return mime;
|
||||
default:
|
||||
return 'image/jpeg';
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,470 @@
|
||||
import type {
|
||||
ClassificationPatchInput,
|
||||
DocumentDto,
|
||||
DocumentListQuery,
|
||||
DocumentPatchInput,
|
||||
ManualDocumentInput,
|
||||
} from '@impuestos/contracts';
|
||||
import { RULES_VERSION, classify, type ClassificationInput } from '@impuestos/rules';
|
||||
import type { Kysely } from 'kysely';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import type { Database, DocumentsTable } from '../../db/schema';
|
||||
import { getProfile } from '../pii';
|
||||
import type { StorageDriver } from '../storage';
|
||||
import { dedupeHash } from './dedupe';
|
||||
import { toDocumentDto, type ClassificationRow, type DocumentRow } from './mapper';
|
||||
|
||||
const PAGE_SIZE = 50;
|
||||
|
||||
export interface DocumentsDeps {
|
||||
db: Kysely<Database>;
|
||||
storage: StorageDriver;
|
||||
}
|
||||
|
||||
export interface NewDocument {
|
||||
userId: string;
|
||||
source: 'scan_qr' | 'scan_ocr' | 'manual';
|
||||
cdc?: string | null;
|
||||
qrUrl?: string | null;
|
||||
docKind: DocumentsTable['doc_kind'];
|
||||
direction: 'purchase' | 'sale';
|
||||
emitterRuc: string;
|
||||
emitterDv?: string | null;
|
||||
emitterName: string;
|
||||
receiverDoc?: string | null;
|
||||
issueDate: string;
|
||||
total: number;
|
||||
amountIva10: number;
|
||||
amountIva5: number;
|
||||
amountExenta: number;
|
||||
iva10: number;
|
||||
iva5: number;
|
||||
supplierRegimeHint: 'normal' | 'resimple' | 'unknown';
|
||||
fileId?: string | null;
|
||||
rawExtraction?: Record<string, unknown> | null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a document, or merges into the one already holding the same comprobante.
|
||||
*
|
||||
* A merge prefers QR sourced data over anything read from a photograph or typed by hand:
|
||||
* the QR carries the numbers the emitter actually reported to DNIT (SPEC.md section 8).
|
||||
*/
|
||||
export async function createOrMerge(
|
||||
deps: DocumentsDeps,
|
||||
input: NewDocument,
|
||||
): Promise<{ document: DocumentDto; merged: boolean }> {
|
||||
const hash = dedupeHash({
|
||||
cdc: input.cdc ?? null,
|
||||
emitterRuc: input.emitterRuc,
|
||||
issueDate: input.issueDate,
|
||||
total: input.total,
|
||||
docKind: input.docKind,
|
||||
});
|
||||
|
||||
const existing = await deps.db
|
||||
.selectFrom('documents')
|
||||
.selectAll()
|
||||
.where('user_id', '=', input.userId)
|
||||
.where('dedupe_hash', '=', hash)
|
||||
.executeTakeFirst();
|
||||
|
||||
if (existing) {
|
||||
const merged = await mergeInto(deps, existing, input);
|
||||
return { document: merged, merged: true };
|
||||
}
|
||||
|
||||
const now = new Date().toISOString();
|
||||
const id = uuidv7();
|
||||
|
||||
await deps.db
|
||||
.insertInto('documents')
|
||||
.values({
|
||||
id,
|
||||
user_id: input.userId,
|
||||
source: input.source,
|
||||
status: 'needs_review',
|
||||
cdc: input.cdc ?? null,
|
||||
qr_url: input.qrUrl ?? null,
|
||||
doc_kind: input.docKind,
|
||||
direction: input.direction,
|
||||
emitter_ruc: input.emitterRuc,
|
||||
emitter_dv: input.emitterDv ?? null,
|
||||
emitter_name: input.emitterName,
|
||||
receiver_doc: input.receiverDoc ?? null,
|
||||
issue_date: input.issueDate,
|
||||
currency: 'PYG',
|
||||
total: input.total,
|
||||
amount_iva10: input.amountIva10,
|
||||
amount_iva5: input.amountIva5,
|
||||
amount_exenta: input.amountExenta,
|
||||
iva10: input.iva10,
|
||||
iva5: input.iva5,
|
||||
supplier_regime_hint: input.supplierRegimeHint,
|
||||
verified_dnit: 0,
|
||||
verification_status: 'unverified',
|
||||
dedupe_hash: hash,
|
||||
file_id: input.fileId ?? null,
|
||||
raw_extraction: input.rawExtraction ? JSON.stringify(input.rawExtraction) : null,
|
||||
created_at: now,
|
||||
confirmed_at: null,
|
||||
})
|
||||
.execute();
|
||||
|
||||
await classifyDocument(deps, id);
|
||||
const document = await getDocument(deps, input.userId, id);
|
||||
if (!document) throw new Error('document disappeared immediately after being written');
|
||||
return { document, merged: false };
|
||||
}
|
||||
|
||||
async function mergeInto(
|
||||
deps: DocumentsDeps,
|
||||
existing: DocumentRow,
|
||||
incoming: NewDocument,
|
||||
): Promise<DocumentDto> {
|
||||
const incomingIsQr = incoming.source === 'scan_qr';
|
||||
const existingIsQr = existing.source === 'scan_qr';
|
||||
|
||||
// Only a QR beats what is already there. Anything else just fills in the blanks.
|
||||
const patch: Partial<DocumentsTable> = incomingIsQr && !existingIsQr
|
||||
? {
|
||||
source: 'scan_qr',
|
||||
cdc: incoming.cdc ?? existing.cdc,
|
||||
qr_url: incoming.qrUrl ?? existing.qr_url,
|
||||
emitter_ruc: incoming.emitterRuc,
|
||||
emitter_dv: incoming.emitterDv ?? existing.emitter_dv,
|
||||
emitter_name: incoming.emitterName,
|
||||
issue_date: incoming.issueDate,
|
||||
total: incoming.total,
|
||||
amount_iva10: incoming.amountIva10,
|
||||
amount_iva5: incoming.amountIva5,
|
||||
amount_exenta: incoming.amountExenta,
|
||||
iva10: incoming.iva10,
|
||||
iva5: incoming.iva5,
|
||||
}
|
||||
: {
|
||||
cdc: existing.cdc ?? incoming.cdc ?? null,
|
||||
file_id: existing.file_id ?? incoming.fileId ?? null,
|
||||
};
|
||||
|
||||
await deps.db.updateTable('documents').set(patch).where('id', '=', existing.id).execute();
|
||||
|
||||
// A user who already judged this document keeps their decision.
|
||||
const classification = await deps.db
|
||||
.selectFrom('classifications')
|
||||
.select('decided_by')
|
||||
.where('document_id', '=', existing.id)
|
||||
.executeTakeFirst();
|
||||
if (classification?.decided_by === 'auto') await classifyDocument(deps, existing.id);
|
||||
|
||||
const document = await getDocument(deps, existing.user_id, existing.id);
|
||||
if (!document) throw new Error('merged document vanished');
|
||||
return document;
|
||||
}
|
||||
|
||||
/** Runs packages/rules over a document and stores the suggestion. Never overrides a user. */
|
||||
export async function classifyDocument(deps: DocumentsDeps, documentId: string): Promise<void> {
|
||||
const row = await deps.db
|
||||
.selectFrom('documents')
|
||||
.selectAll()
|
||||
.where('id', '=', documentId)
|
||||
.executeTakeFirst();
|
||||
if (!row) return;
|
||||
|
||||
const existing = await deps.db
|
||||
.selectFrom('classifications')
|
||||
.selectAll()
|
||||
.where('document_id', '=', documentId)
|
||||
.executeTakeFirst();
|
||||
if (existing && existing.decided_by !== 'auto') return;
|
||||
|
||||
const profile = await getProfile(deps.db, row.user_id);
|
||||
const input: ClassificationInput = {
|
||||
direction: row.direction,
|
||||
docKind: row.doc_kind,
|
||||
emitterName: row.emitter_name,
|
||||
emitterRuc: row.emitter_ruc,
|
||||
supplierRegimeHint: row.supplier_regime_hint,
|
||||
taxpayer: {
|
||||
kind: profile?.taxpayerKind ?? 'individual',
|
||||
hasIva: hasObligation(profile, 'iva_120'),
|
||||
hasIrp: hasObligation(profile, 'irp_515'),
|
||||
},
|
||||
amounts: { total: row.total, iva10: row.iva10, iva5: row.iva5 } as ClassificationInput['amounts'],
|
||||
};
|
||||
|
||||
const suggestion = classify(input);
|
||||
const now = new Date().toISOString();
|
||||
|
||||
const values = {
|
||||
iva_credit_eligible: suggestion.ivaCreditEligible ? 1 : 0,
|
||||
iva_credit_amount: suggestion.ivaCreditAmount,
|
||||
irp_category: suggestion.irpCategory,
|
||||
irp_deductible_amount: suggestion.irpDeductibleAmount,
|
||||
dependent_id: null,
|
||||
confidence: suggestion.confidence,
|
||||
decided_by: 'auto' as const,
|
||||
rules_version: RULES_VERSION,
|
||||
updated_at: now,
|
||||
};
|
||||
|
||||
if (existing) {
|
||||
await deps.db
|
||||
.updateTable('classifications')
|
||||
.set(values)
|
||||
.where('document_id', '=', documentId)
|
||||
.execute();
|
||||
} else {
|
||||
await deps.db.insertInto('classifications').values({ document_id: documentId, ...values }).execute();
|
||||
}
|
||||
|
||||
// Reason codes ride with the document so the detail sheet can explain the decision.
|
||||
const raw = row.raw_extraction ? (JSON.parse(row.raw_extraction) as Record<string, unknown>) : {};
|
||||
await deps.db
|
||||
.updateTable('documents')
|
||||
.set({ raw_extraction: JSON.stringify({ ...raw, reasons: suggestion.reasons }) })
|
||||
.where('id', '=', documentId)
|
||||
.execute();
|
||||
}
|
||||
|
||||
function hasObligation(
|
||||
profile: { obligations: { code: string; active: boolean }[] } | null,
|
||||
code: string,
|
||||
): boolean {
|
||||
return profile?.obligations.some((o) => o.code === code && o.active) ?? false;
|
||||
}
|
||||
|
||||
export async function getDocument(
|
||||
deps: DocumentsDeps,
|
||||
userId: string,
|
||||
id: string,
|
||||
): Promise<DocumentDto | null> {
|
||||
const row = await deps.db
|
||||
.selectFrom('documents')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.where('user_id', '=', userId)
|
||||
.executeTakeFirst();
|
||||
if (!row) return null;
|
||||
|
||||
const classification = await deps.db
|
||||
.selectFrom('classifications')
|
||||
.selectAll()
|
||||
.where('document_id', '=', id)
|
||||
.executeTakeFirst();
|
||||
|
||||
return toDocumentDto(row, classification ?? null, await fileUrlFor(deps, row.file_id));
|
||||
}
|
||||
|
||||
async function fileUrlFor(deps: DocumentsDeps, fileId: string | null): Promise<string | null> {
|
||||
if (!fileId) return null;
|
||||
const file = await deps.db
|
||||
.selectFrom('document_files')
|
||||
.select(['id', 'path'])
|
||||
.where('id', '=', fileId)
|
||||
.executeTakeFirst();
|
||||
if (!file) return null;
|
||||
return deps.storage.kind === 's3' ? deps.storage.url(file.path) : `/api/files/${file.id}`;
|
||||
}
|
||||
|
||||
export async function listDocuments(
|
||||
deps: DocumentsDeps,
|
||||
userId: string,
|
||||
query: DocumentListQuery,
|
||||
): Promise<{ items: DocumentDto[]; total: number; cursor?: string }> {
|
||||
let base = deps.db.selectFrom('documents').where('user_id', '=', userId);
|
||||
|
||||
if (query.month) {
|
||||
base = base.where('issue_date', '>=', `${query.month}-01`).where('issue_date', '<=', `${query.month}-31`);
|
||||
}
|
||||
if (query.direction) base = base.where('direction', '=', query.direction);
|
||||
if (query.status) base = base.where('status', '=', query.status);
|
||||
if (query.q) base = base.where('emitter_name', 'like', `%${query.q.toUpperCase()}%`);
|
||||
if (query.category) {
|
||||
base = base.where(({ exists, selectFrom }) =>
|
||||
exists(
|
||||
selectFrom('classifications')
|
||||
.select('document_id')
|
||||
.whereRef('classifications.document_id', '=', 'documents.id')
|
||||
.where('irp_category', '=', query.category as string),
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
const { total } = await base
|
||||
.select((eb) => eb.fn.countAll<number>().as('total'))
|
||||
.executeTakeFirstOrThrow();
|
||||
|
||||
let page = base.selectAll().orderBy('issue_date', 'desc').orderBy('id', 'desc').limit(PAGE_SIZE + 1);
|
||||
if (query.cursor) page = page.where('id', '<', query.cursor);
|
||||
|
||||
const rows = await page.execute();
|
||||
const hasMore = rows.length > PAGE_SIZE;
|
||||
const visible = hasMore ? rows.slice(0, PAGE_SIZE) : rows;
|
||||
|
||||
const items = await Promise.all(
|
||||
visible.map(async (row) => {
|
||||
const classification = await deps.db
|
||||
.selectFrom('classifications')
|
||||
.selectAll()
|
||||
.where('document_id', '=', row.id)
|
||||
.executeTakeFirst();
|
||||
return toDocumentDto(row, classification ?? null, await fileUrlFor(deps, row.file_id));
|
||||
}),
|
||||
);
|
||||
|
||||
const last = visible.at(-1);
|
||||
return { items, total: Number(total), ...(hasMore && last ? { cursor: last.id } : {}) };
|
||||
}
|
||||
|
||||
export async function confirmDocument(
|
||||
deps: DocumentsDeps,
|
||||
userId: string,
|
||||
id: string,
|
||||
): Promise<DocumentDto | null> {
|
||||
await deps.db
|
||||
.updateTable('documents')
|
||||
.set({ status: 'confirmed', confirmed_at: new Date().toISOString() })
|
||||
.where('id', '=', id)
|
||||
.where('user_id', '=', userId)
|
||||
.execute();
|
||||
return getDocument(deps, userId, id);
|
||||
}
|
||||
|
||||
export async function rejectDocument(
|
||||
deps: DocumentsDeps,
|
||||
userId: string,
|
||||
id: string,
|
||||
reason: string,
|
||||
): Promise<DocumentDto | null> {
|
||||
const row = await deps.db
|
||||
.selectFrom('documents')
|
||||
.select(['id', 'raw_extraction'])
|
||||
.where('id', '=', id)
|
||||
.where('user_id', '=', userId)
|
||||
.executeTakeFirst();
|
||||
if (!row) return null;
|
||||
|
||||
const raw = row.raw_extraction ? (JSON.parse(row.raw_extraction) as Record<string, unknown>) : {};
|
||||
await deps.db
|
||||
.updateTable('documents')
|
||||
.set({ status: 'rejected', raw_extraction: JSON.stringify({ ...raw, rejectReason: reason }) })
|
||||
.where('id', '=', id)
|
||||
.execute();
|
||||
|
||||
return getDocument(deps, userId, id);
|
||||
}
|
||||
|
||||
/** Editing the numbers re-runs the rules, unless the user has already overridden them. */
|
||||
export async function patchDocument(
|
||||
deps: DocumentsDeps,
|
||||
userId: string,
|
||||
id: string,
|
||||
patch: DocumentPatchInput,
|
||||
): Promise<DocumentDto | null> {
|
||||
const row = await deps.db
|
||||
.selectFrom('documents')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.where('user_id', '=', userId)
|
||||
.executeTakeFirst();
|
||||
if (!row) return null;
|
||||
|
||||
const next = {
|
||||
...(patch.direction === undefined ? {} : { direction: patch.direction }),
|
||||
...(patch.docKind === undefined ? {} : { doc_kind: patch.docKind }),
|
||||
...(patch.emitterRuc === undefined ? {} : { emitter_ruc: patch.emitterRuc }),
|
||||
...(patch.emitterDv === undefined ? {} : { emitter_dv: patch.emitterDv }),
|
||||
...(patch.emitterName === undefined ? {} : { emitter_name: patch.emitterName }),
|
||||
...(patch.issueDate === undefined ? {} : { issue_date: patch.issueDate }),
|
||||
...(patch.total === undefined ? {} : { total: patch.total }),
|
||||
...(patch.amountIva10 === undefined ? {} : { amount_iva10: patch.amountIva10 }),
|
||||
...(patch.amountIva5 === undefined ? {} : { amount_iva5: patch.amountIva5 }),
|
||||
...(patch.amountExenta === undefined ? {} : { amount_exenta: patch.amountExenta }),
|
||||
...(patch.iva10 === undefined ? {} : { iva10: patch.iva10 }),
|
||||
...(patch.iva5 === undefined ? {} : { iva5: patch.iva5 }),
|
||||
...(patch.supplierRegimeHint === undefined ? {} : { supplier_regime_hint: patch.supplierRegimeHint }),
|
||||
};
|
||||
|
||||
if (Object.keys(next).length > 0) {
|
||||
const hash = dedupeHash({
|
||||
cdc: row.cdc,
|
||||
emitterRuc: next.emitter_ruc ?? row.emitter_ruc,
|
||||
issueDate: next.issue_date ?? row.issue_date,
|
||||
total: next.total ?? row.total,
|
||||
docKind: next.doc_kind ?? row.doc_kind,
|
||||
});
|
||||
await deps.db
|
||||
.updateTable('documents')
|
||||
.set({ ...next, dedupe_hash: hash })
|
||||
.where('id', '=', id)
|
||||
.execute();
|
||||
}
|
||||
|
||||
await classifyDocument(deps, id);
|
||||
return getDocument(deps, userId, id);
|
||||
}
|
||||
|
||||
/** A user decision. `decided_by` flips to 'user' and the classifier stops overriding it. */
|
||||
export async function patchClassification(
|
||||
deps: DocumentsDeps,
|
||||
userId: string,
|
||||
id: string,
|
||||
patch: ClassificationPatchInput,
|
||||
): Promise<DocumentDto | null> {
|
||||
const row = await deps.db
|
||||
.selectFrom('documents')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.where('user_id', '=', userId)
|
||||
.executeTakeFirst();
|
||||
if (!row) return null;
|
||||
|
||||
const existing = await deps.db
|
||||
.selectFrom('classifications')
|
||||
.selectAll()
|
||||
.where('document_id', '=', id)
|
||||
.executeTakeFirst();
|
||||
if (!existing) {
|
||||
await classifyDocument(deps, id);
|
||||
}
|
||||
|
||||
const current = (existing ??
|
||||
(await deps.db
|
||||
.selectFrom('classifications')
|
||||
.selectAll()
|
||||
.where('document_id', '=', id)
|
||||
.executeTakeFirstOrThrow())) as ClassificationRow;
|
||||
|
||||
const irpCategory = patch.irpCategory ?? current.irp_category;
|
||||
const ivaCreditEligible = patch.ivaCreditEligible ?? current.iva_credit_eligible === 1;
|
||||
|
||||
await deps.db
|
||||
.updateTable('classifications')
|
||||
.set({
|
||||
irp_category: irpCategory,
|
||||
// A category the user chose means the document is deductible at its full total.
|
||||
irp_deductible_amount: irpCategory === 'none' ? 0 : row.total,
|
||||
iva_credit_eligible: ivaCreditEligible ? 1 : 0,
|
||||
iva_credit_amount: ivaCreditEligible ? row.iva10 + row.iva5 : 0,
|
||||
...(patch.dependentId === undefined ? {} : { dependent_id: patch.dependentId }),
|
||||
decided_by: 'user',
|
||||
updated_at: new Date().toISOString(),
|
||||
})
|
||||
.where('document_id', '=', id)
|
||||
.execute();
|
||||
|
||||
return getDocument(deps, userId, id);
|
||||
}
|
||||
|
||||
export async function countNeedsReview(deps: DocumentsDeps, userId: string): Promise<number> {
|
||||
const { total } = await deps.db
|
||||
.selectFrom('documents')
|
||||
.select((eb) => eb.fn.countAll<number>().as('total'))
|
||||
.where('user_id', '=', userId)
|
||||
.where('status', '=', 'needs_review')
|
||||
.executeTakeFirstOrThrow();
|
||||
return Number(total);
|
||||
}
|
||||
|
||||
export type { ManualDocumentInput };
|
||||
@@ -0,0 +1,63 @@
|
||||
import type { IngestErrorDto } from '@impuestos/contracts';
|
||||
import type { Kysely } from 'kysely';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import type { Database } from '../../db/schema';
|
||||
|
||||
export type IngestStage = 'qr_parse' | 'ocr' | 'dedupe' | 'verify' | 'job' | 'other';
|
||||
|
||||
/**
|
||||
* Anything that goes wrong on the way in is recorded rather than swallowed: the user gets
|
||||
* a retry affordance and staff get a queue (SPEC.md section 8, FLOWS.md Flow H).
|
||||
*/
|
||||
export async function recordIngestError(
|
||||
db: Kysely<Database>,
|
||||
entry: {
|
||||
userId?: string | null;
|
||||
documentId?: string | null;
|
||||
stage: IngestStage;
|
||||
message: string;
|
||||
payload?: Record<string, unknown>;
|
||||
},
|
||||
): Promise<string> {
|
||||
const id = uuidv7();
|
||||
await db
|
||||
.insertInto('ingest_errors')
|
||||
.values({
|
||||
id,
|
||||
user_id: entry.userId ?? null,
|
||||
document_id: entry.documentId ?? null,
|
||||
stage: entry.stage,
|
||||
message: entry.message.slice(0, 1000),
|
||||
payload: entry.payload ? JSON.stringify(entry.payload) : null,
|
||||
status: 'open',
|
||||
resolved_by: null,
|
||||
resolved_at: null,
|
||||
created_at: new Date().toISOString(),
|
||||
})
|
||||
.execute();
|
||||
return id;
|
||||
}
|
||||
|
||||
export async function listIngestErrorsForUser(
|
||||
db: Kysely<Database>,
|
||||
userId: string,
|
||||
): Promise<IngestErrorDto[]> {
|
||||
const rows = await db
|
||||
.selectFrom('ingest_errors')
|
||||
.selectAll()
|
||||
.where('user_id', '=', userId)
|
||||
.where('status', '=', 'open')
|
||||
.orderBy('created_at', 'desc')
|
||||
.limit(20)
|
||||
.execute();
|
||||
|
||||
return rows.map((row) => ({
|
||||
id: row.id,
|
||||
userId: row.user_id,
|
||||
documentId: row.document_id,
|
||||
stage: row.stage,
|
||||
message: row.message,
|
||||
status: row.status,
|
||||
createdAt: row.created_at,
|
||||
}));
|
||||
}
|
||||
@@ -0,0 +1,238 @@
|
||||
import { readFileSync } from 'node:fs';
|
||||
import { DocumentDto, NeedsManualDto } from '@impuestos/contracts';
|
||||
import { afterEach, describe, expect, it } from 'vitest';
|
||||
import { FIXTURES } from '../../db/fixtures';
|
||||
import { createHarness, type Harness } from '../../test/harness';
|
||||
import type { OcrExtraction, OcrProvider } from '../documents/ocr';
|
||||
|
||||
const FIXTURE_1 = FIXTURES[0]!;
|
||||
const fixtureBytes = (name: string) =>
|
||||
new Uint8Array(readFileSync(new URL(`../../../../../fixtures/${name}.png`, import.meta.url)));
|
||||
|
||||
let harness: Harness | null = null;
|
||||
|
||||
afterEach(async () => {
|
||||
await harness?.close();
|
||||
harness = null;
|
||||
});
|
||||
|
||||
async function scan(
|
||||
h: Harness,
|
||||
cookie: string,
|
||||
args: { fixture: string; qrPayload?: string },
|
||||
): Promise<Response> {
|
||||
const form = new FormData();
|
||||
form.set('file', new File([fixtureBytes(args.fixture)], `${args.fixture}.png`, { type: 'image/png' }));
|
||||
if (args.qrPayload) form.set('qrPayload', args.qrPayload);
|
||||
return h.app.request('/api/documents/scan', { method: 'POST', body: form, headers: { cookie } });
|
||||
}
|
||||
|
||||
describe('scanning a QR comprobante', () => {
|
||||
it('creates a document prefilled from the CDC', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
const response = await scan(h, cookie, {
|
||||
fixture: FIXTURE_1.name,
|
||||
qrPayload: FIXTURE_1.qrUrl as string,
|
||||
});
|
||||
expect(response.status).toBe(201);
|
||||
|
||||
const document = DocumentDto.parse(await response.json());
|
||||
expect(document.source).toBe('scan_qr');
|
||||
expect(document.status).toBe('needs_review');
|
||||
expect(document.cdc).toBe(FIXTURE_1.cdc);
|
||||
expect(document.issueDate).toBe(FIXTURE_1.issueDate);
|
||||
expect(document.total).toBe(FIXTURE_1.total);
|
||||
expect(document.emitterRuc).toBe(FIXTURE_1.emitterRucBase);
|
||||
expect(document.fileUrl).not.toBeNull();
|
||||
expect(document.classification).not.toBeNull();
|
||||
});
|
||||
|
||||
// CONTRACTS.md 5.1: the same comprobante twice is one document, not two.
|
||||
it('collapses a second scan of the same fixture into the first', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
const before = await countDocuments(h, 'maria@demo.local');
|
||||
|
||||
const first = DocumentDto.parse(
|
||||
await (await scan(h, cookie, { fixture: FIXTURE_1.name, qrPayload: FIXTURE_1.qrUrl as string })).json(),
|
||||
);
|
||||
|
||||
const second = await scan(h, cookie, {
|
||||
fixture: FIXTURE_1.name,
|
||||
qrPayload: FIXTURE_1.qrUrl as string,
|
||||
});
|
||||
expect(second.status).toBe(200);
|
||||
|
||||
const merged = DocumentDto.parse(await second.json());
|
||||
expect(merged.merged).toBe(true);
|
||||
expect(merged.id).toBe(first.id);
|
||||
expect(await countDocuments(h, 'maria@demo.local')).toBe(before + 1);
|
||||
});
|
||||
|
||||
it('keeps two different comprobantes apart', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
const second = FIXTURES[1]!;
|
||||
|
||||
const a = DocumentDto.parse(
|
||||
await (await scan(h, cookie, { fixture: FIXTURE_1.name, qrPayload: FIXTURE_1.qrUrl as string })).json(),
|
||||
);
|
||||
const b = DocumentDto.parse(
|
||||
await (await scan(h, cookie, { fixture: second.name, qrPayload: second.qrUrl as string })).json(),
|
||||
);
|
||||
expect(b.id).not.toBe(a.id);
|
||||
expect(b.merged).toBeUndefined();
|
||||
});
|
||||
|
||||
it('does not merge across users', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const maria = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
const carlos = await h.signIn('carlos@demo.local', 'demo-carlos-1');
|
||||
|
||||
const mine = DocumentDto.parse(
|
||||
await (await scan(h, maria, { fixture: FIXTURE_1.name, qrPayload: FIXTURE_1.qrUrl as string })).json(),
|
||||
);
|
||||
const theirs = DocumentDto.parse(
|
||||
await (await scan(h, carlos, { fixture: FIXTURE_1.name, qrPayload: FIXTURE_1.qrUrl as string })).json(),
|
||||
);
|
||||
expect(theirs.id).not.toBe(mine.id);
|
||||
});
|
||||
|
||||
it('records an ingest error for an unreadable QR and still keeps the file', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
const response = await scan(h, cookie, { fixture: FIXTURE_1.name, qrPayload: 'not-a-cdc' });
|
||||
expect(response.status).toBe(200);
|
||||
NeedsManualDto.parse(await response.json());
|
||||
|
||||
const errors = await h.deps.handle.db
|
||||
.selectFrom('ingest_errors')
|
||||
.selectAll()
|
||||
.where('stage', '=', 'qr_parse')
|
||||
.execute();
|
||||
expect(errors.length).toBeGreaterThan(0);
|
||||
});
|
||||
});
|
||||
|
||||
describe('scanning without a QR', () => {
|
||||
it('asks for manual entry when OCR is not configured', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
const response = await scan(h, cookie, { fixture: 'factura-sin-qr-ferreteria' });
|
||||
expect(response.status).toBe(200);
|
||||
|
||||
const result = NeedsManualDto.parse(await response.json());
|
||||
expect(result.needsManual).toBe(true);
|
||||
expect(result.fileId).toBeTruthy();
|
||||
|
||||
// Nothing was queued, because there is nothing that could read the image.
|
||||
const jobs = await h.deps.handle.db.selectFrom('jobs').selectAll().execute();
|
||||
expect(jobs.filter((job) => job.type === 'ocr_extract')).toEqual([]);
|
||||
});
|
||||
|
||||
it('queues extraction when OCR is configured, and the job creates the document', async () => {
|
||||
const h = (harness = await createHarness({ ocr: stubOcr() }));
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
await scan(h, cookie, { fixture: 'factura-sin-qr-ferreteria' });
|
||||
|
||||
const queued = await h.deps.handle.db
|
||||
.selectFrom('jobs')
|
||||
.selectAll()
|
||||
.where('type', '=', 'ocr_extract')
|
||||
.execute();
|
||||
expect(queued).toHaveLength(1);
|
||||
expect(queued[0]?.status).toBe('pending');
|
||||
|
||||
expect(await h.runJobs()).toBe(1);
|
||||
|
||||
const document = await h.deps.handle.db
|
||||
.selectFrom('documents')
|
||||
.selectAll()
|
||||
.where('emitter_name', '=', 'FERRETERIA SAN MIGUEL')
|
||||
.where('source', '=', 'scan_ocr')
|
||||
.executeTakeFirst();
|
||||
expect(document?.total).toBe(539_000);
|
||||
});
|
||||
});
|
||||
|
||||
describe('manual entry', () => {
|
||||
it('creates a document and classifies it', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
const response = await h.app.request('/api/documents/manual', {
|
||||
method: 'POST',
|
||||
headers: { cookie, 'content-type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
emitterRuc: '80022114',
|
||||
emitterName: 'Farmacia Catedral',
|
||||
issueDate: '2026-08-19',
|
||||
total: 176_000,
|
||||
amountIva10: 160_000,
|
||||
iva10: 16_000,
|
||||
}),
|
||||
});
|
||||
expect(response.status).toBe(201);
|
||||
|
||||
const document = DocumentDto.parse(await response.json());
|
||||
expect(document.source).toBe('manual');
|
||||
expect(document.emitterName).toBe('FARMACIA CATEDRAL');
|
||||
expect(document.classification?.irpCategory).toBe('salud');
|
||||
expect(document.classification?.ivaCreditAmount).toBe(16_000);
|
||||
});
|
||||
|
||||
it('rejects a malformed date and a missing emitter', async () => {
|
||||
const h = (harness = await createHarness());
|
||||
const cookie = await h.signIn('maria@demo.local', 'demo-maria-1');
|
||||
|
||||
for (const body of [
|
||||
{ emitterRuc: '80022114', emitterName: 'X', issueDate: '19/08/2026', total: 1 },
|
||||
{ emitterRuc: '', emitterName: 'X', issueDate: '2026-08-19', total: 1 },
|
||||
]) {
|
||||
const response = await h.app.request('/api/documents/manual', {
|
||||
method: 'POST',
|
||||
headers: { cookie, 'content-type': 'application/json' },
|
||||
body: JSON.stringify(body),
|
||||
});
|
||||
expect(response.status).toBe(400);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
function stubOcr(): OcrProvider {
|
||||
return {
|
||||
extract: async (): Promise<OcrExtraction> => ({
|
||||
emitter_ruc: '80033005',
|
||||
emitter_dv: '1',
|
||||
emitter_name: 'FERRETERIA SAN MIGUEL',
|
||||
receiver_doc: null,
|
||||
doc_number: '0001777',
|
||||
// Deliberately not the seeded FERRETERIA figures: those would dedupe into the
|
||||
// existing document and this test is about the job creating a new one.
|
||||
issue_date: '2026-08-23',
|
||||
total: 539_000,
|
||||
amount_iva10: 490_000,
|
||||
amount_iva5: 0,
|
||||
amount_exenta: 0,
|
||||
iva10: 49_000,
|
||||
iva5: 0,
|
||||
confidence: { total: 0.95, emitter_name: 0.9 },
|
||||
}),
|
||||
};
|
||||
}
|
||||
|
||||
async function countDocuments(h: Harness, email: string): Promise<number> {
|
||||
const { total } = await h.deps.handle.db
|
||||
.selectFrom('documents')
|
||||
.innerJoin('user', 'user.id', 'documents.user_id')
|
||||
.select((eb) => eb.fn.countAll<number>().as('total'))
|
||||
.where('user.email', '=', email)
|
||||
.executeTakeFirstOrThrow();
|
||||
return Number(total);
|
||||
}
|
||||
@@ -0,0 +1,143 @@
|
||||
import type { DocumentDto, ManualDocumentInput, NeedsManualDto } from '@impuestos/contracts';
|
||||
import { docKindFromCdc, isRuleError, parseQrPayload } from '@impuestos/rules';
|
||||
import type { Kysely } from 'kysely';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import type { Database, DocumentsTable } from '../../db/schema';
|
||||
import { createOrMerge, type DocumentsDeps } from '../documents/service';
|
||||
import { sha256 } from '../documents/dedupe';
|
||||
import { enqueue } from '../jobs/queue';
|
||||
import type { StorageDriver } from '../storage';
|
||||
import { storageKey } from '../storage';
|
||||
import { recordIngestError } from './errors';
|
||||
|
||||
export interface IngestDeps extends DocumentsDeps {
|
||||
db: Kysely<Database>;
|
||||
storage: StorageDriver;
|
||||
/** Null when ANTHROPIC_API_KEY is unset: scans without a QR go to the manual form. */
|
||||
ocrEnabled: boolean;
|
||||
}
|
||||
|
||||
export interface StoredFile {
|
||||
id: string;
|
||||
path: string;
|
||||
}
|
||||
|
||||
/** Stores the upload and records it, before anything is known about what it contains. */
|
||||
export async function storeUpload(
|
||||
deps: IngestDeps,
|
||||
userId: string,
|
||||
file: { data: Uint8Array; mime: string },
|
||||
): Promise<StoredFile> {
|
||||
const id = uuidv7();
|
||||
const path = storageKey(userId, id);
|
||||
|
||||
await deps.storage.put(path, file.data, file.mime);
|
||||
await deps.db
|
||||
.insertInto('document_files')
|
||||
.values({
|
||||
id,
|
||||
driver: deps.storage.kind,
|
||||
path,
|
||||
mime: file.mime,
|
||||
size: file.data.byteLength,
|
||||
sha256: sha256(file.data),
|
||||
created_at: new Date().toISOString(),
|
||||
})
|
||||
.execute();
|
||||
|
||||
return { id, path };
|
||||
}
|
||||
|
||||
export type ScanOutcome =
|
||||
| { kind: 'document'; document: DocumentDto; merged: boolean }
|
||||
| { kind: 'needs_manual'; result: NeedsManualDto };
|
||||
|
||||
/**
|
||||
* SPEC.md section 8. Three paths, in order of how much we can trust the result:
|
||||
* QR present -> parse the CDC and prefill from it, which is authoritative.
|
||||
* No QR, OCR on -> queue extraction; the client polls the document.
|
||||
* No QR, OCR off -> tell the client to open the manual form against this file.
|
||||
*/
|
||||
export async function ingestScan(
|
||||
deps: IngestDeps,
|
||||
userId: string,
|
||||
args: { file: { data: Uint8Array; mime: string }; qrPayload?: string | undefined },
|
||||
): Promise<ScanOutcome> {
|
||||
const stored = await storeUpload(deps, userId, args.file);
|
||||
|
||||
if (args.qrPayload && args.qrPayload.trim().length > 0) {
|
||||
try {
|
||||
const qr = parseQrPayload(args.qrPayload);
|
||||
const { document, merged } = await createOrMerge(deps, {
|
||||
userId,
|
||||
source: 'scan_qr',
|
||||
cdc: qr.cdc.cdc,
|
||||
qrUrl: qr.raw,
|
||||
docKind: docKindFromCdc(qr.cdc.tipoDocumento) as DocumentsTable['doc_kind'],
|
||||
// A comprobante someone hands you is a purchase; a sale of your own is entered
|
||||
// deliberately, because getting this backwards moves money in the declaration.
|
||||
direction: 'purchase',
|
||||
emitterRuc: qr.cdc.rucEmisor.replace(/^0+/, ''),
|
||||
emitterDv: qr.cdc.dvEmisor,
|
||||
emitterName: `RUC ${qr.cdc.rucEmisor.replace(/^0+/, '')}`,
|
||||
issueDate: qr.cdc.fechaEmision,
|
||||
total: qr.total ?? 0,
|
||||
amountIva10: 0,
|
||||
amountIva5: 0,
|
||||
amountExenta: 0,
|
||||
iva10: qr.totalIva ?? 0,
|
||||
iva5: 0,
|
||||
supplierRegimeHint: 'unknown',
|
||||
fileId: stored.id,
|
||||
rawExtraction: { qr: qr.raw },
|
||||
});
|
||||
|
||||
await enqueue(deps.db, { type: 'verify_cdc', payload: { documentId: document.id } });
|
||||
return { kind: 'document', document, merged };
|
||||
} catch (error) {
|
||||
// A QR that will not parse is not a dead end: fall through to OCR or manual.
|
||||
await recordIngestError(deps.db, {
|
||||
userId,
|
||||
stage: 'qr_parse',
|
||||
message: isRuleError(error) ? `${error.code}: ${error.message}` : String(error),
|
||||
payload: { qrPayload: args.qrPayload.slice(0, 500) },
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
if (!deps.ocrEnabled) {
|
||||
return { kind: 'needs_manual', result: { needsManual: true, fileId: stored.id } };
|
||||
}
|
||||
|
||||
await enqueue(deps.db, {
|
||||
type: 'ocr_extract',
|
||||
payload: { userId, fileId: stored.id },
|
||||
maxAttempts: 3,
|
||||
});
|
||||
return { kind: 'needs_manual', result: { needsManual: true, fileId: stored.id } };
|
||||
}
|
||||
|
||||
export async function ingestManual(
|
||||
deps: IngestDeps,
|
||||
userId: string,
|
||||
input: ManualDocumentInput,
|
||||
): Promise<{ document: DocumentDto; merged: boolean }> {
|
||||
return createOrMerge(deps, {
|
||||
userId,
|
||||
source: 'manual',
|
||||
docKind: input.docKind,
|
||||
direction: input.direction,
|
||||
emitterRuc: input.emitterRuc,
|
||||
emitterDv: input.emitterDv ?? null,
|
||||
emitterName: input.emitterName.toUpperCase(),
|
||||
issueDate: input.issueDate,
|
||||
total: input.total,
|
||||
amountIva10: input.amountIva10,
|
||||
amountIva5: input.amountIva5,
|
||||
amountExenta: input.amountExenta,
|
||||
iva10: input.iva10,
|
||||
iva5: input.iva5,
|
||||
supplierRegimeHint: input.supplierRegimeHint,
|
||||
fileId: input.fileId ?? null,
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
import { type Kysely, sql } from 'kysely';
|
||||
import type { Dialect } from '../../db/index';
|
||||
import type { Database, JobsTable } from '../../db/schema';
|
||||
|
||||
export type JobRow = JobsTable;
|
||||
|
||||
/**
|
||||
* The one place outside the two driver factories where dialect specific SQL is allowed
|
||||
* (SPEC.md section 5), because claiming a job exactly once is the one thing the two
|
||||
* engines genuinely do differently.
|
||||
*
|
||||
* Postgres: `FOR UPDATE SKIP LOCKED` lets N workers take different rows concurrently.
|
||||
* SQLite: a single writer serialises everything, so a conditional UPDATE that only
|
||||
* matches a still-pending row is already atomic. The `changes` count says who won.
|
||||
*/
|
||||
export async function claimJob(
|
||||
db: Kysely<Database>,
|
||||
dialect: Dialect,
|
||||
args: { instanceId: string; now: string },
|
||||
): Promise<JobRow | null> {
|
||||
return dialect === 'postgres' ? claimPostgres(db, args) : claimSqlite(db, args);
|
||||
}
|
||||
|
||||
async function claimPostgres(
|
||||
db: Kysely<Database>,
|
||||
args: { instanceId: string; now: string },
|
||||
): Promise<JobRow | null> {
|
||||
return db.transaction().execute(async (trx) => {
|
||||
const candidate = await trx
|
||||
.selectFrom('jobs')
|
||||
.selectAll()
|
||||
.where('status', '=', 'pending')
|
||||
.where('run_at', '<=', args.now)
|
||||
.orderBy('run_at')
|
||||
.limit(1)
|
||||
.forUpdate()
|
||||
.skipLocked()
|
||||
.executeTakeFirst();
|
||||
|
||||
if (!candidate) return null;
|
||||
|
||||
await trx
|
||||
.updateTable('jobs')
|
||||
.set({
|
||||
status: 'running',
|
||||
locked_by: args.instanceId,
|
||||
locked_at: args.now,
|
||||
attempts: candidate.attempts + 1,
|
||||
updated_at: args.now,
|
||||
})
|
||||
.where('id', '=', candidate.id)
|
||||
.execute();
|
||||
|
||||
return { ...candidate, status: 'running', locked_by: args.instanceId, locked_at: args.now, attempts: candidate.attempts + 1 };
|
||||
});
|
||||
}
|
||||
|
||||
async function claimSqlite(
|
||||
db: Kysely<Database>,
|
||||
args: { instanceId: string; now: string },
|
||||
): Promise<JobRow | null> {
|
||||
const candidate = await db
|
||||
.selectFrom('jobs')
|
||||
.select('id')
|
||||
.where('status', '=', 'pending')
|
||||
.where('run_at', '<=', args.now)
|
||||
.orderBy('run_at')
|
||||
.limit(1)
|
||||
.executeTakeFirst();
|
||||
|
||||
if (!candidate) return null;
|
||||
|
||||
const result = await db
|
||||
.updateTable('jobs')
|
||||
.set({
|
||||
status: 'running',
|
||||
locked_by: args.instanceId,
|
||||
locked_at: args.now,
|
||||
attempts: sql<number>`attempts + 1`,
|
||||
updated_at: args.now,
|
||||
})
|
||||
// Still pending is the whole guard: a racing claim already flipped it.
|
||||
.where('id', '=', candidate.id)
|
||||
.where('status', '=', 'pending')
|
||||
.executeTakeFirst();
|
||||
|
||||
if (Number(result.numUpdatedRows) === 0) return null;
|
||||
|
||||
return (
|
||||
(await db.selectFrom('jobs').selectAll().where('id', '=', candidate.id).executeTakeFirst()) ??
|
||||
null
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Crash recovery: a worker that died mid job left the row `running` and nobody will ever
|
||||
* finish it. Anything locked longer than the stale window goes back to pending.
|
||||
*/
|
||||
export async function recoverStaleJobs(
|
||||
db: Kysely<Database>,
|
||||
staleMinutes: number,
|
||||
now: Date,
|
||||
): Promise<number> {
|
||||
const cutoff = new Date(now.getTime() - staleMinutes * 60_000).toISOString();
|
||||
|
||||
const result = await db
|
||||
.updateTable('jobs')
|
||||
.set({ status: 'pending', locked_by: null, locked_at: null, updated_at: now.toISOString() })
|
||||
.where('status', '=', 'running')
|
||||
.where('locked_at', '<', cutoff)
|
||||
.executeTakeFirst();
|
||||
|
||||
return Number(result.numUpdatedRows);
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
import type { Kysely } from 'kysely';
|
||||
import type { Database, DocumentsTable } from '../../db/schema';
|
||||
import { classifyDocument, createOrMerge, type DocumentsDeps } from '../documents/service';
|
||||
import type { OcrProvider } from '../documents/ocr';
|
||||
import { recordIngestError } from '../ingest/errors';
|
||||
import type { StorageDriver } from '../storage';
|
||||
import type { JobHandlers } from './poller';
|
||||
|
||||
export interface HandlerDeps extends DocumentsDeps {
|
||||
db: Kysely<Database>;
|
||||
storage: StorageDriver;
|
||||
ocr: OcrProvider | null;
|
||||
}
|
||||
|
||||
export function createHandlers(deps: HandlerDeps): JobHandlers {
|
||||
return {
|
||||
/**
|
||||
* Reads a photographed factura and creates the document from what it found. Throwing
|
||||
* is deliberate: the poller retries with backoff, and only a dead job records an
|
||||
* ingest error, so a transient API blip does not spam the user's error list.
|
||||
*/
|
||||
ocr_extract: async ({ payload }) => {
|
||||
const userId = String(payload['userId'] ?? '');
|
||||
const fileId = String(payload['fileId'] ?? '');
|
||||
if (!userId || !fileId) throw new Error('ocr_extract needs userId and fileId');
|
||||
if (!deps.ocr) throw new Error('OCR is not configured');
|
||||
|
||||
const file = await deps.db
|
||||
.selectFrom('document_files')
|
||||
.selectAll()
|
||||
.where('id', '=', fileId)
|
||||
.executeTakeFirstOrThrow();
|
||||
|
||||
const stream = await deps.storage.getStream(file.path);
|
||||
const data = new Uint8Array(await new Response(stream).arrayBuffer());
|
||||
const extraction = await deps.ocr.extract({ data, mime: file.mime });
|
||||
|
||||
if (!extraction.emitter_ruc || !extraction.issue_date || extraction.total === null) {
|
||||
// Not an error to retry: the photograph genuinely does not carry the fields.
|
||||
await recordIngestError(deps.db, {
|
||||
userId,
|
||||
stage: 'ocr',
|
||||
message: 'the image did not yield an emitter, a date and a total',
|
||||
payload: { fileId, extraction },
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
await createOrMerge(deps, {
|
||||
userId,
|
||||
source: 'scan_ocr',
|
||||
docKind: 'factura',
|
||||
direction: 'purchase',
|
||||
emitterRuc: extraction.emitter_ruc,
|
||||
emitterDv: extraction.emitter_dv,
|
||||
emitterName: (extraction.emitter_name ?? extraction.emitter_ruc).toUpperCase(),
|
||||
receiverDoc: extraction.receiver_doc,
|
||||
issueDate: extraction.issue_date,
|
||||
total: extraction.total,
|
||||
amountIva10: extraction.amount_iva10 ?? 0,
|
||||
amountIva5: extraction.amount_iva5 ?? 0,
|
||||
amountExenta: extraction.amount_exenta ?? 0,
|
||||
iva10: extraction.iva10 ?? 0,
|
||||
iva5: extraction.iva5 ?? 0,
|
||||
supplierRegimeHint: 'unknown',
|
||||
fileId,
|
||||
rawExtraction: { ocr: extraction },
|
||||
});
|
||||
},
|
||||
|
||||
classify_document: async ({ payload }) => {
|
||||
const documentId = String(payload['documentId'] ?? '');
|
||||
if (!documentId) throw new Error('classify_document needs documentId');
|
||||
await classifyDocument(deps, documentId);
|
||||
},
|
||||
|
||||
/**
|
||||
* Checking a CDC against DNIT is out of scope for v1 (SPEC.md section 10), so this
|
||||
* records that no verification was attempted rather than pretending one succeeded.
|
||||
*/
|
||||
verify_cdc: async ({ payload }) => {
|
||||
const documentId = String(payload['documentId'] ?? '');
|
||||
if (!documentId) return;
|
||||
await deps.db
|
||||
.updateTable('documents')
|
||||
.set({ verification_status: 'unverified', verified_dnit: 0 } satisfies Partial<DocumentsTable>)
|
||||
.where('id', '=', documentId)
|
||||
.execute();
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,4 @@
|
||||
export { enqueue, nextRunAt, BACKOFF_MS, type EnqueueOptions, type JobType } from './queue';
|
||||
export { createPoller, type JobHandlers, type JobContext, type Poller } from './poller';
|
||||
export { createHandlers, type HandlerDeps } from './handlers';
|
||||
export { claimJob, recoverStaleJobs, type JobRow } from './claim';
|
||||
@@ -0,0 +1,224 @@
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { createHarness } from '../../test/harness';
|
||||
import { recoverStaleJobs } from './claim';
|
||||
import { enqueue, nextRunAt } from './queue';
|
||||
|
||||
describe('enqueue', () => {
|
||||
it('queues a job that is due now', async () => {
|
||||
const h = await createHarness();
|
||||
const { id, deduped } = await enqueue(h.deps.handle.db, {
|
||||
type: 'classify_document',
|
||||
payload: { documentId: 'x' },
|
||||
});
|
||||
|
||||
expect(deduped).toBe(false);
|
||||
const row = await h.deps.handle.db
|
||||
.selectFrom('jobs')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.executeTakeFirstOrThrow();
|
||||
expect(row.status).toBe('pending');
|
||||
expect(row.attempts).toBe(0);
|
||||
await h.close();
|
||||
});
|
||||
|
||||
// Sweeps run daily and a restart must not queue a second copy (SPEC.md section 10).
|
||||
it('is idempotent for a given dedupe key', async () => {
|
||||
const h = await createHarness();
|
||||
const first = await enqueue(h.deps.handle.db, {
|
||||
type: 'deadline_sweep',
|
||||
payload: {},
|
||||
dedupeKey: 'sweep:2026-09-04',
|
||||
});
|
||||
const second = await enqueue(h.deps.handle.db, {
|
||||
type: 'deadline_sweep',
|
||||
payload: {},
|
||||
dedupeKey: 'sweep:2026-09-04',
|
||||
});
|
||||
|
||||
expect(second.deduped).toBe(true);
|
||||
expect(second.id).toBe(first.id);
|
||||
await h.close();
|
||||
});
|
||||
});
|
||||
|
||||
describe('backoff', () => {
|
||||
it('follows the schedule in SPEC.md section 10', () => {
|
||||
const now = new Date('2026-09-04T00:00:00Z');
|
||||
const minutesLater = (attempts: number) =>
|
||||
(new Date(nextRunAt(attempts, now)).getTime() - now.getTime()) / 60_000;
|
||||
|
||||
expect(minutesLater(1)).toBe(1);
|
||||
expect(minutesLater(2)).toBe(5);
|
||||
expect(minutesLater(3)).toBe(25);
|
||||
expect(minutesLater(4)).toBe(120);
|
||||
expect(minutesLater(5)).toBe(720);
|
||||
// Past the schedule the job is dead, but the helper must still return something sane.
|
||||
expect(minutesLater(6)).toBe(720);
|
||||
});
|
||||
});
|
||||
|
||||
describe('failure handling', () => {
|
||||
it('retries a failing job with backoff, then marks it dead', async () => {
|
||||
const h = await createHarness();
|
||||
await enqueue(h.deps.handle.db, {
|
||||
type: 'classify_document',
|
||||
payload: {}, // no documentId: the handler throws
|
||||
maxAttempts: 2,
|
||||
});
|
||||
|
||||
await h.runJobs();
|
||||
let row = await h.deps.handle.db.selectFrom('jobs').selectAll().executeTakeFirstOrThrow();
|
||||
expect(row.status).toBe('pending');
|
||||
expect(row.attempts).toBe(1);
|
||||
expect(row.last_error).toContain('documentId');
|
||||
// It is scheduled into the future, so the next tick does not pick it straight back up.
|
||||
expect(new Date(row.run_at).getTime()).toBeGreaterThan(Date.now());
|
||||
|
||||
// Make it due again and let it exhaust its attempts.
|
||||
await h.deps.handle.db
|
||||
.updateTable('jobs')
|
||||
.set({ run_at: new Date(Date.now() - 1000).toISOString() })
|
||||
.execute();
|
||||
await h.runJobs();
|
||||
|
||||
row = await h.deps.handle.db.selectFrom('jobs').selectAll().executeTakeFirstOrThrow();
|
||||
expect(row.status).toBe('dead');
|
||||
expect(row.attempts).toBe(2);
|
||||
await h.close();
|
||||
});
|
||||
});
|
||||
|
||||
/**
|
||||
* The phase 3 acceptance case: a process that dies mid job leaves the row locked and
|
||||
* `running`, and nobody would ever finish it. The next poller to come up takes it back.
|
||||
*/
|
||||
describe('crash recovery', () => {
|
||||
it('returns a job abandoned by a dead worker to the queue', async () => {
|
||||
const h = await createHarness();
|
||||
const { id } = await enqueue(h.deps.handle.db, {
|
||||
type: 'classify_document',
|
||||
payload: { documentId: 'abc' },
|
||||
});
|
||||
|
||||
// Exactly the row a killed process leaves behind.
|
||||
const abandonedAt = new Date(Date.now() - 60 * 60_000).toISOString();
|
||||
await h.deps.handle.db
|
||||
.updateTable('jobs')
|
||||
.set({ status: 'running', locked_by: 'worker-that-died', locked_at: abandonedAt, attempts: 1 })
|
||||
.where('id', '=', id)
|
||||
.execute();
|
||||
|
||||
const recovered = await recoverStaleJobs(h.deps.handle.db, h.env.JOBS_STALE_MINUTES, new Date());
|
||||
expect(recovered).toBe(1);
|
||||
|
||||
const row = await h.deps.handle.db
|
||||
.selectFrom('jobs')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.executeTakeFirstOrThrow();
|
||||
expect(row.status).toBe('pending');
|
||||
expect(row.locked_by).toBeNull();
|
||||
await h.close();
|
||||
});
|
||||
|
||||
it('leaves a job that is merely slow alone', async () => {
|
||||
const h = await createHarness();
|
||||
const { id } = await enqueue(h.deps.handle.db, {
|
||||
type: 'classify_document',
|
||||
payload: { documentId: 'abc' },
|
||||
});
|
||||
await h.deps.handle.db
|
||||
.updateTable('jobs')
|
||||
.set({ status: 'running', locked_by: 'busy', locked_at: new Date().toISOString() })
|
||||
.where('id', '=', id)
|
||||
.execute();
|
||||
|
||||
expect(await recoverStaleJobs(h.deps.handle.db, 10, new Date())).toBe(0);
|
||||
await h.close();
|
||||
});
|
||||
|
||||
it('a restarted poller picks the recovered job up and runs it', async () => {
|
||||
const h = await createHarness();
|
||||
const document = await h.deps.handle.db
|
||||
.selectFrom('documents')
|
||||
.select('id')
|
||||
.where('status', '=', 'needs_review')
|
||||
.executeTakeFirstOrThrow();
|
||||
|
||||
const { id } = await enqueue(h.deps.handle.db, {
|
||||
type: 'classify_document',
|
||||
payload: { documentId: document.id },
|
||||
});
|
||||
await h.deps.handle.db
|
||||
.updateTable('jobs')
|
||||
.set({
|
||||
status: 'running',
|
||||
locked_by: 'worker-that-died',
|
||||
locked_at: new Date(Date.now() - 60 * 60_000).toISOString(),
|
||||
attempts: 1,
|
||||
})
|
||||
.where('id', '=', id)
|
||||
.execute();
|
||||
|
||||
// runOnce recovers stale rows first, which is what a fresh poller does on boot.
|
||||
expect(await h.runJobs()).toBe(1);
|
||||
|
||||
const row = await h.deps.handle.db
|
||||
.selectFrom('jobs')
|
||||
.selectAll()
|
||||
.where('id', '=', id)
|
||||
.executeTakeFirstOrThrow();
|
||||
expect(row.status).toBe('done');
|
||||
expect(row.attempts).toBe(2);
|
||||
await h.close();
|
||||
});
|
||||
});
|
||||
|
||||
describe('claiming', () => {
|
||||
it('hands a job to exactly one caller', async () => {
|
||||
const h = await createHarness();
|
||||
for (let i = 0; i < 5; i++) {
|
||||
await enqueue(h.deps.handle.db, { type: 'classify_document', payload: { documentId: `d${i}` } });
|
||||
}
|
||||
|
||||
const { claimJob } = await import('./claim');
|
||||
const claimed = new Set<string>();
|
||||
for (let i = 0; i < 5; i++) {
|
||||
const job = await claimJob(h.deps.handle.db, h.deps.handle.dialect, {
|
||||
instanceId: `worker-${i}`,
|
||||
now: new Date().toISOString(),
|
||||
});
|
||||
expect(job).not.toBeNull();
|
||||
expect(claimed.has(job!.id)).toBe(false);
|
||||
claimed.add(job!.id);
|
||||
}
|
||||
|
||||
// Nothing left to claim.
|
||||
expect(
|
||||
await claimJob(h.deps.handle.db, h.deps.handle.dialect, {
|
||||
instanceId: 'worker-late',
|
||||
now: new Date().toISOString(),
|
||||
}),
|
||||
).toBeNull();
|
||||
await h.close();
|
||||
});
|
||||
|
||||
it('does not claim a job scheduled for the future', async () => {
|
||||
const h = await createHarness();
|
||||
await enqueue(h.deps.handle.db, {
|
||||
type: 'classify_document',
|
||||
payload: {},
|
||||
runAt: new Date(Date.now() + 60_000),
|
||||
});
|
||||
|
||||
const { claimJob } = await import('./claim');
|
||||
expect(
|
||||
await claimJob(h.deps.handle.db, h.deps.handle.dialect, {
|
||||
instanceId: 'w',
|
||||
now: new Date().toISOString(),
|
||||
}),
|
||||
).toBeNull();
|
||||
await h.close();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,135 @@
|
||||
import type { Kysely } from 'kysely';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import type { Dialect } from '../../db/index';
|
||||
import type { Database } from '../../db/schema';
|
||||
import { claimJob, recoverStaleJobs, type JobRow } from './claim';
|
||||
import { nextRunAt, type JobType } from './queue';
|
||||
|
||||
export interface JobContext {
|
||||
id: string;
|
||||
type: JobType;
|
||||
payload: Record<string, unknown>;
|
||||
attempts: number;
|
||||
}
|
||||
|
||||
export type JobHandler = (context: JobContext) => Promise<void>;
|
||||
export type JobHandlers = Partial<Record<JobType, JobHandler>>;
|
||||
|
||||
export interface PollerOptions {
|
||||
db: Kysely<Database>;
|
||||
dialect: Dialect;
|
||||
handlers: JobHandlers;
|
||||
intervalMs: number;
|
||||
staleMinutes: number;
|
||||
/** Called for every terminal failure, so the admin error queue sees dead jobs. */
|
||||
onDead?: (job: JobRow, error: unknown) => Promise<void>;
|
||||
}
|
||||
|
||||
export interface Poller {
|
||||
start(): void;
|
||||
stop(): Promise<void>;
|
||||
/** Drains everything currently due. Used by tests and by the first tick after boot. */
|
||||
runOnce(): Promise<number>;
|
||||
readonly instanceId: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* One poller per process, guarded by the caller (SPEC.md section 10). Claiming is what
|
||||
* makes N of these safe against each other, not this loop.
|
||||
*/
|
||||
export function createPoller(options: PollerOptions): Poller {
|
||||
const instanceId = uuidv7();
|
||||
let timer: NodeJS.Timeout | null = null;
|
||||
let running = false;
|
||||
let stopped = false;
|
||||
let inFlight: Promise<unknown> = Promise.resolve();
|
||||
|
||||
async function runOne(): Promise<boolean> {
|
||||
const now = new Date();
|
||||
const job = await claimJob(options.db, options.dialect, {
|
||||
instanceId,
|
||||
now: now.toISOString(),
|
||||
});
|
||||
if (!job) return false;
|
||||
|
||||
const handler = options.handlers[job.type as JobType];
|
||||
try {
|
||||
if (!handler) throw new Error(`no handler registered for job type: ${job.type}`);
|
||||
await handler({
|
||||
id: job.id,
|
||||
type: job.type as JobType,
|
||||
payload: JSON.parse(job.payload) as Record<string, unknown>,
|
||||
attempts: job.attempts,
|
||||
});
|
||||
await options.db
|
||||
.updateTable('jobs')
|
||||
.set({ status: 'done', locked_by: null, locked_at: null, updated_at: new Date().toISOString() })
|
||||
.where('id', '=', job.id)
|
||||
.execute();
|
||||
} catch (error) {
|
||||
await fail(job, error);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
async function fail(job: JobRow, error: unknown): Promise<void> {
|
||||
const now = new Date();
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
const exhausted = job.attempts >= job.max_attempts;
|
||||
|
||||
await options.db
|
||||
.updateTable('jobs')
|
||||
.set({
|
||||
status: exhausted ? 'dead' : 'pending',
|
||||
run_at: exhausted ? job.run_at : nextRunAt(job.attempts, now),
|
||||
locked_by: null,
|
||||
locked_at: null,
|
||||
last_error: message.slice(0, 2000),
|
||||
updated_at: now.toISOString(),
|
||||
})
|
||||
.where('id', '=', job.id)
|
||||
.execute();
|
||||
|
||||
if (exhausted) await options.onDead?.(job, error);
|
||||
}
|
||||
|
||||
async function runOnce(): Promise<number> {
|
||||
await recoverStaleJobs(options.db, options.staleMinutes, new Date());
|
||||
|
||||
let processed = 0;
|
||||
// Bounded so one tick cannot monopolise the process when the queue is deep.
|
||||
while (processed < 50 && !stopped && (await runOne())) processed += 1;
|
||||
return processed;
|
||||
}
|
||||
|
||||
async function tick(): Promise<void> {
|
||||
if (running || stopped) return;
|
||||
running = true;
|
||||
inFlight = runOnce().catch((error: unknown) => {
|
||||
console.error('[jobs] poll failed', error);
|
||||
});
|
||||
try {
|
||||
await inFlight;
|
||||
} finally {
|
||||
running = false;
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
instanceId,
|
||||
start() {
|
||||
if (timer) return;
|
||||
timer = setInterval(() => void tick(), options.intervalMs);
|
||||
// Node should not stay alive purely to keep polling.
|
||||
timer.unref?.();
|
||||
void tick();
|
||||
},
|
||||
async stop() {
|
||||
stopped = true;
|
||||
if (timer) clearInterval(timer);
|
||||
timer = null;
|
||||
await inFlight;
|
||||
},
|
||||
runOnce,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
import type { Kysely } from 'kysely';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import type { Database } from '../../db/schema';
|
||||
|
||||
export type JobType =
|
||||
| 'ocr_extract'
|
||||
| 'classify_document'
|
||||
| 'verify_cdc'
|
||||
| 'generate_declaration_pdf'
|
||||
| 'send_notification'
|
||||
| 'deadline_sweep'
|
||||
| 'auto_confirm_sweep'
|
||||
| 'digest_sweep'
|
||||
| 'purge_user';
|
||||
|
||||
/** Retry schedule from SPEC.md section 10: 1m, 5m, 25m, 2h, 12h, then dead. */
|
||||
export const BACKOFF_MS = [60_000, 300_000, 1_500_000, 7_200_000, 43_200_000] as const;
|
||||
|
||||
export function nextRunAt(attempts: number, now: Date): string {
|
||||
const delay = BACKOFF_MS[Math.min(attempts, BACKOFF_MS.length) - 1] ?? BACKOFF_MS[0];
|
||||
return new Date(now.getTime() + delay).toISOString();
|
||||
}
|
||||
|
||||
export interface EnqueueOptions {
|
||||
type: JobType;
|
||||
payload: Record<string, unknown>;
|
||||
runAt?: Date;
|
||||
maxAttempts?: number;
|
||||
/**
|
||||
* Makes the enqueue idempotent. A job with the same type and key that is already
|
||||
* pending, running or done is not queued again, which is what keeps the daily sweeps
|
||||
* from piling up when a poller restarts (SPEC.md section 10).
|
||||
*/
|
||||
dedupeKey?: string;
|
||||
}
|
||||
|
||||
export async function enqueue(
|
||||
db: Kysely<Database>,
|
||||
options: EnqueueOptions,
|
||||
): Promise<{ id: string; deduped: boolean }> {
|
||||
const now = new Date();
|
||||
const payload = options.dedupeKey
|
||||
? { ...options.payload, dedupeKey: options.dedupeKey }
|
||||
: options.payload;
|
||||
|
||||
if (options.dedupeKey) {
|
||||
const existing = await db
|
||||
.selectFrom('jobs')
|
||||
.select('id')
|
||||
.where('type', '=', options.type)
|
||||
.where('status', 'in', ['pending', 'running', 'done'])
|
||||
.where('payload', 'like', `%"dedupeKey":"${options.dedupeKey}"%`)
|
||||
.executeTakeFirst();
|
||||
if (existing) return { id: existing.id, deduped: true };
|
||||
}
|
||||
|
||||
const id = uuidv7();
|
||||
await db
|
||||
.insertInto('jobs')
|
||||
.values({
|
||||
id,
|
||||
type: options.type,
|
||||
payload: JSON.stringify(payload),
|
||||
status: 'pending',
|
||||
run_at: (options.runAt ?? now).toISOString(),
|
||||
attempts: 0,
|
||||
max_attempts: options.maxAttempts ?? 5,
|
||||
locked_by: null,
|
||||
locked_at: null,
|
||||
last_error: null,
|
||||
created_at: now.toISOString(),
|
||||
updated_at: now.toISOString(),
|
||||
})
|
||||
.execute();
|
||||
|
||||
return { id, deduped: false };
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
/**
|
||||
* Everything that touches a stored file goes through this (SPEC.md section 7). Keys never
|
||||
* contain a user supplied filename: a scan is `<userId>/<uuid>`, so a hostile name cannot
|
||||
* escape a prefix or collide with someone else's file.
|
||||
*/
|
||||
export interface StorageDriver {
|
||||
readonly kind: 'local' | 's3';
|
||||
put(key: string, data: Uint8Array, mime: string): Promise<void>;
|
||||
getStream(key: string): Promise<ReadableStream<Uint8Array>>;
|
||||
delete(key: string): Promise<void>;
|
||||
/** A signed URL on S3, or the authenticated API route for local disk. */
|
||||
url(key: string): Promise<string>;
|
||||
/** Readiness probe: cheap, and it must fail when the backing store is unreachable. */
|
||||
check(): Promise<void>;
|
||||
}
|
||||
|
||||
export function storageKey(userId: string, id: string): string {
|
||||
return `${userId}/${id}`;
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
import type { Env } from '../../lib/env';
|
||||
import type { StorageDriver } from './driver';
|
||||
import { createLocalDriver } from './local';
|
||||
import { createS3Driver } from './s3';
|
||||
|
||||
export { type StorageDriver, storageKey } from './driver';
|
||||
|
||||
/** One switch, driven by STORAGE_DRIVER. Nothing else in the codebase branches on it. */
|
||||
export function createStorage(env: Env): StorageDriver {
|
||||
return env.STORAGE_DRIVER === 's3' ? createS3Driver(env) : createLocalDriver(env.STORAGE_LOCAL_PATH);
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
import { createReadStream } from 'node:fs';
|
||||
import { mkdir, rm, stat, writeFile } from 'node:fs/promises';
|
||||
import { dirname, join, resolve, sep } from 'node:path';
|
||||
import { Readable } from 'node:stream';
|
||||
import type { StorageDriver } from './driver';
|
||||
|
||||
/**
|
||||
* Local disk. Requires one shared volume across every replica, which the boot check in
|
||||
* lib/env.ts says out loud. Fine for the compose stack, not for k8s.
|
||||
*/
|
||||
export function createLocalDriver(root: string): StorageDriver {
|
||||
const base = resolve(root);
|
||||
|
||||
/** Refuses any key that would resolve outside the storage root. */
|
||||
function pathFor(key: string): string {
|
||||
const target = resolve(join(base, key));
|
||||
if (target !== base && !target.startsWith(base + sep)) {
|
||||
throw new Error(`storage key escapes the storage root: ${key}`);
|
||||
}
|
||||
return target;
|
||||
}
|
||||
|
||||
return {
|
||||
kind: 'local',
|
||||
|
||||
async put(key, data, _mime) {
|
||||
const target = pathFor(key);
|
||||
await mkdir(dirname(target), { recursive: true });
|
||||
await writeFile(target, data);
|
||||
},
|
||||
|
||||
async getStream(key) {
|
||||
const target = pathFor(key);
|
||||
await stat(target);
|
||||
return Readable.toWeb(createReadStream(target)) as ReadableStream<Uint8Array>;
|
||||
},
|
||||
|
||||
async delete(key) {
|
||||
await rm(pathFor(key), { force: true });
|
||||
},
|
||||
|
||||
async url(key) {
|
||||
// Local files are private, so this is the authenticated route rather than a path.
|
||||
return `/api/files/${encodeURIComponent(key)}`;
|
||||
},
|
||||
|
||||
async check() {
|
||||
await mkdir(base, { recursive: true });
|
||||
await stat(base);
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
import {
|
||||
DeleteObjectCommand,
|
||||
GetObjectCommand,
|
||||
HeadBucketCommand,
|
||||
PutObjectCommand,
|
||||
S3Client,
|
||||
} from '@aws-sdk/client-s3';
|
||||
import { getSignedUrl } from '@aws-sdk/s3-request-presigner';
|
||||
import type { Env } from '../../lib/env';
|
||||
import type { StorageDriver } from './driver';
|
||||
|
||||
const SIGNED_URL_TTL_SECONDS = 600;
|
||||
|
||||
/** Any S3 compatible endpoint: AWS, MinIO, R2, Backblaze. Path style for the rest. */
|
||||
export function createS3Driver(env: Env): StorageDriver {
|
||||
const bucket = env.S3_BUCKET;
|
||||
if (!bucket) throw new Error('S3_BUCKET is required for the s3 storage driver');
|
||||
|
||||
const client = new S3Client({
|
||||
region: env.S3_REGION ?? 'us-east-1',
|
||||
forcePathStyle: env.S3_FORCE_PATH_STYLE,
|
||||
...(env.S3_ENDPOINT ? { endpoint: env.S3_ENDPOINT } : {}),
|
||||
...(env.S3_ACCESS_KEY_ID && env.S3_SECRET_ACCESS_KEY
|
||||
? {
|
||||
credentials: {
|
||||
accessKeyId: env.S3_ACCESS_KEY_ID,
|
||||
secretAccessKey: env.S3_SECRET_ACCESS_KEY,
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
});
|
||||
|
||||
return {
|
||||
kind: 's3',
|
||||
|
||||
async put(key, data, mime) {
|
||||
await client.send(
|
||||
new PutObjectCommand({ Bucket: bucket, Key: key, Body: data, ContentType: mime }),
|
||||
);
|
||||
},
|
||||
|
||||
async getStream(key) {
|
||||
const result = await client.send(new GetObjectCommand({ Bucket: bucket, Key: key }));
|
||||
if (!result.Body) throw new Error(`object has no body: ${key}`);
|
||||
return result.Body.transformToWebStream();
|
||||
},
|
||||
|
||||
async delete(key) {
|
||||
await client.send(new DeleteObjectCommand({ Bucket: bucket, Key: key }));
|
||||
},
|
||||
|
||||
async url(key) {
|
||||
return getSignedUrl(client, new GetObjectCommand({ Bucket: bucket, Key: key }), {
|
||||
expiresIn: SIGNED_URL_TTL_SECONDS,
|
||||
});
|
||||
},
|
||||
|
||||
async check() {
|
||||
await client.send(new HeadBucketCommand({ Bucket: bucket }));
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -1,3 +1,6 @@
|
||||
import { mkdtempSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { createAuth } from '../auth/options';
|
||||
import { createDb } from '../db/index';
|
||||
import { migrateToLatest } from '../db/migrator';
|
||||
@@ -5,6 +8,9 @@ import { seed } from '../db/seed';
|
||||
import { createApp, type AppHandle } from '../http/app';
|
||||
import type { AppDeps } from '../http/context';
|
||||
import { type Env, parseEnv } from '../lib/env';
|
||||
import { createHandlers, createPoller, type Poller } from '../modules/jobs';
|
||||
import type { OcrProvider } from '../modules/documents/ocr';
|
||||
import { createStorage, type StorageDriver } from '../modules/storage';
|
||||
|
||||
export const TEST_ENV: Record<string, string> = {
|
||||
NODE_ENV: 'test',
|
||||
@@ -14,16 +20,27 @@ export const TEST_ENV: Record<string, string> = {
|
||||
APP_PUBLIC_URL: 'http://localhost:3000',
|
||||
};
|
||||
|
||||
export interface HarnessOptions {
|
||||
env?: Record<string, string>;
|
||||
/** Stand in for the Anthropic call, so OCR paths are testable without a key. */
|
||||
ocr?: OcrProvider | null;
|
||||
storage?: StorageDriver;
|
||||
}
|
||||
|
||||
export interface Harness extends AppHandle {
|
||||
deps: AppDeps;
|
||||
env: Env;
|
||||
storage: StorageDriver;
|
||||
/** Drains the job queue. Tests call this instead of waiting on the interval. */
|
||||
runJobs: () => Promise<number>;
|
||||
poller: Poller;
|
||||
close: () => Promise<void>;
|
||||
/** Signs in a seeded account and returns the cookie header for later requests. */
|
||||
signIn: (email: string, password: string) => Promise<string>;
|
||||
}
|
||||
|
||||
export async function createHarness(overrides: Record<string, string> = {}): Promise<Harness> {
|
||||
const parsed = parseEnv({ ...TEST_ENV, ...overrides });
|
||||
export async function createHarness(options: HarnessOptions = {}): Promise<Harness> {
|
||||
const parsed = parseEnv({ ...TEST_ENV, ...(options.env ?? {}) });
|
||||
if (!parsed.ok || !parsed.env) throw new Error(parsed.message);
|
||||
const env = parsed.env;
|
||||
|
||||
@@ -33,14 +50,34 @@ export async function createHarness(overrides: Record<string, string> = {}): Pro
|
||||
const auth = createAuth({ db: handle.db, dialect: handle.dialect, env, sendOtp: async () => undefined });
|
||||
await seed(handle, auth);
|
||||
|
||||
const deps: AppDeps = { env, handle, auth };
|
||||
const storage =
|
||||
options.storage ??
|
||||
createStorage({ ...env, STORAGE_DRIVER: 'local', STORAGE_LOCAL_PATH: mkdtempSync(join(tmpdir(), 'impuestos-test-')) });
|
||||
const ocr = options.ocr ?? null;
|
||||
const ingest = { db: handle.db, storage, ocrEnabled: ocr !== null };
|
||||
|
||||
const deps: AppDeps = { env, handle, auth, storage, ingest };
|
||||
const appHandle = createApp(deps);
|
||||
|
||||
const poller = createPoller({
|
||||
db: handle.db,
|
||||
dialect: handle.dialect,
|
||||
handlers: createHandlers({ db: handle.db, storage, ocr }),
|
||||
intervalMs: env.JOBS_POLL_INTERVAL_MS,
|
||||
staleMinutes: env.JOBS_STALE_MINUTES,
|
||||
});
|
||||
|
||||
return {
|
||||
...appHandle,
|
||||
deps,
|
||||
env,
|
||||
close: () => handle.close(),
|
||||
storage,
|
||||
poller,
|
||||
runJobs: () => poller.runOnce(),
|
||||
close: async () => {
|
||||
await poller.stop();
|
||||
await handle.close();
|
||||
},
|
||||
signIn: async (email, password) => {
|
||||
const response = await appHandle.app.request('/api/auth/sign-in/email', {
|
||||
method: 'POST',
|
||||
|
||||
Reference in New Issue
Block a user