phase-6: the staff console, and a log that says who did what

Flow H, three screens behind a role check: find an account, work the
ingestion error queue, read and export the audit log. Superadmins can
change a role, never their own.

The error queue merges ingest errors and dead jobs into one table with a
cursor that pages both sources; only a job can be retried and only an
ingest row resolved, with a note that migration 003 gives it somewhere
to live.

writeAudit no longer defaults a missing subject to the actor, which had
been recording a user search as staff looking themselves up. Omitting
the subject still means acting on yourself; null now means the action
has no subject, which is what a search, a retry and an export are.

Reading the log is not audited. Exporting it is: a copy leaving the
building is a different act from looking.

e2e/global-setup.ts asks for every screen once before the suite starts,
so a dev server's first-request compile is paid before the first test
rather than by it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Michilis
2026-09-04 21:55:59 +00:00
co-authored by Claude Opus 5
parent 6620650e9e
commit 4c39926483
39 changed files with 2767 additions and 11 deletions
@@ -0,0 +1,15 @@
import type { Kysely } from 'kysely';
/**
* SPEC-GAP: FLOWS.md Flow H has staff resolve an error "with note" and SPEC.md section 5
* gives `ingest_errors` nowhere to put it. The note is what the next person reads to find
* out what happened, so it is a column of its own rather than another key inside the
* payload blob the failure wrote.
*/
export async function up(db: Kysely<unknown>): Promise<void> {
await db.schema.alterTable('ingest_errors').addColumn('resolution_note', 'text').execute();
}
export async function down(db: Kysely<unknown>): Promise<void> {
await db.schema.alterTable('ingest_errors').dropColumn('resolution_note').execute();
}
+2
View File
@@ -1,6 +1,7 @@
import type { Migration, MigrationProvider } from 'kysely/migration';
import * as core from './001_core';
import * as insightDismissals from './002_insight_dismissals';
import * as errorResolutionNote from './003_error_resolution_note';
/**
* Migrations are listed statically rather than read from disk: the production image
@@ -9,6 +10,7 @@ import * as insightDismissals from './002_insight_dismissals';
const migrations: Record<string, Migration> = {
'001_core': core,
'002_insight_dismissals': insightDismissals,
'003_error_resolution_note': errorResolutionNote,
};
export const migrationProvider: MigrationProvider = {
+2
View File
@@ -228,6 +228,8 @@ export interface IngestErrorsTable {
status: 'open' | 'resolved';
resolved_by: string | null;
resolved_at: string | null;
/** What the staff member who closed the row said about it. */
resolution_note: string | null;
created_at: string;
}
+37
View File
@@ -80,6 +80,7 @@ export async function seed(
await seedDocuments(handle, options.now ?? new Date());
await seedDeclarations(handle, options.now ?? new Date());
await seedAuditTrail(handle);
await seedDeadJob(handle, options.now ?? new Date());
return result;
}
@@ -248,6 +249,42 @@ async function upsertProfileRow(handle: DbHandle, seed: ProfileSeed): Promise<vo
.execute();
}
/**
* One dead job, so the error queue has something from both of its sources and the retry
* action has something to act on. The two open ingest errors CONTRACTS.md section 4 asks
* for are written by seedDocuments, next to the documents they failed on.
*/
async function seedDeadJob(handle: DbHandle, now: Date): Promise<void> {
const db = handle.db;
const existing = await db
.selectFrom('jobs')
.select('id')
.where('status', '=', 'dead')
.executeTakeFirst();
if (existing) return;
const maria = await userIdFor(handle, 'maria@demo.local');
const deadAt = new Date(now.getTime() - 48 * 60 * 60 * 1000).toISOString();
await db
.insertInto('jobs')
.values({
id: uuidv7(),
type: 'verify_cdc',
payload: JSON.stringify({ userId: maria, cdc: '01801234567001001000000012024011512345678901' }),
status: 'dead',
run_at: deadAt,
attempts: 5,
max_attempts: 5,
locked_by: null,
locked_at: null,
last_error: 'dnit lookup timed out',
created_at: deadAt,
updated_at: deadAt,
})
.execute();
}
async function userIdFor(handle: DbHandle, email: string): Promise<string> {
const row = await handle.db
.selectFrom('user')
+2
View File
@@ -3,6 +3,7 @@ import type { AppDeps, AppEnv } from './context';
import { HttpError, toEnvelope } from './errors';
import { liveness, readiness } from './health';
import { localeMiddleware, sessionMiddleware } from './middleware';
import { adminRoutes } from './routes/admin';
import { dashboardRoutes, deadlineRoutes } from './routes/dashboard';
import { declarationRoutes } from './routes/declarations';
import { documentRoutes } from './routes/documents';
@@ -57,6 +58,7 @@ export function createApp(deps: AppDeps): AppHandle {
api.route('/dashboard', dashboardRoutes(deps));
api.route('/deadlines', deadlineRoutes(deps));
api.route('/declarations', declarationRoutes(deps));
api.route('/admin', adminRoutes(deps));
app.route('/api', api);
+224
View File
@@ -0,0 +1,224 @@
import {
AdminAuditQuery,
AdminErrorListQuery,
ResolveErrorInput,
RoleChangeInput,
} from '@impuestos/contracts';
import { Hono } from 'hono';
import type { z } from 'zod';
import { ADMIN_ROLES } from '../../auth/options';
import {
auditRowsForExport,
changeRole,
listAdminErrors,
listAudit,
resolveIngestError,
retryJob,
searchUsers,
toCsv,
userOverview,
} from '../../modules/admin';
import { type AuditAction, writeAudit } from '../../modules/audit';
import type { AppDeps, AppEnv, SessionUser } from '../context';
import { HttpError } from '../errors';
import { requireRole } from '../middleware';
const SUPERADMIN_ONLY = ['superadmin'] as const;
/**
* FLOWS.md Flow H. Two rules hold for every handler here: the role is checked in the
* handler and not only at the router, and anything that reads or changes a user's data
* writes an audit row before it answers.
*/
export function adminRoutes(deps: AppDeps): Hono<AppEnv> {
const routes = new Hono<AppEnv>();
const db = deps.handle.db;
routes.get('/users/search', async (c) => {
const staff = requireRole(c, ADMIN_ROLES);
const term = c.req.query('q') ?? '';
const items = await searchUsers(db, term);
await audit(c, deps, staff, {
action: 'admin.user_search',
subjectUserId: null,
resource: 'user',
detail: { q: term, results: items.length },
});
return c.json({ items });
});
routes.get('/users/:id/overview', async (c) => {
const staff = requireRole(c, ADMIN_ROLES);
const id = c.req.param('id');
const overview = await userOverview(db, id);
if (!overview) throw new HttpError('not_found');
// CONTRACTS.md section 5.5 pins this: one view, exactly one `admin.user_lookup` row.
await audit(c, deps, staff, {
action: 'admin.user_lookup',
subjectUserId: id,
resource: `user/${id}`,
});
return c.json(overview);
});
routes.post('/users/:id/role', async (c) => {
const actor = requireRole(c, SUPERADMIN_ONLY);
const id = c.req.param('id');
const input = parse(RoleChangeInput, await body(c));
const result = await changeRole(db, {
actorUserId: actor.id,
targetUserId: id,
role: input.role,
});
if (!result.ok) {
if (result.reason === 'not_found') throw new HttpError('not_found');
throw new HttpError('conflict', { field: 'role', detail: { reason: result.reason } });
}
await audit(c, deps, actor, {
action: 'admin.role_change',
subjectUserId: id,
resource: `user/${id}`,
detail: { from: result.previous, to: input.role },
});
return c.json({ ok: true } as const);
});
routes.get('/errors', async (c) => {
requireRole(c, ADMIN_ROLES);
const query = parse(AdminErrorListQuery, {
stage: c.req.query('stage'),
status: c.req.query('status'),
cursor: c.req.query('cursor'),
});
return c.json(await listAdminErrors(db, query));
});
routes.post('/errors/:id/resolve', async (c) => {
const staff = requireRole(c, ADMIN_ROLES);
const id = c.req.param('id');
const input = parse(ResolveErrorInput, await body(c));
const result = await resolveIngestError(db, id, { userId: staff.id, note: input.note });
if (!result.ok) {
if (result.reason === 'not_found') throw new HttpError('not_found');
throw new HttpError('conflict', { detail: { reason: result.reason } });
}
await audit(c, deps, staff, {
action: 'admin.error_resolve',
subjectUserId: null,
resource: `ingest_errors/${id}`,
detail: { note: input.note },
});
return c.json({ ok: true } as const);
});
routes.post('/jobs/:id/retry', async (c) => {
const staff = requireRole(c, ADMIN_ROLES);
const id = c.req.param('id');
const result = await retryJob(db, id);
if (!result.ok) {
if (result.reason === 'not_found') throw new HttpError('not_found');
throw new HttpError('conflict', { detail: { reason: result.reason } });
}
await audit(c, deps, staff, {
action: 'admin.job_retry',
subjectUserId: null,
resource: `jobs/${id}`,
});
return c.json({ ok: true } as const);
});
/**
* Reading the log is not itself audited. A row for every scroll of the audit screen
* would bury the accesses that matter under the act of looking for them. Taking a copy
* out of the building is a different thing, so the CSV export below is audited.
*/
routes.get('/audit', async (c) => {
requireRole(c, ADMIN_ROLES);
return c.json(await listAudit(db, auditQuery(c)));
});
routes.get('/audit/export.csv', async (c) => {
const staff = requireRole(c, ADMIN_ROLES);
const query = auditQuery(c);
const rows = await auditRowsForExport(db, query);
await audit(c, deps, staff, {
action: 'admin.audit_export',
subjectUserId: null,
resource: 'audit_log',
detail: { rows: rows.length, ...query },
});
return new Response(toCsv(rows), {
headers: {
'content-type': 'text/csv; charset=utf-8',
'content-disposition': `attachment; filename="audit-${new Date().toISOString().slice(0, 10)}.csv"`,
},
});
});
return routes;
}
function auditQuery(c: { req: { query: (name: string) => string | undefined } }): AdminAuditQuery {
return parse(AdminAuditQuery, {
actor: c.req.query('actor'),
action: c.req.query('action'),
subject: c.req.query('subject'),
from: c.req.query('from'),
to: c.req.query('to'),
cursor: c.req.query('cursor'),
});
}
async function body(c: { req: { json: () => Promise<unknown> } }): Promise<unknown> {
try {
return await c.req.json();
} catch {
throw new HttpError('validation_error');
}
}
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;
}
function audit(
c: { req: { header: (name: string) => string | undefined } },
deps: AppDeps,
actor: SessionUser,
entry: {
action: AuditAction;
subjectUserId: string | null;
resource: string;
detail?: Record<string, unknown>;
},
): Promise<void> {
return writeAudit(deps.handle.db, {
actorUserId: actor.id,
actorRole: actor.role,
ip: c.req.header('x-forwarded-for')?.split(',')[0]?.trim() ?? null,
...entry,
});
}
+11
View File
@@ -4,6 +4,7 @@ import {
DependentInput,
NotificationPrefsInput,
ProfileInput,
UserRole,
} from '@impuestos/contracts';
import { Hono } from 'hono';
import type { z } from 'zod';
@@ -28,6 +29,16 @@ export function meRoutes(deps: AppDeps): Hono<AppEnv> {
const routes = new Hono<AppEnv>();
const db = deps.handle.db;
// Who you are, for a client that needs the role before it renders (the admin console).
routes.get('/session', (c) => {
const user = requireUser(c);
// The session carries whatever string the user row holds. A value outside the enum
// would fail the client's schema and take a page down, so it reads as the least
// privileged role instead.
const role = UserRole.safeParse(user.role);
return c.json({ id: user.id, email: user.email, role: role.success ? role.data : 'user' });
});
// 404 until setup is complete: the client routes to onboarding (CONTRACTS.md section 3).
routes.get('/profile', async (c) => {
const user = requireUser(c);
+348
View File
@@ -0,0 +1,348 @@
import {
AdminAuditListDto,
AdminErrorListDto,
AdminUserOverviewDto,
AdminUserSearchDto,
} from '@impuestos/contracts';
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { createHarness, type Harness } from '../../test/harness';
let h: Harness;
let staff: string;
let superadmin: string;
let maria: string;
let mariaId: string;
beforeAll(async () => {
h = await createHarness();
staff = await h.signIn('staff@demo.local', 'demo-staff-1');
superadmin = await h.signIn('superadmin@demo.local', 'demo-superadmin-1');
maria = await h.signIn('maria@demo.local', 'demo-maria-1');
const row = await h.deps.handle.db
.selectFrom('user')
.select('id')
.where('email', '=', 'maria@demo.local')
.executeTakeFirstOrThrow();
mariaId = row.id;
});
afterAll(async () => {
await h.close();
});
const as = (cookie: string) => (path: string, init: RequestInit = {}) =>
h.app.request(path, {
...init,
headers: { cookie, 'content-type': 'application/json', ...(init.headers ?? {}) },
});
describe('who may reach the console', () => {
it('turns away an anonymous request', async () => {
expect((await h.app.request('/api/admin/users/search?q=maria')).status).toBe(401);
});
it('turns away a signed in user', async () => {
expect((await as(maria)('/api/admin/users/search?q=maria')).status).toBe(403);
expect((await as(maria)(`/api/admin/users/${mariaId}/overview`)).status).toBe(403);
expect((await as(maria)('/api/admin/errors')).status).toBe(403);
expect((await as(maria)('/api/admin/audit')).status).toBe(403);
});
});
describe('GET /admin/users/search', () => {
it('finds a user by email, by name and by document number', async () => {
for (const term of ['maria@demo', 'gonzalez', '4123456']) {
const response = await as(staff)(`/api/admin/users/search?q=${encodeURIComponent(term)}`);
expect(response.status, term).toBe(200);
const { items } = AdminUserSearchDto.parse(await response.json());
expect(items.map((item) => item.email), term).toContain('maria@demo.local');
}
});
it('shows the RUC with its check digit and the role', async () => {
const response = await as(staff)('/api/admin/users/search?q=maria@demo.local');
const { items } = AdminUserSearchDto.parse(await response.json());
expect(items).toHaveLength(1);
expect(items[0]?.doc).toBe('4123456-1');
expect(items[0]?.role).toBe('user');
expect(items[0]?.fullName).toBe('Maria Gonzalez');
});
it('returns nothing for a term too short to be a search', async () => {
const response = await as(staff)('/api/admin/users/search?q=a');
const { items } = AdminUserSearchDto.parse(await response.json());
expect(items).toEqual([]);
});
});
describe('GET /admin/users/:id/overview', () => {
it('counts what the account holds', async () => {
const response = await as(staff)(`/api/admin/users/${mariaId}/overview`);
expect(response.status).toBe(200);
const overview = AdminUserOverviewDto.parse(await response.json());
expect(overview.user.email).toBe('maria@demo.local');
expect(overview.profile?.fullName).toBe('Maria Gonzalez');
expect(overview.counts.documents).toBeGreaterThan(0);
expect(overview.counts.needsReview).toBeGreaterThan(0);
expect(overview.counts.declarations).toBeGreaterThan(0);
// The two seeded ingest errors belong to her.
expect(overview.recentErrors.length).toBeGreaterThanOrEqual(2);
expect(overview.lastActivityAt).not.toBeNull();
});
it('is a 404 for an id nobody has', async () => {
expect((await as(staff)('/api/admin/users/nope/overview')).status).toBe(404);
});
/** CONTRACTS.md section 5.5. */
it('writes exactly one admin.user_lookup row, visible on the audit endpoint', async () => {
const before = await auditCount('admin.user_lookup');
expect((await as(staff)(`/api/admin/users/${mariaId}/overview`)).status).toBe(200);
const response = await as(staff)('/api/admin/audit?action=admin.user_lookup');
const list = AdminAuditListDto.parse(await response.json());
expect(list.total).toBe(before + 1);
const newest = list.items[0];
expect(newest?.actorEmail).toBe('staff@demo.local');
expect(newest?.actorRole).toBe('staff');
expect(newest?.subjectEmail).toBe('maria@demo.local');
expect(newest?.resource).toBe(`user/${mariaId}`);
});
it('records a search separately from a view', async () => {
await as(staff)('/api/admin/users/search?q=maria');
const response = await as(staff)('/api/admin/audit?action=admin.user_search');
const list = AdminAuditListDto.parse(await response.json());
expect(list.total).toBeGreaterThan(0);
expect(list.items[0]?.detail?.['q']).toBe('maria');
// A search is about nobody in particular. Recording the searcher as its own subject
// would make the log say staff looked themselves up.
expect(list.items[0]?.subjectUserId).toBeNull();
});
});
describe('GET /admin/errors', () => {
it('merges ingest errors and dead jobs into one queue', async () => {
const response = await as(staff)('/api/admin/errors');
expect(response.status).toBe(200);
const list = AdminErrorListDto.parse(await response.json());
const stages = list.items.map((item) => item.stage);
expect(stages).toContain('ocr');
expect(stages).toContain('qr_parse');
expect(stages).toContain('job');
// The row expands to show what the failure captured.
const ocr = list.items.find((item) => item.stage === 'ocr');
expect(ocr?.payload?.['seeded']).toBe(true);
expect(ocr?.userEmail).toBe('maria@demo.local');
const job = list.items.find((item) => item.source === 'job');
expect(job?.attempts).toBe(5);
expect(job?.message).toContain('dnit lookup timed out');
});
it('filters by stage', async () => {
const response = await as(staff)('/api/admin/errors?stage=job');
const list = AdminErrorListDto.parse(await response.json());
expect(list.items.length).toBeGreaterThan(0);
expect(list.items.every((item) => item.source === 'job')).toBe(true);
});
it('rejects a stage that is not one of ours', async () => {
expect((await as(staff)('/api/admin/errors?stage=whatever')).status).toBe(400);
});
});
describe('POST /admin/errors/:id/resolve', () => {
it('closes the row with a note, once', async () => {
const list = AdminErrorListDto.parse(
await (await as(staff)('/api/admin/errors?stage=qr_parse&status=open')).json(),
);
const target = list.items[0];
expect(target).toBeDefined();
const response = await as(staff)(`/api/admin/errors/${target?.id}/resolve`, {
method: 'POST',
body: JSON.stringify({ note: 'supplier reissued the comprobante' }),
});
expect(response.status).toBe(200);
const row = await h.deps.handle.db
.selectFrom('ingest_errors')
.selectAll()
.where('id', '=', target?.id ?? '')
.executeTakeFirstOrThrow();
expect(row.status).toBe('resolved');
expect(row.resolution_note).toBe('supplier reissued the comprobante');
expect(row.resolved_by).not.toBeNull();
// Resolving it again is a conflict, not a silent no-op.
const again = await as(staff)(`/api/admin/errors/${target?.id}/resolve`, {
method: 'POST',
body: JSON.stringify({ note: 'again' }),
});
expect(again.status).toBe(409);
expect(await auditCount('admin.error_resolve')).toBe(1);
});
it('needs a note', async () => {
const response = await as(staff)('/api/admin/errors/whatever/resolve', {
method: 'POST',
body: JSON.stringify({ note: '' }),
});
expect(response.status).toBe(400);
});
});
describe('POST /admin/jobs/:id/retry', () => {
it('puts a dead job back in the queue and out of the error list', async () => {
const list = AdminErrorListDto.parse(
await (await as(staff)('/api/admin/errors?stage=job')).json(),
);
const job = list.items[0];
expect(job).toBeDefined();
const response = await as(staff)(`/api/admin/jobs/${job?.id}/retry`, { method: 'POST' });
expect(response.status).toBe(200);
const row = await h.deps.handle.db
.selectFrom('jobs')
.selectAll()
.where('id', '=', job?.id ?? '')
.executeTakeFirstOrThrow();
expect(row.status).toBe('pending');
// The full backoff schedule again, rather than dying on the first stumble.
expect(row.attempts).toBe(0);
const after = AdminErrorListDto.parse(
await (await as(staff)('/api/admin/errors?stage=job')).json(),
);
expect(after.items).toHaveLength(0);
expect(await auditCount('admin.job_retry')).toBe(1);
// Retrying something that is not dead is a conflict.
const again = await as(staff)(`/api/admin/jobs/${job?.id}/retry`, { method: 'POST' });
expect(again.status).toBe(409);
});
});
describe('POST /admin/users/:id/role', () => {
it('is refused to staff', async () => {
const response = await as(staff)(`/api/admin/users/${mariaId}/role`, {
method: 'POST',
body: JSON.stringify({ role: 'accountant' }),
});
expect(response.status).toBe(403);
});
it('refuses a superadmin changing their own role', async () => {
const self = await h.deps.handle.db
.selectFrom('user')
.select('id')
.where('email', '=', 'superadmin@demo.local')
.executeTakeFirstOrThrow();
const response = await as(superadmin)(`/api/admin/users/${self.id}/role`, {
method: 'POST',
body: JSON.stringify({ role: 'user' }),
});
expect(response.status).toBe(409);
});
it('promotes a user and records who did it', async () => {
const carlos = await h.deps.handle.db
.selectFrom('user')
.select('id')
.where('email', '=', 'carlos@demo.local')
.executeTakeFirstOrThrow();
const response = await as(superadmin)(`/api/admin/users/${carlos.id}/role`, {
method: 'POST',
body: JSON.stringify({ role: 'accountant' }),
});
expect(response.status).toBe(200);
const audit = AdminAuditListDto.parse(
await (await as(superadmin)('/api/admin/audit?action=admin.role_change')).json(),
);
const newest = audit.items[0];
expect(newest?.actorEmail).toBe('superadmin@demo.local');
expect(newest?.subjectEmail).toBe('carlos@demo.local');
expect(newest?.detail).toEqual({ from: 'user', to: 'accountant' });
// Setting the role it already has changes nothing and says so.
const again = await as(superadmin)(`/api/admin/users/${carlos.id}/role`, {
method: 'POST',
body: JSON.stringify({ role: 'accountant' }),
});
expect(again.status).toBe(409);
});
});
describe('the audit view', () => {
it('filters by actor email and by date', async () => {
const byActor = AdminAuditListDto.parse(
await (await as(staff)('/api/admin/audit?actor=staff@demo.local')).json(),
);
expect(byActor.total).toBeGreaterThan(0);
expect(byActor.items.every((item) => item.actorEmail === 'staff@demo.local')).toBe(true);
// A window that closed before anything happened holds nothing.
const empty = AdminAuditListDto.parse(
await (await as(staff)('/api/admin/audit?from=2020-01-01&to=2020-01-02')).json(),
);
expect(empty.items).toEqual([]);
expect(empty.total).toBe(0);
});
it('exports CSV with a header row and quoted fields', async () => {
const response = await as(staff)('/api/admin/audit/export.csv?action=admin.user_lookup');
expect(response.status).toBe(200);
expect(response.headers.get('content-type')).toContain('text/csv');
expect(response.headers.get('content-disposition')).toContain('attachment');
const lines = (await response.text()).split('\r\n');
expect(lines[0]).toBe(
'"id","created_at","actor_email","actor_role","action","subject_email","resource","ip","detail"',
);
expect(lines.length).toBeGreaterThan(1);
expect(lines[1]).toContain('"admin.user_lookup"');
});
it('audits the export, since a copy leaves the building', async () => {
const before = await auditCount('admin.audit_export');
await as(staff)('/api/admin/audit/export.csv');
expect(await auditCount('admin.audit_export')).toBe(before + 1);
});
it('does not audit merely reading the log', async () => {
const before = await h.deps.handle.db
.selectFrom('audit_log')
.select((eb) => eb.fn.countAll<number>().as('total'))
.executeTakeFirstOrThrow();
await as(staff)('/api/admin/audit');
const after = await h.deps.handle.db
.selectFrom('audit_log')
.select((eb) => eb.fn.countAll<number>().as('total'))
.executeTakeFirstOrThrow();
expect(Number(after.total)).toBe(Number(before.total));
});
});
async function auditCount(action: string): Promise<number> {
const row = await h.deps.handle.db
.selectFrom('audit_log')
.select((eb) => eb.fn.countAll<number>().as('total'))
.where('action', '=', action)
.executeTakeFirstOrThrow();
return Number(row.total);
}
+169
View File
@@ -0,0 +1,169 @@
import { type AdminAuditDto, type AdminAuditQuery, PAGE_SIZE } from '@impuestos/contracts';
import type { Kysely } from 'kysely';
import type { Database } from '../../db/schema';
/** The CSV export walks the whole result set, so it needs a ceiling that a browser can open. */
const EXPORT_LIMIT = 10_000;
/**
* Read only view over the append only log (SPEC.md section 14). `actor` and `subject`
* match on email or id, because staff have an email in front of them and the log stores
* an id.
*/
export async function listAudit(
db: Kysely<Database>,
query: AdminAuditQuery,
): Promise<{ items: AdminAuditDto[]; total: number; cursor?: string }> {
const base = filtered(db, query);
const { total } = await base
.select((eb) => eb.fn.countAll<number>().as('total'))
.executeTakeFirstOrThrow();
let page = selectRows(base).limit(PAGE_SIZE + 1);
if (query.cursor) page = page.where('audit_log.id', '<', query.cursor);
const rows = await page.execute();
const hasMore = rows.length > PAGE_SIZE;
const visible = hasMore ? rows.slice(0, PAGE_SIZE) : rows;
const last = visible.at(-1);
return {
items: visible.map(toDto),
total: Number(total),
...(hasMore && last ? { cursor: last.id } : {}),
};
}
export async function auditRowsForExport(
db: Kysely<Database>,
query: AdminAuditQuery,
): Promise<AdminAuditDto[]> {
const rows = await selectRows(filtered(db, query)).limit(EXPORT_LIMIT).execute();
return rows.map(toDto);
}
/** RFC 4180: quote every field, double the quotes inside it. Excel opens this. */
export function toCsv(rows: AdminAuditDto[]): string {
const header = [
'id',
'created_at',
'actor_email',
'actor_role',
'action',
'subject_email',
'resource',
'ip',
'detail',
];
const lines = rows.map((row) =>
[
row.id,
row.createdAt,
row.actorEmail ?? '',
row.actorRole,
row.action,
row.subjectEmail ?? '',
row.resource,
row.ip ?? '',
row.detail ? JSON.stringify(row.detail) : '',
]
.map(escape)
.join(','),
);
return [header.map(escape).join(','), ...lines].join('\r\n');
}
function escape(value: string): string {
return `"${value.replaceAll('"', '""')}"`;
}
function filtered(db: Kysely<Database>, query: AdminAuditQuery) {
let base = db.selectFrom('audit_log');
if (query.action) base = base.where('audit_log.action', '=', query.action);
if (query.actor) {
const actor = query.actor;
base = base.where((eb) =>
eb.or([
eb('audit_log.actor_user_id', '=', actor),
eb(
'audit_log.actor_user_id',
'in',
eb.selectFrom('user').select('user.id').where('user.email', '=', actor),
),
]),
);
}
if (query.subject) {
const subject = query.subject;
base = base.where((eb) =>
eb.or([
eb('audit_log.subject_user_id', '=', subject),
eb(
'audit_log.subject_user_id',
'in',
eb.selectFrom('user').select('user.id').where('user.email', '=', subject),
),
]),
);
}
// Timestamps are ISO-8601 text, so a date bound is a string comparison. `to` is
// inclusive of the whole day, which is what a person picking a date means.
if (query.from) base = base.where('audit_log.created_at', '>=', `${query.from}T00:00:00.000Z`);
if (query.to) base = base.where('audit_log.created_at', '<=', `${query.to}T23:59:59.999Z`);
return base;
}
function selectRows(base: ReturnType<typeof filtered>) {
return base
.leftJoin('user as actor', 'actor.id', 'audit_log.actor_user_id')
.leftJoin('user as subject', 'subject.id', 'audit_log.subject_user_id')
.selectAll('audit_log')
.select(['actor.email as actorEmail', 'subject.email as subjectEmail'])
.orderBy('audit_log.created_at', 'desc')
.orderBy('audit_log.id', 'desc');
}
interface AuditRow {
id: string;
actor_user_id: string;
actorEmail: string | null;
actor_role: string;
action: string;
subject_user_id: string | null;
subjectEmail: string | null;
resource: string;
detail: string | null;
ip: string | null;
created_at: string;
}
function toDto(row: AuditRow): AdminAuditDto {
let detail: Record<string, unknown> | null = null;
if (row.detail) {
try {
const parsed: unknown = JSON.parse(row.detail);
if (typeof parsed === 'object' && parsed !== null && !Array.isArray(parsed)) {
detail = parsed as Record<string, unknown>;
}
} catch {
detail = null;
}
}
return {
id: row.id,
actorUserId: row.actor_user_id,
actorEmail: row.actorEmail,
actorRole: row.actor_role,
action: row.action,
subjectUserId: row.subject_user_id,
subjectEmail: row.subjectEmail,
resource: row.resource,
detail,
ip: row.ip,
createdAt: row.created_at,
};
}
+232
View File
@@ -0,0 +1,232 @@
import { type AdminErrorDto, type AdminErrorListQuery, PAGE_SIZE } from '@impuestos/contracts';
import type { Kysely } from 'kysely';
import type { Database } from '../../db/schema';
/**
* The error queue (FLOWS.md Flow H). Two sources feed one table: rows the ingestion
* pipeline wrote, and jobs that exhausted their retries. They are merged here rather than
* in the UI so that paging and filtering mean the same thing for both.
*
* The cursor is `createdAt|id` of the last row on the page. Both sources are filtered by
* it before the merge, so a row is never shown twice and never skipped.
*/
export async function listAdminErrors(
db: Kysely<Database>,
query: AdminErrorListQuery,
): Promise<{ items: AdminErrorDto[]; total: number; cursor?: string }> {
const after = parseCursor(query.cursor);
const wantsIngest = query.stage !== 'job';
// A dead job is a live problem, so it never appears under the resolved filter: retrying
// it takes it out of the queue entirely.
const wantsJobs =
(query.stage === undefined || query.stage === 'job') && query.status !== 'resolved';
const [ingest, ingestTotal] = wantsIngest
? await ingestPage(db, query, after)
: [[] as AdminErrorDto[], 0];
const [jobs, jobTotal] = wantsJobs ? await deadJobPage(db, after) : [[] as AdminErrorDto[], 0];
const merged = [...ingest, ...jobs].sort(byNewest);
const visible = merged.slice(0, PAGE_SIZE);
const last = visible.at(-1);
const hasMore = merged.length > PAGE_SIZE;
return {
items: visible,
total: ingestTotal + jobTotal,
...(hasMore && last ? { cursor: `${last.createdAt}|${last.id}` } : {}),
};
}
async function ingestPage(
db: Kysely<Database>,
query: AdminErrorListQuery,
after: { createdAt: string; id: string } | null,
): Promise<[AdminErrorDto[], number]> {
let base = db.selectFrom('ingest_errors');
if (query.stage) base = base.where('stage', '=', query.stage);
if (query.status) base = base.where('status', '=', query.status);
const { total } = await base
.select((eb) => eb.fn.countAll<number>().as('total'))
.executeTakeFirstOrThrow();
let page = base
.leftJoin('user', 'user.id', 'ingest_errors.user_id')
.selectAll('ingest_errors')
.select('user.email as userEmail')
.orderBy('ingest_errors.created_at', 'desc')
.orderBy('ingest_errors.id', 'desc')
.limit(PAGE_SIZE + 1);
if (after) {
page = page.where((eb) =>
eb.or([
eb('ingest_errors.created_at', '<', after.createdAt),
eb.and([
eb('ingest_errors.created_at', '=', after.createdAt),
eb('ingest_errors.id', '<', after.id),
]),
]),
);
}
const rows = await page.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,
source: 'ingest' as const,
payload: parseJson(row.payload),
userEmail: row.userEmail ?? null,
attempts: null,
resolvedAt: row.resolved_at,
})),
Number(total),
];
}
/** A job that ran out of attempts is a support problem, so it joins the same queue. */
async function deadJobPage(
db: Kysely<Database>,
after: { createdAt: string; id: string } | null,
): Promise<[AdminErrorDto[], number]> {
const base = db.selectFrom('jobs').where('status', '=', 'dead');
const { total } = await base
.select((eb) => eb.fn.countAll<number>().as('total'))
.executeTakeFirstOrThrow();
let page = base.selectAll().orderBy('created_at', 'desc').orderBy('id', 'desc').limit(PAGE_SIZE + 1);
if (after) {
page = page.where((eb) =>
eb.or([
eb('created_at', '<', after.createdAt),
eb.and([eb('created_at', '=', after.createdAt), eb('id', '<', after.id)]),
]),
);
}
const rows = await page.execute();
const items = await Promise.all(
rows.map(async (row): Promise<AdminErrorDto> => {
const payload = parseJson(row.payload);
const userId = typeof payload?.['userId'] === 'string' ? payload['userId'] : null;
const documentId = typeof payload?.['documentId'] === 'string' ? payload['documentId'] : null;
return {
id: row.id,
userId,
documentId,
stage: 'job',
// The job type alone when nothing was captured: an invented sentence here would
// be copy, and copy belongs in the catalogs.
message: row.last_error ? `${row.type}: ${row.last_error}` : row.type,
status: 'open',
createdAt: row.created_at,
source: 'job',
payload,
userEmail: userId ? await emailFor(db, userId) : null,
attempts: row.attempts,
resolvedAt: null,
};
}),
);
return [items, Number(total)];
}
export interface ResolveResult {
ok: boolean;
reason?: 'not_found' | 'already_resolved';
}
export async function resolveIngestError(
db: Kysely<Database>,
id: string,
by: { userId: string; note: string },
): Promise<ResolveResult> {
const row = await db
.selectFrom('ingest_errors')
.select(['id', 'status'])
.where('id', '=', id)
.executeTakeFirst();
if (!row) return { ok: false, reason: 'not_found' };
if (row.status === 'resolved') return { ok: false, reason: 'already_resolved' };
await db
.updateTable('ingest_errors')
.set({
status: 'resolved',
resolved_by: by.userId,
resolved_at: new Date().toISOString(),
resolution_note: by.note,
})
.where('id', '=', id)
.execute();
return { ok: true };
}
export interface RetryResult {
ok: boolean;
reason?: 'not_found' | 'not_retryable';
}
/**
* Puts a dead job back at the front of the queue. Attempts reset to zero so the retry
* gets the whole backoff schedule again rather than dying on its first stumble.
*/
export async function retryJob(db: Kysely<Database>, id: string): Promise<RetryResult> {
const row = await db
.selectFrom('jobs')
.select(['id', 'status'])
.where('id', '=', id)
.executeTakeFirst();
if (!row) return { ok: false, reason: 'not_found' };
if (row.status !== 'dead' && row.status !== 'failed') return { ok: false, reason: 'not_retryable' };
const now = new Date().toISOString();
await db
.updateTable('jobs')
.set({ status: 'pending', run_at: now, attempts: 0, locked_by: null, locked_at: null, updated_at: now })
.where('id', '=', id)
.execute();
return { ok: true };
}
async function emailFor(db: Kysely<Database>, userId: string): Promise<string | null> {
const row = await db
.selectFrom('user')
.select('email')
.where('id', '=', userId)
.executeTakeFirst();
return row?.email ?? null;
}
function byNewest(a: AdminErrorDto, b: AdminErrorDto): number {
if (a.createdAt !== b.createdAt) return a.createdAt < b.createdAt ? 1 : -1;
return a.id < b.id ? 1 : -1;
}
function parseCursor(cursor: string | undefined): { createdAt: string; id: string } | null {
if (!cursor) return null;
const [createdAt, id] = cursor.split('|');
return createdAt && id ? { createdAt, id } : null;
}
function parseJson(value: string | null): Record<string, unknown> | null {
if (!value) return null;
try {
const parsed: unknown = JSON.parse(value);
return typeof parsed === 'object' && parsed !== null && !Array.isArray(parsed)
? (parsed as Record<string, unknown>)
: null;
} catch {
return null;
}
}
+10
View File
@@ -0,0 +1,10 @@
export { searchUsers, findUser, userOverview } from './users';
export {
listAdminErrors,
resolveIngestError,
retryJob,
type ResolveResult,
type RetryResult,
} from './errors';
export { listAudit, auditRowsForExport, toCsv } from './audit';
export { changeRole, type RoleChangeResult } from './roles';
+41
View File
@@ -0,0 +1,41 @@
import type { UserRole } from '@impuestos/contracts';
import type { Kysely } from 'kysely';
import type { Database } from '../../db/schema';
export interface RoleChangeResult {
ok: boolean;
reason?: 'not_found' | 'self' | 'unchanged';
previous?: string;
}
/**
* Superadmin only, and never on yourself: the one thing worse than an account with too
* much power is the last superadmin demoting themselves out of the console.
*
* No session shuffling is needed. `getSession` reads the role off the user row on every
* request, so a demotion takes effect on the next call the demoted user makes.
*/
export async function changeRole(
db: Kysely<Database>,
args: { actorUserId: string; targetUserId: string; role: UserRole },
): Promise<RoleChangeResult> {
if (args.actorUserId === args.targetUserId) return { ok: false, reason: 'self' };
const target = await db
.selectFrom('user')
.select(['id', 'role'])
.where('id', '=', args.targetUserId)
.executeTakeFirst();
if (!target) return { ok: false, reason: 'not_found' };
const previous = target.role ?? 'user';
if (previous === args.role) return { ok: false, reason: 'unchanged', previous };
await db
.updateTable('user')
.set({ role: args.role, updatedAt: new Date().toISOString() })
.where('id', '=', args.targetUserId)
.execute();
return { ok: true, previous };
}
+152
View File
@@ -0,0 +1,152 @@
import type { AdminUserDto, AdminUserOverviewDto } from '@impuestos/contracts';
import type { Kysely } from 'kysely';
import type { Database } from '../../db/schema';
import { listIngestErrorsForUser } from '../ingest/errors';
import { getProfile } from '../pii';
/** A search that returns everyone is not a search: staff have to name who they are looking for. */
const MIN_QUERY_LENGTH = 2;
const MAX_RESULTS = 25;
/**
* Finds users by email, name or document number (FLOWS.md Flow H). Matching is
* case insensitive on both sides so that "GONZALEZ" and "gonzalez" behave the same, and
* punctuation in a RUC is ignored: staff read numbers off a screen, not out of the table.
*/
export async function searchUsers(db: Kysely<Database>, term: string): Promise<AdminUserDto[]> {
const trimmed = term.trim();
if (trimmed.length < MIN_QUERY_LENGTH) return [];
const like = `%${trimmed.toLowerCase()}%`;
const digits = trimmed.replace(/\D/g, '');
const rows = await db
.selectFrom('user')
.leftJoin('profiles', 'profiles.user_id', 'user.id')
.select([
'user.id as id',
'user.email as email',
'user.name as name',
'user.role as role',
'user.createdAt as createdAt',
'profiles.full_name as fullName',
'profiles.ruc as ruc',
'profiles.ruc_dv as rucDv',
'profiles.ci as ci',
])
.where((eb) => {
const clauses = [
eb(eb.fn('lower', ['user.email']), 'like', like),
eb(eb.fn('lower', ['user.name']), 'like', like),
eb(eb.fn('lower', ['profiles.full_name']), 'like', like),
];
// An empty digit string would match every RUC, which is the opposite of a search.
if (digits.length > 0) {
clauses.push(eb('profiles.ruc', 'like', `%${digits}%`));
clauses.push(eb('profiles.ci', 'like', `%${digits}%`));
}
return eb.or(clauses);
})
.orderBy('user.createdAt', 'asc')
.limit(MAX_RESULTS)
.execute();
return rows.map(toAdminUser);
}
export async function findUser(db: Kysely<Database>, id: string): Promise<AdminUserDto | null> {
const row = await db
.selectFrom('user')
.leftJoin('profiles', 'profiles.user_id', 'user.id')
.select([
'user.id as id',
'user.email as email',
'user.name as name',
'user.role as role',
'user.createdAt as createdAt',
'profiles.full_name as fullName',
'profiles.ruc as ruc',
'profiles.ruc_dv as rucDv',
'profiles.ci as ci',
])
.where('user.id', '=', id)
.executeTakeFirst();
return row ? toAdminUser(row) : null;
}
/** Everything the support screen shows about one account, in one round trip. */
export async function userOverview(
db: Kysely<Database>,
id: string,
): Promise<AdminUserOverviewDto | null> {
const user = await findUser(db, id);
if (!user) return null;
const counts = await db
.selectFrom('documents')
.select((eb) => [
eb.fn.countAll<number>().as('documents'),
eb.fn
.sum<number>(eb.case().when('status', '=', 'needs_review').then(1).else(0).end())
.as('needsReview'),
])
.where('user_id', '=', id)
.executeTakeFirstOrThrow();
const declarations = await db
.selectFrom('declarations')
.select((eb) => eb.fn.countAll<number>().as('total'))
.where('user_id', '=', id)
.executeTakeFirstOrThrow();
const lastDocument = await db
.selectFrom('documents')
.select('created_at')
.where('user_id', '=', id)
.orderBy('created_at', 'desc')
.executeTakeFirst();
return {
user,
profile: await getProfile(db, id),
counts: {
documents: Number(counts.documents),
needsReview: Number(counts.needsReview ?? 0),
declarations: Number(declarations.total),
},
recentErrors: await listIngestErrorsForUser(db, id),
lastActivityAt: lastDocument?.created_at ?? null,
};
}
interface UserRow {
id: string;
email: string;
name: string;
role: string | null;
createdAt: string;
fullName: string | null;
ruc: string | null;
rucDv: string | null;
ci: string | null;
}
function toAdminUser(row: UserRow): AdminUserDto {
const doc = row.ruc ? `${row.ruc}-${row.rucDv ?? ''}`.replace(/-$/, '') : row.ci;
return {
id: row.id,
email: row.email,
fullName: row.fullName ?? row.name,
doc: doc ?? null,
// better-auth leaves the column null for an account created before a role was set.
role: isRole(row.role) ? row.role : 'user',
createdAt: row.createdAt,
};
}
const ROLE_VALUES = ['user', 'accountant', 'staff', 'superadmin'] as const;
function isRole(value: string | null): value is (typeof ROLE_VALUES)[number] {
return value !== null && (ROLE_VALUES as readonly string[]).includes(value);
}
+15 -3
View File
@@ -17,16 +17,28 @@ export type AuditAction =
| 'notification_prefs.update'
| 'data.export'
| 'account.delete'
| 'admin.user_search'
/**
* Reading one user's data. CONTRACTS.md section 5.5 pins this action to the overview
* endpoint, so the name stays even though `admin.user_search` reads more naturally
* next to it.
*/
| 'admin.user_lookup'
| 'admin.user_view'
| 'admin.role_change'
| 'admin.error_resolve'
| 'admin.job_retry'
| 'admin.audit_export'
| 'admin.file_access';
export interface AuditEntry {
actorUserId: string;
actorRole: string;
action: AuditAction;
/** Whose data this touched. Equal to the actor for a user acting on themselves. */
/**
* Whose data this touched. Omit it for a user acting on themselves and it defaults to
* the actor; pass `null` when the action has no subject at all, such as a search or an
* export. The two are different facts and the log must not blur them.
*/
subjectUserId?: string | null;
resource: string;
detail?: Record<string, unknown> | undefined;
@@ -41,7 +53,7 @@ export async function writeAudit(db: Kysely<Database>, entry: AuditEntry): Promi
actor_user_id: entry.actorUserId,
actor_role: entry.actorRole,
action: entry.action,
subject_user_id: entry.subjectUserId ?? entry.actorUserId,
subject_user_id: entry.subjectUserId === undefined ? entry.actorUserId : entry.subjectUserId,
resource: entry.resource,
detail: entry.detail === undefined ? null : JSON.stringify(entry.detail),
ip: entry.ip ?? null,
+1
View File
@@ -32,6 +32,7 @@ export async function recordIngestError(
status: 'open',
resolved_by: null,
resolved_at: null,
resolution_note: null,
created_at: new Date().toISOString(),
})
.execute();
+14 -3
View File
@@ -61,14 +61,20 @@ describe('backoff', () => {
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, {
// The seed leaves a dead job behind for the admin error queue, so this test names the
// row it queued rather than assuming it is the only one in the table.
const { id } = 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();
let row = await h.deps.handle.db
.selectFrom('jobs')
.selectAll()
.where('id', '=', id)
.executeTakeFirstOrThrow();
expect(row.status).toBe('pending');
expect(row.attempts).toBe(1);
expect(row.last_error).toContain('documentId');
@@ -79,10 +85,15 @@ describe('failure handling', () => {
await h.deps.handle.db
.updateTable('jobs')
.set({ run_at: new Date(Date.now() - 1000).toISOString() })
.where('id', '=', id)
.execute();
await h.runJobs();
row = await h.deps.handle.db.selectFrom('jobs').selectAll().executeTakeFirstOrThrow();
row = await h.deps.handle.db
.selectFrom('jobs')
.selectAll()
.where('id', '=', id)
.executeTakeFirstOrThrow();
expect(row.status).toBe('dead');
expect(row.attempts).toBe(2);
await h.close();