diff --git a/packages/orchestrator/lib/backpressure-monitor.integration.test.ts b/packages/orchestrator/lib/backpressure-monitor.integration.test.ts index 4d814d7ba5..13e1683e4c 100644 --- a/packages/orchestrator/lib/backpressure-monitor.integration.test.ts +++ b/packages/orchestrator/lib/backpressure-monitor.integration.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { getTestDbClient, Scheduler } from '@nangohq/scheduler'; import { metrics, nanoid } from '@nangohq/utils'; @@ -17,7 +17,7 @@ const noopCallbacks = { }; describe('BackpressureMonitor', () => { - const dbClient = getTestDbClient(); + const dbClient = getTestDbClient('orchestrator_backpressure'); let scheduler: Scheduler; beforeEach(async () => { @@ -30,6 +30,10 @@ describe('BackpressureMonitor', () => { await dbClient.clearDatabase(); }); + afterAll(async () => { + await dbClient.destroy(); + }); + it('should emit ORCH_QUEUE_BACKPRESSURE for groups exceeding their max concurrency', async () => { const groupKey = `sync:env:${nanoid()}`; // Cap of 2, three queued -> backpressure diff --git a/packages/orchestrator/lib/clients/client.integration.test.ts b/packages/orchestrator/lib/clients/client.integration.test.ts index 432b4457d7..2397990055 100644 --- a/packages/orchestrator/lib/clients/client.integration.test.ts +++ b/packages/orchestrator/lib/clients/client.integration.test.ts @@ -12,7 +12,7 @@ import type { PostImmediate } from '../routes/v1/postImmediate.js'; import type { Task } from '@nangohq/scheduler'; import type { Result } from '@nangohq/utils'; -const dbClient = getTestDbClient(); +const dbClient = getTestDbClient('orchestrator_client'); const eventsHandler = new TaskEventsHandler(dbClient.db); const scheduler = new Scheduler({ db: dbClient.db, @@ -33,6 +33,7 @@ describe('OrchestratorClient', async () => { afterAll(async () => { scheduler.stop(); await dbClient.clearDatabase(); + await dbClient.destroy(); }); describe('recurring schedule', () => { diff --git a/packages/orchestrator/lib/clients/processor.integration.test.ts b/packages/orchestrator/lib/clients/processor.integration.test.ts index d6ea6126e7..8a9a0f420b 100644 --- a/packages/orchestrator/lib/clients/processor.integration.test.ts +++ b/packages/orchestrator/lib/clients/processor.integration.test.ts @@ -15,7 +15,7 @@ import type { OrchestratorTask } from './types.js'; import type { Task } from '@nangohq/scheduler'; import type { Result } from '@nangohq/utils'; -const dbClient = getTestDbClient(); +const dbClient = getTestDbClient('orchestrator_processor'); const taskEventsHandler = new TaskEventsHandler(dbClient.db); const scheduler = new Scheduler({ db: dbClient.db, @@ -37,6 +37,7 @@ describe('OrchestratorProcessor', () => { scheduler.stop(); await setTimeout(100); // wait for the scheduler to stop await dbClient.clearDatabase(); + await dbClient.destroy(); }); it('should process tasks', async () => { diff --git a/packages/orchestrator/lib/events.integration.test.ts b/packages/orchestrator/lib/events.integration.test.ts index 7e42af919f..fa28a8496d 100644 --- a/packages/orchestrator/lib/events.integration.test.ts +++ b/packages/orchestrator/lib/events.integration.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { getTestDbClient } from '@nangohq/scheduler'; @@ -7,7 +7,7 @@ import { taskEvents, TaskEventsHandler } from './events.js'; import type { Task } from '@nangohq/scheduler'; -const dbClient = getTestDbClient(); +const dbClient = getTestDbClient('orchestrator_events'); const mockActionTask: Task = { id: '00000000-0000-0000-0000-000000000000', name: 'task-name', @@ -57,6 +57,10 @@ describe('TaskEventsHandler', () => { await dbClient.clearDatabase(); }); + afterAll(async () => { + await dbClient.destroy(); + }); + describe('emit', () => { it('should throw error for empty event name', () => { expect(() => eventsHandler.emit('')).toThrow('Event name must be a non-empty string'); diff --git a/packages/orchestrator/lib/helpers.test.ts b/packages/orchestrator/lib/helpers.test.ts index 7380de5784..8ce92bbd94 100644 --- a/packages/orchestrator/lib/helpers.test.ts +++ b/packages/orchestrator/lib/helpers.test.ts @@ -14,8 +14,8 @@ export class TestOrchestratorService { private scheduler: Scheduler | null; private eventsHandler: TaskEventsHandler; - constructor({ port }: { port: number }) { - this.dbClient = getTestDbClient(); + constructor({ port, schema }: { port: number; schema: string }) { + this.dbClient = getTestDbClient(schema); this.eventsHandler = new TaskEventsHandler(this.dbClient.db); this.port = port; this.scheduler = null; diff --git a/packages/scheduler/lib/daemons/daemon.integration.test.ts b/packages/scheduler/lib/daemons/daemon.integration.test.ts index f35809fbda..3a03f2b0d9 100644 --- a/packages/scheduler/lib/daemons/daemon.integration.test.ts +++ b/packages/scheduler/lib/daemons/daemon.integration.test.ts @@ -1,15 +1,20 @@ import { setTimeout } from 'node:timers/promises'; -import { describe, expect, it, vi } from 'vitest'; +import { afterAll, describe, expect, it, vi } from 'vitest'; import { getTestDbClient } from '../db/helpers.test.js'; import { SchedulerDaemon } from './daemon.js'; import type knex from 'knex'; -const db = getTestDbClient().db; +const dbClient = getTestDbClient('scheduler_daemon'); +const db = dbClient.db; describe('SchedulerDaemon', () => { + afterAll(async () => { + await dbClient.destroy(); + }); + it('should be abortable', async () => { class TestDaemon extends SchedulerDaemon { constructor({ db, abortSignal, onError }: { db: knex.Knex; abortSignal: AbortSignal; onError: (err: Error) => void }) { diff --git a/packages/scheduler/lib/daemons/scheduling/scheduling.integration.test.ts b/packages/scheduler/lib/daemons/scheduling/scheduling.integration.test.ts index e80f71864f..85498643f7 100644 --- a/packages/scheduler/lib/daemons/scheduling/scheduling.integration.test.ts +++ b/packages/scheduler/lib/daemons/scheduling/scheduling.integration.test.ts @@ -1,8 +1,7 @@ import { uuidv7 } from 'uuidv7'; -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterAll, afterEach, beforeEach, describe, expect, it } from 'vitest'; import { defaultSchedulerConfig } from '../../config.js'; -import { DatabaseClient, defaultDatabaseClientOptions } from '../../db/client.js'; import { getTestDbClient } from '../../db/helpers.test.js'; import { DbSchedule, SCHEDULES_TABLE } from '../../models/schedules.js'; import * as schedules from '../../models/schedules.js'; @@ -16,7 +15,7 @@ import type { Schedule, ScheduleState, Task, TaskState } from '../../types.js'; import type knex from 'knex'; describe('dueSchedules', () => { - const dbClient = getTestDbClient(); + const dbClient = getTestDbClient('scheduler_scheduling'); const db = dbClient.db; beforeEach(async () => { @@ -27,6 +26,10 @@ describe('dueSchedules', () => { await dbClient.clearDatabase(); }); + afterAll(async () => { + await dbClient.destroy(); + }); + it('should not return schedule that is deleted', async () => { await addSchedule(db, { state: 'DELETED', frequency: '3 minutes' }); const due = await dueSchedules(db); @@ -93,13 +96,7 @@ describe('dueSchedules', () => { }); describe('SchedulingDaemon', () => { - // Dedicated schema: running the daemon against the shared 'scheduler' schema races with the - // looping daemons in scheduler.integration.test.ts via SKIP LOCKED. - const dbClient = new DatabaseClient({ - ...defaultDatabaseClientOptions, - url: `postgres://${process.env['NANGO_DB_USER']}:${process.env['NANGO_DB_PASSWORD']}@${process.env['NANGO_DB_HOST']}:${process.env['NANGO_DB_PORT']}/${process.env['NANGO_DB_NAME']}`, - schema: 'scheduler_daemon' - }); + const dbClient = getTestDbClient('scheduler_scheduling_daemon'); const db = dbClient.db; beforeEach(async () => { @@ -110,6 +107,10 @@ describe('SchedulingDaemon', () => { await dbClient.clearDatabase(); }); + afterAll(async () => { + await dbClient.destroy(); + }); + it('should stamp materialized tasks with the configured recurringGroupMaxConcurrency', async () => { const schedule = await addSchedule(db); const daemon = new SchedulingDaemon({ diff --git a/packages/scheduler/lib/db/helpers.test.ts b/packages/scheduler/lib/db/helpers.test.ts index 6fb50e4da3..d0b3a74ff3 100644 --- a/packages/scheduler/lib/db/helpers.test.ts +++ b/packages/scheduler/lib/db/helpers.test.ts @@ -1,8 +1,10 @@ import { DatabaseClient, defaultDatabaseClientOptions } from './client.js'; -export const getTestDbClient = () => +// Every suite passes its own schema. `migrate()` creates it and `clearDatabase()` drops it, +// so sharing one means a teardown in one file destroys the schema another file is still using. +export const getTestDbClient = (schema: string) => new DatabaseClient({ ...defaultDatabaseClientOptions, url: `postgres://${process.env['NANGO_DB_USER']}:${process.env['NANGO_DB_PASSWORD']}@${process.env['NANGO_DB_HOST']}:${process.env['NANGO_DB_PORT']}/${process.env['NANGO_DB_NAME']}`, - schema: 'scheduler' + schema }); diff --git a/packages/scheduler/lib/models/schedules.integration.test.ts b/packages/scheduler/lib/models/schedules.integration.test.ts index 78e061c876..8de903dfd2 100644 --- a/packages/scheduler/lib/models/schedules.integration.test.ts +++ b/packages/scheduler/lib/models/schedules.integration.test.ts @@ -1,7 +1,7 @@ import { setTimeout } from 'timers/promises'; import { uuidv7 } from 'uuidv7'; -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterAll, afterEach, beforeEach, describe, expect, it } from 'vitest'; import { getTestDbClient } from '../db/helpers.test.js'; import * as schedules from './schedules.js'; @@ -10,7 +10,7 @@ import type { Schedule } from '../types.js'; import type knex from 'knex'; describe('Schedules', () => { - const dbClient = getTestDbClient(); + const dbClient = getTestDbClient('scheduler_schedules'); const db = dbClient.db; beforeEach(async () => { await dbClient.migrate(); @@ -19,6 +19,12 @@ describe('Schedules', () => { await dbClient.clearDatabase(); }); + // Close the knex pool. Nine of these suites leak one otherwise, which exhausts Postgres + // once they share a process with the rest of the suite. + afterAll(async () => { + await dbClient.destroy(); + }); + it('should be successfully created', async () => { const schedule = await createSchedule(db); expect(schedule).toMatchObject({ diff --git a/packages/scheduler/lib/models/tasks.integration.test.ts b/packages/scheduler/lib/models/tasks.integration.test.ts index 3af6195344..21b36b6b78 100644 --- a/packages/scheduler/lib/models/tasks.integration.test.ts +++ b/packages/scheduler/lib/models/tasks.integration.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterAll, afterEach, beforeEach, describe, expect, it } from 'vitest'; import { nanoid } from '@nangohq/utils'; @@ -27,7 +27,7 @@ const props = { }; describe('Task', () => { - const dbClient = getTestDbClient(); + const dbClient = getTestDbClient('scheduler_tasks'); const db = dbClient.db; beforeEach(async () => { await dbClient.migrate(); @@ -37,6 +37,12 @@ describe('Task', () => { await dbClient.clearDatabase(); }); + // Close the knex pool. Nine of these suites leak one otherwise, which exhausts Postgres + // once they share a process with the rest of the suite. + afterAll(async () => { + await dbClient.destroy(); + }); + it('should be successfully created', async () => { const res = (await tasks.create(db, [props])).unwrap(); expect(res.created[0]).toMatchObject({ diff --git a/packages/scheduler/lib/scheduler.integration.test.ts b/packages/scheduler/lib/scheduler.integration.test.ts index e9f18d0877..52eee5309f 100644 --- a/packages/scheduler/lib/scheduler.integration.test.ts +++ b/packages/scheduler/lib/scheduler.integration.test.ts @@ -12,7 +12,7 @@ import type { TaskProps } from './models/tasks.js'; import type { Schedule, ScheduleState, Task } from './types.js'; describe('Scheduler', () => { - const dbClient = getTestDbClient(); + const dbClient = getTestDbClient('scheduler'); const callbacks = { CREATED: vi.fn((task: Task) => expect(task.state).toBe('CREATED')), STARTED: vi.fn((task: Task) => expect(task.state).toBe('STARTED')), @@ -40,6 +40,7 @@ describe('Scheduler', () => { afterAll(async () => { await scheduler.stop(); await dbClient.clearDatabase(); + await dbClient.destroy(); }); it('mark task as SUCCEEDED', async () => { diff --git a/vite.integration.config.ts b/vite.integration.config.ts index 58a15ed816..a7d3b082cb 100644 --- a/vite.integration.config.ts +++ b/vite.integration.config.ts @@ -34,15 +34,7 @@ function findSuitesUsingViMock(root: string): string[] { return found; } -const usesViMock = findSuitesUsingViMock('packages'); - -// These all run live polling daemons against the fixed `scheduler` schema, so a daemon from -// one file steals dequeues and transitions tasks in the next. Orchestrator counts because its -// harness builds a Scheduler on `getTestDbClient` from @nangohq/scheduler, which pins that same -// schema, and tears down with `clearDatabase()`: -const ownsTheSchedulerSchema = ['**/packages/scheduler/**/*.integration.test.ts', '**/packages/orchestrator/**/*.integration.test.ts']; - -const needsOwnProcess = [...usesViMock, ...ownsTheSchedulerSchema]; +const needsOwnProcess = findSuitesUsingViMock('packages'); const shared = { include: ['**/*.integration.{test,spec}.?(c|m)[jt]s?(x)'],