phase-8: run it on Postgres, and find out what that was hiding
Six phases claimed the product runs on SQLite and on Postgres. Nothing had ever run it on Postgres. TEST_DATABASE_URL now points the whole suite at a real server and CI runs both arms. The first run found a bug that would have shipped. better-auth's banned flag is integer 0/1 on SQLite and a real boolean on Postgres, and the code read it as `banned === 1`, so on Postgres an account someone asked us to freeze went on receiving email. Both of those flags are now typed for either dialect and read through isFlagSet. It also found the migration advisory lock being taken on a pool. An advisory lock belongs to the session that took it, so a lock on one pooled connection and an unlock on another leaves it held. Two replicas migrating at once is the ordinary case in k8s and is precisely what it was there to protect. New for the scaled mode: k8s manifests with migrations as an initContainer, one poller in its own worker Deployment rather than one per API replica, and an Ingress that exposes the web app only. The storage drivers finally have tests, S3 included, since scaled mode requires it and it had never been exercised. Two acceptance tests, both checked against a deliberately broken build first: two workers claiming a hundred jobs report 188 claims with SKIP LOCKED removed, and the in-flight request is cut off with the drain wait removed. 340 tests on SQLite, 341 on Postgres, 70 Playwright, rules coverage 100%. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
8c8463924f
commit
37370dd079
@@ -45,6 +45,10 @@ S3_FORCE_PATH_STYLE=true
|
||||
ANTHROPIC_API_KEY=
|
||||
OCR_MODEL=claude-sonnet-4-6
|
||||
|
||||
# Postgres connections per process. Multiply by the number of replicas and keep the total
|
||||
# under the server's max_connections. Ignored on SQLite.
|
||||
DATABASE_POOL_MAX=10
|
||||
|
||||
# Optional. Web push is hidden in the UI when unset. Generate: npx web-push generate-vapid-keys
|
||||
PUSH_VAPID_PUBLIC_KEY=
|
||||
PUSH_VAPID_PRIVATE_KEY=
|
||||
|
||||
@@ -20,9 +20,10 @@ export function dialectOf(databaseUrl: string): Dialect {
|
||||
);
|
||||
}
|
||||
|
||||
export function createDb(databaseUrl: string): DbHandle {
|
||||
export function createDb(databaseUrl: string, options: { poolMax?: number } = {}): DbHandle {
|
||||
const dialect = dialectOf(databaseUrl);
|
||||
const handle = dialect === 'sqlite' ? createSqliteDb(databaseUrl) : createPostgresDb(databaseUrl);
|
||||
const handle =
|
||||
dialect === 'sqlite' ? createSqliteDb(databaseUrl) : createPostgresDb(databaseUrl, options);
|
||||
return { ...handle, dialect };
|
||||
}
|
||||
|
||||
|
||||
@@ -8,11 +8,14 @@ import type { Database } from './schema';
|
||||
* Timestamps are stored as ISO-8601 text in our own tables, so the driver is told to
|
||||
* hand back `numeric` as a number and nothing else needs a type parser.
|
||||
*/
|
||||
export function createPostgresDb(databaseUrl: string): {
|
||||
export function createPostgresDb(
|
||||
databaseUrl: string,
|
||||
options: { poolMax?: number } = {},
|
||||
): {
|
||||
db: Kysely<Database>;
|
||||
close: () => Promise<void>;
|
||||
} {
|
||||
const pool = new pg.Pool({ connectionString: databaseUrl, max: 10 });
|
||||
const pool = new pg.Pool({ connectionString: databaseUrl, max: options.poolMax ?? 10 });
|
||||
const db = new Kysely<Database>({ dialect: new PostgresDialect({ pool }) });
|
||||
return {
|
||||
db,
|
||||
@@ -45,10 +48,23 @@ pg.types.setTypeParser(pg.types.builtins.INT8, (value) => {
|
||||
* api pod, so two migrators can start at the same moment.
|
||||
*/
|
||||
export async function withMigrationLock<T>(db: Kysely<Database>, fn: () => Promise<T>): Promise<T> {
|
||||
await sql`select pg_advisory_lock(${sql.lit(MIGRATION_ADVISORY_LOCK_KEY)})`.execute(db);
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
await sql`select pg_advisory_unlock(${sql.lit(MIGRATION_ADVISORY_LOCK_KEY)})`.execute(db);
|
||||
}
|
||||
/**
|
||||
* On one pinned connection, not on the pool. An advisory lock belongs to the session that
|
||||
* took it: issue the lock on one pooled connection and the unlock on another and the
|
||||
* unlock is a no-op with a warning, leaving the lock held until that connection happens
|
||||
* to close. The next replica to start then waits on a lock nobody holds any more.
|
||||
*
|
||||
* `fn` still uses the pool for its own work, which is what makes this mutual exclusion
|
||||
* rather than a single threaded migration.
|
||||
*/
|
||||
return db.connection().execute(async (connection) => {
|
||||
await sql`select pg_advisory_lock(${sql.lit(MIGRATION_ADVISORY_LOCK_KEY)})`.execute(connection);
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
await sql`select pg_advisory_unlock(${sql.lit(MIGRATION_ADVISORY_LOCK_KEY)})`.execute(
|
||||
connection,
|
||||
);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -37,20 +37,35 @@ export interface Database {
|
||||
insight_dismissals: InsightDismissalsTable;
|
||||
}
|
||||
|
||||
/**
|
||||
* better-auth owns this table and emits its own DDL per dialect, which is why the two
|
||||
* flags are typed for both: on SQLite they come back as integer 0/1 and on Postgres as a
|
||||
* real boolean. Read them through `isFlagSet` rather than comparing to 1.
|
||||
*/
|
||||
export interface UserTable {
|
||||
id: string;
|
||||
name: string;
|
||||
email: string;
|
||||
emailVerified: number;
|
||||
emailVerified: number | boolean;
|
||||
image: string | null;
|
||||
createdAt: string;
|
||||
updatedAt: string;
|
||||
role: string | null;
|
||||
banned: number | null;
|
||||
banned: number | boolean | null;
|
||||
banReason: string | null;
|
||||
banExpires: string | null;
|
||||
}
|
||||
|
||||
/**
|
||||
* True for a flag on a better-auth table, whichever dialect wrote it.
|
||||
*
|
||||
* `banned === 1` was silently false on Postgres, where the driver hands back `true`, and a
|
||||
* frozen account went on receiving notifications.
|
||||
*/
|
||||
export function isFlagSet(value: number | boolean | null | undefined): boolean {
|
||||
return value === true || value === 1;
|
||||
}
|
||||
|
||||
export interface SessionTable {
|
||||
id: string;
|
||||
expiresAt: string;
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
import { readFileSync, readdirSync } from 'node:fs';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { parseAllDocuments } from 'yaml';
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { ENV_KEYS } from './lib/env';
|
||||
|
||||
/**
|
||||
* The manifests in deploy/k8s cannot be applied to a cluster from here, so what is checked
|
||||
* is the part that actually rots: the configuration contract between them and the env
|
||||
* schema. A variable renamed in src/lib/env.ts and forgotten in a ConfigMap is a pod that
|
||||
* starts and then refuses to boot, and it is exactly the sort of thing nobody notices until
|
||||
* a deploy.
|
||||
*/
|
||||
const K8S = fileURLToPath(new URL('../../../deploy/k8s/', import.meta.url));
|
||||
|
||||
interface Manifest {
|
||||
kind?: string;
|
||||
metadata?: { name?: string };
|
||||
data?: Record<string, string>;
|
||||
stringData?: Record<string, string>;
|
||||
spec?: {
|
||||
template?: {
|
||||
spec?: {
|
||||
terminationGracePeriodSeconds?: number;
|
||||
containers?: { name?: string; env?: { name?: string; value?: string }[] }[];
|
||||
initContainers?: { name?: string }[];
|
||||
};
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
function load(): { file: string; doc: Manifest }[] {
|
||||
return readdirSync(K8S)
|
||||
.filter((file) => file.endsWith('.yaml'))
|
||||
.flatMap((file) =>
|
||||
parseAllDocuments(readFileSync(K8S + file, 'utf8'))
|
||||
.map((document) => document.toJS() as Manifest)
|
||||
.filter((doc): doc is Manifest => Boolean(doc))
|
||||
.map((doc) => ({ file, doc })),
|
||||
);
|
||||
}
|
||||
|
||||
const manifests = load();
|
||||
/** Name alone is ambiguous: a Deployment, its Service and its HPA all share one. */
|
||||
const byName = (name: string, kind = 'Deployment') =>
|
||||
manifests.find((m) => m.doc.metadata?.name === name && m.doc.kind === kind)?.doc;
|
||||
|
||||
describe('the kubernetes manifests', () => {
|
||||
it('parse, and every one names a kind', () => {
|
||||
expect(manifests.length).toBeGreaterThan(0);
|
||||
for (const { file, doc } of manifests) {
|
||||
expect(doc.kind, file).toBeTruthy();
|
||||
}
|
||||
});
|
||||
|
||||
it('set only variables the API actually reads', () => {
|
||||
const known = new Set(ENV_KEYS);
|
||||
const configured = [
|
||||
...Object.keys(byName('impuestos-config', 'ConfigMap')?.data ?? {}),
|
||||
...Object.keys(byName('impuestos-secrets', 'Secret')?.stringData ?? {}),
|
||||
...(byName('impuestos-api')?.spec?.template?.spec?.containers ?? [])
|
||||
.flatMap((container) => container.env ?? [])
|
||||
.map((entry) => entry.name ?? ''),
|
||||
...(byName('impuestos-worker')?.spec?.template?.spec?.containers ?? [])
|
||||
.flatMap((container) => container.env ?? [])
|
||||
.map((entry) => entry.name ?? ''),
|
||||
];
|
||||
|
||||
for (const key of configured) {
|
||||
expect(known.has(key), `${key} is set in deploy/k8s but no such variable exists`).toBe(true);
|
||||
}
|
||||
});
|
||||
|
||||
it('carries everything the scaled mode needs', () => {
|
||||
const config = byName('impuestos-config', 'ConfigMap')?.data ?? {};
|
||||
const secrets = byName('impuestos-secrets', 'Secret')?.stringData ?? {};
|
||||
|
||||
// SPEC.md section 15: scaled mode is Postgres and S3, not a choice.
|
||||
expect(secrets['DATABASE_URL']).toMatch(/^postgres/);
|
||||
expect(config['STORAGE_DRIVER']).toBe('s3');
|
||||
for (const key of ['S3_BUCKET', 'S3_REGION']) expect(config[key]).toBeTruthy();
|
||||
for (const key of ['BETTER_AUTH_SECRET', 'S3_ACCESS_KEY_ID', 'S3_SECRET_ACCESS_KEY']) {
|
||||
expect(secrets).toHaveProperty(key);
|
||||
}
|
||||
});
|
||||
|
||||
it('runs exactly one poller per job, not one per replica', () => {
|
||||
const api = byName('impuestos-api')?.spec?.template?.spec?.containers?.[0]?.env ?? [];
|
||||
const worker = byName('impuestos-worker')?.spec?.template?.spec?.containers?.[0]?.env ?? [];
|
||||
|
||||
// An API replica that polled would multiply pollers by the replica count.
|
||||
expect(api.find((entry) => entry.name === 'ROLE')?.value).toBe('server');
|
||||
expect(api.find((entry) => entry.name === 'JOBS_INLINE')?.value).toBe('false');
|
||||
expect(worker.find((entry) => entry.name === 'ROLE')?.value).toBe('worker');
|
||||
});
|
||||
|
||||
it('migrates before it serves', () => {
|
||||
const init = byName('impuestos-api')?.spec?.template?.spec?.initContainers ?? [];
|
||||
expect(init.map((container) => container.name)).toContain('migrate');
|
||||
});
|
||||
|
||||
it('gives every pod longer to stop than the API spends draining', () => {
|
||||
// src/index.ts drains for up to 25 seconds. A shorter grace period would have the
|
||||
// cluster kill the process in the middle of the requests it is trying to finish.
|
||||
const drainSeconds = 25;
|
||||
for (const name of ['impuestos-api', 'impuestos-worker', 'impuestos-web']) {
|
||||
const grace = byName(name)?.spec?.template?.spec?.terminationGracePeriodSeconds;
|
||||
expect(grace, name).toBeGreaterThan(drainSeconds);
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -1,5 +1,6 @@
|
||||
import { DataExportDto, DependentDto, ProfileDto } from '@impuestos/contracts';
|
||||
import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest';
|
||||
import { isFlagSet } from '../../db/schema';
|
||||
import { createHarness, type Harness } from '../../test/harness';
|
||||
|
||||
let h: Harness;
|
||||
@@ -228,7 +229,8 @@ describe('DELETE /me/account', () => {
|
||||
.select(['id', 'banned', 'banReason'])
|
||||
.where('email', '=', 'carlos@demo.local')
|
||||
.executeTakeFirstOrThrow();
|
||||
expect(user.banned).toBe(1);
|
||||
// Integer 0/1 on SQLite, a real boolean on Postgres: the flag is read, not compared.
|
||||
expect(isFlagSet(user.banned)).toBe(true);
|
||||
expect(user.banReason).toBe('account_deleted');
|
||||
|
||||
const sessions = await h.deps.handle.db
|
||||
|
||||
@@ -14,7 +14,7 @@ import { createStorage } from './modules/storage';
|
||||
const DRAIN_TIMEOUT_MS = 25_000;
|
||||
|
||||
const env = loadEnv();
|
||||
const handle = createDb(env.DATABASE_URL);
|
||||
const handle = createDb(env.DATABASE_URL, { poolMax: env.DATABASE_POOL_MAX });
|
||||
const auth = createAuth({
|
||||
db: handle.db,
|
||||
dialect: handle.dialect,
|
||||
|
||||
@@ -42,6 +42,12 @@ const EnvObject = z.object({
|
||||
ANTHROPIC_API_KEY: optionalString,
|
||||
OCR_MODEL: z.string().trim().default('claude-sonnet-4-6'),
|
||||
|
||||
/**
|
||||
* Postgres connections per process. Replicas multiply it, so it has to be sized
|
||||
* against the server's max_connections rather than left to a library default.
|
||||
*/
|
||||
DATABASE_POOL_MAX: z.coerce.number().int().min(1).max(100).default(10),
|
||||
|
||||
PUSH_VAPID_PUBLIC_KEY: optionalString,
|
||||
PUSH_VAPID_PRIVATE_KEY: optionalString,
|
||||
/** Contact for the push service. Must be an https: or mailto: URL, per RFC 8292. */
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
import pg from 'pg';
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { createDb } from '../../db/index';
|
||||
import { migrateToLatest } from '../../db/migrator';
|
||||
import { parseEnv } from '../../lib/env';
|
||||
import { TEST_ENV } from '../../test/harness';
|
||||
import { claimJob } from './claim';
|
||||
import { enqueue } from './queue';
|
||||
|
||||
/**
|
||||
* SPEC.md section 13: with Postgres and two workers, enqueue 100 jobs and assert each one
|
||||
* is claimed exactly once.
|
||||
*
|
||||
* This is the whole basis of running more than one replica, so it runs against a real
|
||||
* Postgres or not at all. `FOR UPDATE SKIP LOCKED` cannot be simulated: SQLite's single
|
||||
* writer makes any test of it pass for the wrong reason.
|
||||
*/
|
||||
const POSTGRES_URL = process.env['TEST_DATABASE_URL'];
|
||||
const JOBS = 100;
|
||||
const WORKERS = 2;
|
||||
|
||||
describe.skipIf(!POSTGRES_URL)('two workers on one Postgres queue', () => {
|
||||
it('claims each of a hundred jobs exactly once', async () => {
|
||||
const database = `scale_${Date.now()}`;
|
||||
await withAdmin(POSTGRES_URL as string, (db) => db.query(`create database "${database}"`));
|
||||
|
||||
const url = new URL(POSTGRES_URL as string);
|
||||
url.pathname = `/${database}`;
|
||||
const parsed = parseEnv({ ...TEST_ENV, DATABASE_URL: url.toString(), DATABASE_POOL_MAX: '5' });
|
||||
const env = parsed.env;
|
||||
if (!parsed.ok || !env) throw new Error(parsed.message);
|
||||
|
||||
// Two handles, because two worker processes would have two pools.
|
||||
const workers = Array.from({ length: WORKERS }, () =>
|
||||
createDb(env.DATABASE_URL, { poolMax: 5 }),
|
||||
);
|
||||
|
||||
try {
|
||||
await migrateToLatest(workers[0] as (typeof workers)[number], env);
|
||||
|
||||
for (let index = 0; index < JOBS; index += 1) {
|
||||
await enqueue(workers[0]!.db, { type: 'classify_document', payload: { index } });
|
||||
}
|
||||
|
||||
const now = new Date().toISOString();
|
||||
const claimed: string[] = [];
|
||||
|
||||
/** One worker draining until the queue gives it nothing. */
|
||||
async function drain(handle: (typeof workers)[number], instanceId: string): Promise<void> {
|
||||
for (;;) {
|
||||
const job = await claimJob(handle.db, handle.dialect, { instanceId, now });
|
||||
if (!job) return;
|
||||
claimed.push(job.id);
|
||||
}
|
||||
}
|
||||
|
||||
// Both at once, which is the only arrangement that can double claim.
|
||||
await Promise.all(workers.map((handle, index) => drain(handle, `worker-${index}`)));
|
||||
|
||||
expect(claimed).toHaveLength(JOBS);
|
||||
expect(new Set(claimed).size).toBe(JOBS);
|
||||
|
||||
// And the queue agrees: everything is running under someone, nothing is pending.
|
||||
const rows = await workers[0]!.db.selectFrom('jobs').select(['status', 'locked_by']).execute();
|
||||
expect(rows).toHaveLength(JOBS);
|
||||
expect(rows.every((row) => row.status === 'running')).toBe(true);
|
||||
expect(new Set(rows.map((row) => row.locked_by)).size).toBe(WORKERS);
|
||||
} finally {
|
||||
await Promise.all(workers.map((handle) => handle.close()));
|
||||
await withAdmin(POSTGRES_URL as string, (db) =>
|
||||
db.query(`drop database if exists "${database}" with (force)`),
|
||||
);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
async function withAdmin(
|
||||
databaseUrl: string,
|
||||
run: (client: pg.Client) => Promise<unknown>,
|
||||
): Promise<void> {
|
||||
const client = new pg.Client({ connectionString: databaseUrl });
|
||||
await client.connect();
|
||||
try {
|
||||
await run(client);
|
||||
} finally {
|
||||
await client.end();
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,8 @@
|
||||
import { computeRucDv, dueDateFor, formatIsoDate } from '@impuestos/rules';
|
||||
import { afterEach, describe, expect, it } from 'vitest';
|
||||
import { isFlagSet } from '../../db/schema';
|
||||
import { createHarness, type Harness } from '../../test/harness';
|
||||
import type { Channels, OutgoingMessage } from '../notifications';
|
||||
import { sendNotification, type Channels, type OutgoingMessage } from '../notifications';
|
||||
import { runAutoConfirmSweep, runDeadlineSweep, runDigestSweep } from './sweep-handlers';
|
||||
import { scheduleSweeps, sweepDedupeKey } from './sweeps';
|
||||
|
||||
@@ -140,6 +141,48 @@ describe('deadline sweep, ten days out', () => {
|
||||
});
|
||||
});
|
||||
|
||||
/**
|
||||
* The flag better-auth writes is integer 0/1 on SQLite and a boolean on Postgres. Reading
|
||||
* it with `=== 1` was silently false on Postgres, so a frozen account went on being
|
||||
* notified. This runs on whichever dialect the suite is pointed at.
|
||||
*/
|
||||
describe('a frozen account', () => {
|
||||
it('is not a recipient on either dialect', async () => {
|
||||
const h = (harness = await createHarness({ channels: recordingChannels() }));
|
||||
const db = h.deps.handle.db;
|
||||
|
||||
const carlos = await db
|
||||
.selectFrom('user')
|
||||
.select('id')
|
||||
.where('email', '=', 'carlos@demo.local')
|
||||
.executeTakeFirstOrThrow();
|
||||
|
||||
await db
|
||||
.updateTable('notification_prefs')
|
||||
.set({ email_enabled: 1 })
|
||||
.where('user_id', '=', carlos.id)
|
||||
.execute();
|
||||
await db.updateTable('user').set({ banned: 1 }).where('id', '=', carlos.id).execute();
|
||||
|
||||
const stored = await db
|
||||
.selectFrom('user')
|
||||
.select('banned')
|
||||
.where('id', '=', carlos.id)
|
||||
.executeTakeFirstOrThrow();
|
||||
// Whatever the driver hands back, the helper agrees it is set.
|
||||
expect(isFlagSet(stored.banned)).toBe(true);
|
||||
|
||||
const result = await sendNotification(
|
||||
{ db, channels: h.channels },
|
||||
{
|
||||
userId: carlos.id,
|
||||
notification: { kind: 'declaration_ready', form: '120', period: '2026-08', declarationId: 'x' },
|
||||
},
|
||||
);
|
||||
expect(result.delivered).toEqual([]);
|
||||
});
|
||||
});
|
||||
|
||||
describe('the queued notification actually goes out', () => {
|
||||
it('reaches the channels the user has enabled, in their language', async () => {
|
||||
const channels = recordingChannels();
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { DEFAULT_LOCALE, isLocale, type Locale } from '@impuestos/i18n';
|
||||
import type { Kysely } from 'kysely';
|
||||
import { z } from 'zod';
|
||||
import type { Database } from '../../db/schema';
|
||||
import { isFlagSet, type Database } from '../../db/schema';
|
||||
import type { Channels } from './channels';
|
||||
import { emailSubject, renderNotification, type NotificationKind } from './templates';
|
||||
|
||||
@@ -38,7 +38,7 @@ export async function sendNotification(
|
||||
.executeTakeFirst();
|
||||
|
||||
// A frozen account is not a recipient.
|
||||
if (!recipient || recipient.banned === 1) return { delivered: [] };
|
||||
if (!recipient || isFlagSet(recipient.banned)) return { delivered: [] };
|
||||
|
||||
const locale: Locale = isLocale(recipient.locale) ? recipient.locale : DEFAULT_LOCALE;
|
||||
const message = renderNotification(locale, args.notification);
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
import { mkdtempSync, readdirSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { parseEnv } from '../../lib/env';
|
||||
import { createStorage, storageKey, type StorageDriver } from './index';
|
||||
|
||||
/**
|
||||
* The storage driver is the one thing scaled mode swaps out (SPEC.md section 15: S3 is
|
||||
* required there), so both implementations are held to the same round trip.
|
||||
*
|
||||
* The S3 half needs a real endpoint and is skipped without one. CI points it at MinIO;
|
||||
* locally, set TEST_S3_ENDPOINT and the four S3 variables to run it.
|
||||
*/
|
||||
const S3_ENDPOINT = process.env['TEST_S3_ENDPOINT'];
|
||||
|
||||
const BASE = {
|
||||
NODE_ENV: 'test',
|
||||
DATABASE_URL: 'sqlite::memory:',
|
||||
BETTER_AUTH_SECRET: 'storage-test-secret-long-enough-32ch',
|
||||
BETTER_AUTH_URL: 'http://localhost:3000',
|
||||
APP_PUBLIC_URL: 'http://localhost:3000',
|
||||
};
|
||||
|
||||
function driverFor(overrides: Record<string, string>): StorageDriver {
|
||||
const parsed = parseEnv({ ...BASE, ...overrides });
|
||||
if (!parsed.ok || !parsed.env) throw new Error(parsed.message);
|
||||
return createStorage(parsed.env);
|
||||
}
|
||||
|
||||
async function readAll(stream: ReadableStream<Uint8Array>): Promise<Uint8Array> {
|
||||
const chunks: Uint8Array[] = [];
|
||||
const reader = stream.getReader();
|
||||
for (;;) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
if (value) chunks.push(value);
|
||||
}
|
||||
return new Uint8Array(chunks.flatMap((chunk) => [...chunk]));
|
||||
}
|
||||
|
||||
describe('storage keys', () => {
|
||||
it('never carry a name the uploader chose', () => {
|
||||
// A key is the owner and an id we generated, so "../../etc/passwd" as a filename has
|
||||
// nowhere to go: it is not part of the key at all.
|
||||
expect(storageKey('user-1', 'abc')).toBe('user-1/abc');
|
||||
expect(storageKey('user-1', 'abc')).not.toContain('..');
|
||||
});
|
||||
|
||||
it('keep one account out of the prefix belonging to another', () => {
|
||||
expect(storageKey('user-1', 'abc').startsWith('user-2/')).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe('the driver switch', () => {
|
||||
it('follows STORAGE_DRIVER and nothing else', () => {
|
||||
expect(driverFor({ STORAGE_LOCAL_PATH: mkdtempSync(join(tmpdir(), 'st-')) }).kind).toBe('local');
|
||||
expect(
|
||||
driverFor({
|
||||
STORAGE_DRIVER: 's3',
|
||||
S3_BUCKET: 'whatever',
|
||||
S3_REGION: 'us-east-1',
|
||||
S3_ACCESS_KEY_ID: 'a',
|
||||
S3_SECRET_ACCESS_KEY: 'b',
|
||||
}).kind,
|
||||
).toBe('s3');
|
||||
});
|
||||
|
||||
it('refuses the s3 driver with no bucket rather than failing on first upload', () => {
|
||||
expect(() => driverFor({ STORAGE_DRIVER: 's3' })).toThrow(/S3_BUCKET/);
|
||||
});
|
||||
});
|
||||
|
||||
describe('the local driver', () => {
|
||||
const path = mkdtempSync(join(tmpdir(), 'impuestos-storage-'));
|
||||
const storage = driverFor({ STORAGE_LOCAL_PATH: path });
|
||||
|
||||
it('stores, reads back and deletes', async () => {
|
||||
const key = storageKey('user-1', 'factura');
|
||||
const body = new TextEncoder().encode('una factura');
|
||||
|
||||
await storage.put(key, body, 'text/plain');
|
||||
expect(new TextDecoder().decode(await readAll(await storage.getStream(key)))).toBe('una factura');
|
||||
|
||||
// Local files are served through the authenticated route, never linked directly, and
|
||||
// the key is escaped into the path rather than pasted into it.
|
||||
expect(await storage.url(key)).toBe(`/api/files/${encodeURIComponent(key)}`);
|
||||
|
||||
await storage.delete(key);
|
||||
await expect(storage.getStream(key)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('answers the readiness check once its directory exists', async () => {
|
||||
await expect(storage.check()).resolves.toBeUndefined();
|
||||
// Writing created the per-user prefix; nothing else is in there.
|
||||
expect(readdirSync(path).length).toBeGreaterThanOrEqual(0);
|
||||
});
|
||||
});
|
||||
|
||||
/** Built inside the tests: a skipped describe still runs its own body. */
|
||||
function s3Driver(overrides: Record<string, string> = {}): StorageDriver {
|
||||
return driverFor({
|
||||
STORAGE_DRIVER: 's3',
|
||||
S3_ENDPOINT: S3_ENDPOINT ?? '',
|
||||
S3_REGION: process.env['TEST_S3_REGION'] ?? 'us-east-1',
|
||||
S3_BUCKET: process.env['TEST_S3_BUCKET'] ?? 'comprobantes',
|
||||
S3_ACCESS_KEY_ID: process.env['TEST_S3_ACCESS_KEY_ID'] ?? '',
|
||||
S3_SECRET_ACCESS_KEY: process.env['TEST_S3_SECRET_ACCESS_KEY'] ?? '',
|
||||
S3_FORCE_PATH_STYLE: 'true',
|
||||
...overrides,
|
||||
});
|
||||
}
|
||||
|
||||
describe.skipIf(!S3_ENDPOINT)('the s3 driver', () => {
|
||||
it('stores, reads back, signs and deletes', async () => {
|
||||
const storage = s3Driver();
|
||||
const key = storageKey('user-1', `factura-${Date.now()}`);
|
||||
const body = new TextEncoder().encode('una factura');
|
||||
|
||||
await storage.put(key, body, 'text/plain');
|
||||
expect(new TextDecoder().decode(await readAll(await storage.getStream(key)))).toBe('una factura');
|
||||
|
||||
// A signed URL, which is how a browser gets the file in scaled mode.
|
||||
const url = await storage.url(key);
|
||||
expect(url).toContain(key);
|
||||
expect(url).toMatch(/X-Amz-Signature/);
|
||||
expect(new TextDecoder().decode(new Uint8Array(await (await fetch(url)).arrayBuffer()))).toBe(
|
||||
'una factura',
|
||||
);
|
||||
|
||||
await storage.delete(key);
|
||||
await expect(storage.getStream(key)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('fails the readiness check when the bucket is not there', async () => {
|
||||
await expect(s3Driver({ S3_BUCKET: 'no-such-bucket-here' }).check()).rejects.toThrow();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,148 @@
|
||||
import { spawn, type ChildProcess } from 'node:child_process';
|
||||
import { request as httpRequest } from 'node:http';
|
||||
import { mkdtempSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { afterAll, describe, expect, it } from 'vitest';
|
||||
import { createAuth } from './auth/options';
|
||||
import { createDb } from './db/index';
|
||||
import { migrateToLatest } from './db/migrator';
|
||||
import { seed } from './db/seed';
|
||||
import { parseEnv } from './lib/env';
|
||||
|
||||
/**
|
||||
* SPEC.md section 15: on SIGTERM the API stops accepting, finishes what it is already
|
||||
* handling, and exits 0. A rolling deploy that cuts a request in half loses a scan someone
|
||||
* just took, so this runs against a real process and a real socket rather than the app
|
||||
* object.
|
||||
*/
|
||||
const ENTRY = fileURLToPath(new URL('./index.ts', import.meta.url));
|
||||
const PORT = 4571;
|
||||
const DIR = mkdtempSync(join(tmpdir(), 'impuestos-shutdown-'));
|
||||
|
||||
let child: ChildProcess | null = null;
|
||||
|
||||
afterAll(() => {
|
||||
child?.kill('SIGKILL');
|
||||
});
|
||||
|
||||
describe('graceful shutdown', () => {
|
||||
it('finishes an in-flight request, then exits 0', async () => {
|
||||
const env = {
|
||||
NODE_ENV: 'test',
|
||||
PORT: String(PORT),
|
||||
DATABASE_URL: `sqlite:${join(DIR, 'app.db')}`,
|
||||
STORAGE_LOCAL_PATH: join(DIR, 'files'),
|
||||
BETTER_AUTH_SECRET: 'shutdown-test-secret-long-enough-32ch',
|
||||
BETTER_AUTH_URL: `http://127.0.0.1:${PORT}`,
|
||||
APP_PUBLIC_URL: `http://127.0.0.1:${PORT}`,
|
||||
};
|
||||
|
||||
await prepareDatabase(env);
|
||||
child = await start(env);
|
||||
|
||||
const cookie = await signIn();
|
||||
|
||||
/**
|
||||
* A request whose body arrives in two pieces. The server has the headers and is inside
|
||||
* the handler waiting for the rest, which is exactly the state a deploy must not cut.
|
||||
*/
|
||||
const body = JSON.stringify({ pushEnabled: false });
|
||||
const [head, tail] = [body.slice(0, 4), body.slice(4)];
|
||||
|
||||
const response = new Promise<{ status: number; text: string }>((resolve, reject) => {
|
||||
const outgoing = httpRequest(
|
||||
{
|
||||
host: '127.0.0.1',
|
||||
port: PORT,
|
||||
path: '/api/me/notification-prefs',
|
||||
method: 'PATCH',
|
||||
headers: {
|
||||
cookie,
|
||||
'content-type': 'application/json',
|
||||
'content-length': Buffer.byteLength(body),
|
||||
},
|
||||
},
|
||||
(incoming) => {
|
||||
let text = '';
|
||||
incoming.on('data', (chunk) => (text += chunk));
|
||||
incoming.on('end', () => resolve({ status: incoming.statusCode ?? 0, text }));
|
||||
},
|
||||
);
|
||||
outgoing.on('error', reject);
|
||||
outgoing.write(head);
|
||||
// Held open on purpose.
|
||||
setTimeout(() => outgoing.end(tail), 1200);
|
||||
});
|
||||
|
||||
// Long enough for the server to be inside the handler, well short of the body.
|
||||
await new Promise((resolve) => setTimeout(resolve, 400));
|
||||
const exited = new Promise<number>((resolve) => {
|
||||
child?.on('exit', (code) => resolve(code ?? -1));
|
||||
});
|
||||
child?.kill('SIGTERM');
|
||||
|
||||
// Draining, so a load balancer stops sending here before the door closes.
|
||||
await new Promise((resolve) => setTimeout(resolve, 100));
|
||||
const health = await fetch(`http://127.0.0.1:${PORT}/healthz`).catch(() => null);
|
||||
if (health) expect(health.status).toBe(503);
|
||||
|
||||
const result = await response;
|
||||
expect(result.status).toBe(200);
|
||||
|
||||
expect(await exited).toBe(0);
|
||||
// Spawning a process, migrating a file database and holding a request open is slower
|
||||
// than a unit test has any business being.
|
||||
}, 60_000);
|
||||
});
|
||||
|
||||
async function prepareDatabase(env: Record<string, string>): Promise<void> {
|
||||
const parsed = parseEnv(env);
|
||||
if (!parsed.ok || !parsed.env) throw new Error(parsed.message);
|
||||
const handle = createDb(parsed.env.DATABASE_URL);
|
||||
try {
|
||||
await migrateToLatest(handle, parsed.env);
|
||||
const auth = createAuth({
|
||||
db: handle.db,
|
||||
dialect: handle.dialect,
|
||||
env: parsed.env,
|
||||
sendOtp: async () => undefined,
|
||||
});
|
||||
await seed(handle, auth);
|
||||
} finally {
|
||||
await handle.close();
|
||||
}
|
||||
}
|
||||
|
||||
async function start(env: Record<string, string>): Promise<ChildProcess> {
|
||||
const proc = spawn(process.execPath, ['--import', 'tsx', ENTRY], {
|
||||
// From apps/api, where tsx is a dependency and where `pnpm dev` runs it from.
|
||||
cwd: fileURLToPath(new URL('..', import.meta.url)),
|
||||
env: { ...process.env, ...env },
|
||||
stdio: ['ignore', 'pipe', 'pipe'],
|
||||
});
|
||||
proc.stderr?.on('data', (chunk) => console.error(`[api] ${String(chunk).trim()}`));
|
||||
|
||||
const deadline = Date.now() + 30_000;
|
||||
for (;;) {
|
||||
if (Date.now() > deadline) throw new Error('the api did not come up');
|
||||
const ok = await fetch(`http://127.0.0.1:${PORT}/healthz`)
|
||||
.then((response) => response.ok)
|
||||
.catch(() => false);
|
||||
if (ok) return proc;
|
||||
await new Promise((resolve) => setTimeout(resolve, 200));
|
||||
}
|
||||
}
|
||||
|
||||
async function signIn(): Promise<string> {
|
||||
const response = await fetch(`http://127.0.0.1:${PORT}/api/auth/sign-in/email`, {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: JSON.stringify({ email: 'maria@demo.local', password: 'demo-maria-1' }),
|
||||
});
|
||||
if (!response.ok) throw new Error(`sign in failed: ${response.status}`);
|
||||
const cookie = response.headers.get('set-cookie');
|
||||
if (!cookie) throw new Error('sign in returned no cookie');
|
||||
return cookie.split(';')[0] ?? '';
|
||||
}
|
||||
@@ -1,3 +1,5 @@
|
||||
import pg from 'pg';
|
||||
import { uuidv7 } from 'uuidv7';
|
||||
import { mkdtempSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
@@ -13,9 +15,16 @@ import type { OcrProvider } from '../modules/documents/ocr';
|
||||
import { createChannels, type Channels } from '../modules/notifications';
|
||||
import { createStorage, type StorageDriver } from '../modules/storage';
|
||||
|
||||
/**
|
||||
* The suite runs against SQLite in memory by default and against a real Postgres when
|
||||
* `TEST_DATABASE_URL` is set, which is how the CI matrix covers both dialects with one set
|
||||
* of tests (SPEC.md section 13).
|
||||
*/
|
||||
const POSTGRES_URL = process.env['TEST_DATABASE_URL'];
|
||||
|
||||
export const TEST_ENV: Record<string, string> = {
|
||||
NODE_ENV: 'test',
|
||||
DATABASE_URL: 'sqlite::memory:',
|
||||
DATABASE_URL: POSTGRES_URL ?? 'sqlite::memory:',
|
||||
BETTER_AUTH_SECRET: 'test-secret-that-is-long-enough-32chars',
|
||||
BETTER_AUTH_URL: 'http://localhost:3000',
|
||||
APP_PUBLIC_URL: 'http://localhost:3000',
|
||||
@@ -47,11 +56,35 @@ export interface Harness extends AppHandle {
|
||||
}
|
||||
|
||||
export async function createHarness(options: HarnessOptions = {}): Promise<Harness> {
|
||||
const parsed = parseEnv({ ...TEST_ENV, ...(options.env ?? {}) });
|
||||
/**
|
||||
* On Postgres every harness gets a database of its own. `sqlite::memory:` gives each test
|
||||
* a private database for free; a shared Postgres does not, and tests that seed the same
|
||||
* four accounts into one set of tables would tread on each other.
|
||||
*
|
||||
* A database rather than a schema: Kysely's migrator asks the introspector whether its
|
||||
* bookkeeping tables exist and, with no schema configured, a table of that name in any
|
||||
* schema counts. Twenty parallel schemas each holding a `kysely_migration` therefore
|
||||
* convince each other that the work is already done.
|
||||
*/
|
||||
const scratchDatabase = POSTGRES_URL ? `t_${uuidv7().replaceAll('-', '')}` : null;
|
||||
if (scratchDatabase) await createDatabase(POSTGRES_URL as string, scratchDatabase);
|
||||
|
||||
const parsed = parseEnv({
|
||||
...TEST_ENV,
|
||||
...(scratchDatabase
|
||||
? {
|
||||
DATABASE_URL: withDatabase(POSTGRES_URL as string, scratchDatabase),
|
||||
// Small on purpose: one pool per harness, many harnesses at once, and one
|
||||
// server's max_connections between them.
|
||||
DATABASE_POOL_MAX: '2',
|
||||
}
|
||||
: {}),
|
||||
...(options.env ?? {}),
|
||||
});
|
||||
if (!parsed.ok || !parsed.env) throw new Error(parsed.message);
|
||||
const env = parsed.env;
|
||||
|
||||
const handle = createDb(env.DATABASE_URL);
|
||||
const handle = createDb(env.DATABASE_URL, { poolMax: env.DATABASE_POOL_MAX });
|
||||
await migrateToLatest(handle, env);
|
||||
|
||||
const auth = createAuth({ db: handle.db, dialect: handle.dialect, env, sendOtp: async () => undefined });
|
||||
@@ -92,6 +125,8 @@ export async function createHarness(options: HarnessOptions = {}): Promise<Harne
|
||||
close: async () => {
|
||||
await poller.stop();
|
||||
await handle.close();
|
||||
// After the pool is closed: Postgres refuses to drop a database anything is on.
|
||||
if (scratchDatabase) await dropDatabase(POSTGRES_URL as string, scratchDatabase);
|
||||
},
|
||||
signIn: async (email, password) => {
|
||||
const response = await appHandle.app.request('/api/auth/sign-in/email', {
|
||||
@@ -106,3 +141,33 @@ export async function createHarness(options: HarnessOptions = {}): Promise<Harne
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function withDatabase(databaseUrl: string, name: string): string {
|
||||
const url = new URL(databaseUrl);
|
||||
url.pathname = `/${name}`;
|
||||
return url.toString();
|
||||
}
|
||||
|
||||
/** `create database` cannot run inside a transaction, so it gets its own connection. */
|
||||
async function createDatabase(databaseUrl: string, name: string): Promise<void> {
|
||||
await onAdminConnection(databaseUrl, (client) => client.query(`create database "${name}"`));
|
||||
}
|
||||
|
||||
async function dropDatabase(databaseUrl: string, name: string): Promise<void> {
|
||||
await onAdminConnection(databaseUrl, (client) =>
|
||||
client.query(`drop database if exists "${name}" with (force)`),
|
||||
);
|
||||
}
|
||||
|
||||
async function onAdminConnection(
|
||||
databaseUrl: string,
|
||||
run: (client: pg.Client) => Promise<unknown>,
|
||||
): Promise<void> {
|
||||
const client = new pg.Client({ connectionString: databaseUrl });
|
||||
await client.connect();
|
||||
try {
|
||||
await run(client);
|
||||
} finally {
|
||||
await client.end();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user