diff --git a/pkgs/edge-worker/src/core/Queries.ts b/pkgs/edge-worker/src/core/Queries.ts index f0463aa17..73918572e 100644 --- a/pkgs/edge-worker/src/core/Queries.ts +++ b/pkgs/edge-worker/src/core/Queries.ts @@ -82,6 +82,13 @@ export class Queries { * pinging during startup (debounce). */ async trackWorkerFunction(functionName: string, startMode: WorkerStartMode = 'http'): Promise { + if (startMode === 'http') { + await this.sql` + SELECT pgflow.track_worker_function(${functionName}) + `; + return; + } + await this.sql` SELECT pgflow.track_worker_function(${functionName}, ${startMode}) `; diff --git a/pkgs/edge-worker/src/core/Worker.ts b/pkgs/edge-worker/src/core/Worker.ts index 11dccfac3..da3ea960d 100644 --- a/pkgs/edge-worker/src/core/Worker.ts +++ b/pkgs/edge-worker/src/core/Worker.ts @@ -67,7 +67,7 @@ export class Worker { if (!this.isMainLoopActive) { this.logDeprecation(); - if (this.lifecycle.isDeprecated) { + if (this.isDeprecated) { this.deprecationHandler?.(); } break; @@ -142,7 +142,7 @@ export class Worker { } get isStarting() { - return this.lifecycle.isStarting; + return this.lifecycle.isStarting ?? false; } get isRunning() { @@ -150,7 +150,7 @@ export class Worker { } get isDeprecated() { - return this.lifecycle.isDeprecated; + return this.lifecycle.isDeprecated ?? false; } get isStopped() { diff --git a/pkgs/edge-worker/src/core/WorkerLifecycle.ts b/pkgs/edge-worker/src/core/WorkerLifecycle.ts index dbc3e693a..37b08aa0f 100644 --- a/pkgs/edge-worker/src/core/WorkerLifecycle.ts +++ b/pkgs/edge-worker/src/core/WorkerLifecycle.ts @@ -1,6 +1,6 @@ import type { Queries } from './Queries.js'; import type { Queue } from '../queue/Queue.js'; -import type { ILifecycle, Json, WorkerBootstrap, WorkerRow } from './types.js'; +import type { InternalLifecycle, Json, WorkerBootstrap, WorkerRow } from './types.js'; import { States, WorkerState } from './WorkerState.js'; import type { Logger } from '../platform/types.js'; @@ -9,7 +9,7 @@ export interface LifecycleConfig { heartbeatInterval?: number; } -export class WorkerLifecycle implements ILifecycle { +export class WorkerLifecycle implements InternalLifecycle { private workerState: WorkerState; private logger: Logger; private queries: Queries; diff --git a/pkgs/edge-worker/src/core/types.ts b/pkgs/edge-worker/src/core/types.ts index be30b8058..e74b1e1dc 100644 --- a/pkgs/edge-worker/src/core/types.ts +++ b/pkgs/edge-worker/src/core/types.ts @@ -27,15 +27,20 @@ export interface ILifecycle { get edgeFunctionName(): string | undefined; get queueName(): string; get isCreated(): boolean; - get isStarting(): boolean; + readonly isStarting?: boolean; get isRunning(): boolean; - get isDeprecated(): boolean; + readonly isDeprecated?: boolean; get isStopping(): boolean; get isStopped(): boolean; transitionToStopping(): void; } +export interface InternalLifecycle extends ILifecycle { + readonly isStarting: boolean; + readonly isDeprecated: boolean; +} + export interface IBatchProcessor { processBatch(): Promise; awaitCompletion(): Promise; diff --git a/pkgs/edge-worker/src/flow/FlowWorkerLifecycle.ts b/pkgs/edge-worker/src/flow/FlowWorkerLifecycle.ts index 68f6e4084..a6608fd8e 100644 --- a/pkgs/edge-worker/src/flow/FlowWorkerLifecycle.ts +++ b/pkgs/edge-worker/src/flow/FlowWorkerLifecycle.ts @@ -1,5 +1,5 @@ import type { Queries } from '../core/Queries.js'; -import type { ILifecycle, WorkerBootstrap, WorkerRow } from '../core/types.js'; +import type { InternalLifecycle, WorkerBootstrap, WorkerRow } from '../core/types.js'; import type { Logger, StartupContext } from '../platform/types.js'; import { States, WorkerState } from '../core/WorkerState.js'; import type { AnyFlow } from '@pgflow/dsl'; @@ -21,7 +21,7 @@ type CompilationStatus = 'compiled' | 'verified' | 'recompiled' | 'mismatch'; /** * A specialized WorkerLifecycle for Flow-based workers that is aware of the Flow's step types */ -export class FlowWorkerLifecycle implements ILifecycle { +export class FlowWorkerLifecycle implements InternalLifecycle { private workerState: WorkerState; private logger: Logger; private queries: Queries; diff --git a/pkgs/edge-worker/src/platform/ProcessPlatformAdapter.ts b/pkgs/edge-worker/src/platform/ProcessPlatformAdapter.ts index 2b6201a12..c9cece8a4 100644 --- a/pkgs/edge-worker/src/platform/ProcessPlatformAdapter.ts +++ b/pkgs/edge-worker/src/platform/ProcessPlatformAdapter.ts @@ -37,9 +37,13 @@ export class ProcessPlatformAdapter implements PlatformAdapter | null = null; private stopPromise: Promise | null = null; private gracefulExitPromise: Promise | null = null; + private cleanupPromise: Promise | null = null; + private signalHandlersRegistered = false; private signalCount = 0; + private readonly signalHandler = () => this.handleSignal(); constructor( options?: ProcessAdapterOptions, @@ -48,41 +52,34 @@ export class ProcessPlatformAdapter implements PlatformAdapter { - const workerName = this.validatedEnv.WORKER_NAME || 'pgflow-worker'; - const workerId = this.deps.randomUUID(); - - this.workerId = workerId; - this.loggingFactory.setWorkerId(workerId); - this.loggingFactory.setWorkerName(workerName); - - this.worker = createWorkerFn(this.loggingFactory.createLogger); - - this.worker.onDeprecated(() => { - this.handleDeprecation().catch(() => undefined); - }); - - await this.worker.startOnlyOnce({ - edgeFunctionName: workerName, - workerId, - startMode: 'process', - }); - + startWorker(createWorkerFn: CreateWorkerFn): Promise { this.registerSignalHandlers(); + this.startupPromise ??= this.performStartWorker(createWorkerFn).catch( + async (error) => { + await this.cleanup(); + throw error; + } + ); + return this.startupPromise; } stopWorker(): Promise { @@ -126,9 +123,41 @@ export class ProcessPlatformAdapter implements PlatformAdapter { + const workerName = this.validatedEnv.WORKER_NAME || 'pgflow-worker'; + const workerId = this.deps.randomUUID(); + + this.workerId = workerId; + this.loggingFactory.setWorkerId(workerId); + this.loggingFactory.setWorkerName(workerName); + + this.worker = createWorkerFn(this.loggingFactory.createLogger); + this.worker.onDeprecated(() => { + this.handleDeprecation().catch(() => undefined); + }); + + await this.worker.startOnlyOnce({ + edgeFunctionName: workerName, + workerId, + startMode: 'process', + }); + } + private registerSignalHandlers(): void { + if (this.signalHandlersRegistered) return; + + this.signalHandlersRegistered = true; for (const signal of ['SIGTERM', 'SIGINT', 'SIGQUIT'] satisfies ProcessSignal[]) { - this.deps.onSignal(signal, () => this.handleSignal()); + this.deps.onSignal(signal, this.signalHandler); + } + } + + private removeSignalHandlers(): void { + if (!this.signalHandlersRegistered) return; + + this.signalHandlersRegistered = false; + for (const signal of ['SIGTERM', 'SIGINT', 'SIGQUIT'] satisfies ProcessSignal[]) { + this.deps.offSignal?.(signal, this.signalHandler); } } @@ -136,6 +165,7 @@ export class ProcessPlatformAdapter implements PlatformAdapter { + this.cleanupPromise ??= (async () => { + this.removeSignalHandlers(); if (this.ownsSql) { await this._platformResources.sql.end(); } - } + })(); + return this.cleanupPromise; } private gracefulExit(): Promise { diff --git a/pkgs/edge-worker/src/platform/SupabasePlatformAdapter.ts b/pkgs/edge-worker/src/platform/SupabasePlatformAdapter.ts index 000a1734b..426670108 100644 --- a/pkgs/edge-worker/src/platform/SupabasePlatformAdapter.ts +++ b/pkgs/edge-worker/src/platform/SupabasePlatformAdapter.ts @@ -222,15 +222,20 @@ export class SupabasePlatformAdapter implements PlatformAdapter; onSignal: (signal: ProcessSignal, handler: () => void | Promise) => void; + offSignal?: (signal: ProcessSignal, handler: () => void | Promise) => void; exit: (code: number) => never; setExitCode: (code: number) => void; randomUUID: () => string; @@ -11,6 +12,7 @@ export type ProcessDeps = { type ProcessLike = { env?: Record; on?: (signal: ProcessSignal, handler: () => void | Promise) => void; + off?: (signal: ProcessSignal, handler: () => void | Promise) => void; exit?: (code: number) => never; exitCode?: number; }; @@ -23,13 +25,20 @@ export function getProcessDeps(): ProcessDeps { const processLike = (globalThis as { process?: ProcessLike }).process; const cryptoLike = (globalThis as { crypto?: CryptoLike }).crypto; - if (!processLike?.env || !processLike.on || !processLike.exit || !cryptoLike?.randomUUID) { + if ( + !processLike?.env || + !processLike.on || + !processLike.off || + !processLike.exit || + !cryptoLike?.randomUUID + ) { throw new Error('Process runtime is not available'); } return { env: processLike.env, onSignal: (signal, handler) => processLike.on?.(signal, handler), + offSignal: (signal, handler) => processLike.off?.(signal, handler), exit: (code) => processLike.exit!(code), setExitCode: (code) => { processLike.exitCode = code; diff --git a/pkgs/edge-worker/src/platform/resolveConnection.ts b/pkgs/edge-worker/src/platform/resolveConnection.ts index a72fdb6a9..4c6bd0ff6 100644 --- a/pkgs/edge-worker/src/platform/resolveConnection.ts +++ b/pkgs/edge-worker/src/platform/resolveConnection.ts @@ -20,36 +20,37 @@ export interface ConnectionEnv extends Record { export interface ConnectionOptions { hasSql?: boolean; connectionString?: string; + allowDatabaseUrl?: boolean; } -export interface SqlConnectionOptions { +export interface SqlConnectionOptions extends ConnectionOptions { sql?: postgres.Sql; - connectionString?: string; maxPgConnections?: number; } /** - * Resolves the connection string based on priority: - * config.sql -> config.connectionString -> DATABASE_URL -> EDGE_WORKER_DB_URL -> local fallback + * Resolves the connection string based on priority. Supabase workers ignore + * DATABASE_URL by default; process workers opt in to DATABASE_URL priority. */ export function resolveConnectionString( env: ConnectionEnv, options?: ConnectionOptions ): string | undefined { - const isLocal = isLocalSupabaseEnv(env); + const envConnectionString = options?.allowDatabaseUrl + ? env.DATABASE_URL || env.EDGE_WORKER_DB_URL + : env.EDGE_WORKER_DB_URL; // Zero-config local dev: use docker pooler when nothing else is configured if ( - isLocal && + isLocalSupabaseEnv(env) && !options?.hasSql && !options?.connectionString && - !env.DATABASE_URL && - !env.EDGE_WORKER_DB_URL + !envConnectionString ) { return DOCKER_TRANSACTION_POOLER_URL; } - return options?.connectionString || env.DATABASE_URL || env.EDGE_WORKER_DB_URL; + return options?.connectionString || envConnectionString; } /** @@ -68,12 +69,7 @@ export function assertConnectionAvailable( } /** - * Resolves and creates the SQL connection based on priority: - * 1. config.sql - User-provided SQL client (highest priority) - * 2. config.connectionString - User-provided connection string - * 3. DATABASE_URL - Environment variable - * 4. EDGE_WORKER_DB_URL - Environment variable - * 5. Local Supabase detection + Docker URL (lowest priority) + * Resolves and creates the SQL connection using resolveConnectionString(). * * @throws Error if no connection source is available */ @@ -81,31 +77,16 @@ export function resolveSqlConnection( env: ConnectionEnv, options?: SqlConnectionOptions ): postgres.Sql { - // 1. config.sql - highest priority if (options?.sql) { return options.sql; } - const max = options?.maxPgConnections ?? 4; - - // 2. config.connectionString - if (options?.connectionString) { - return postgres(options.connectionString, { prepare: false, max }); - } - - // 3. DATABASE_URL - if (env.DATABASE_URL) { - return postgres(env.DATABASE_URL, { prepare: false, max }); - } - - // 4. EDGE_WORKER_DB_URL - if (env.EDGE_WORKER_DB_URL) { - return postgres(env.EDGE_WORKER_DB_URL, { prepare: false, max }); - } - - // 5. Local Supabase detection + docker URL - if (isLocalSupabaseEnv(env)) { - return postgres(DOCKER_TRANSACTION_POOLER_URL, { prepare: false, max }); + const connectionString = resolveConnectionString(env, options); + if (connectionString) { + return postgres(connectionString, { + prepare: false, + max: options?.maxPgConnections ?? 4, + }); } throw new Error( diff --git a/pkgs/edge-worker/tests/types/platform-adapter.test-d.ts b/pkgs/edge-worker/tests/types/platform-adapter.test-d.ts index 429844168..de2f5fa33 100644 --- a/pkgs/edge-worker/tests/types/platform-adapter.test-d.ts +++ b/pkgs/edge-worker/tests/types/platform-adapter.test-d.ts @@ -1,5 +1,9 @@ import type { PlatformAdapter } from '../../src/platform/types.ts'; -import type { WorkerBootstrap, WorkerStartMode } from '../../src/core/types.ts'; +import type { + ILifecycle, + WorkerBootstrap, + WorkerStartMode, +} from '../../src/core/types.ts'; const adapterWithoutRequestShutdown: PlatformAdapter = { async startWorker() {}, @@ -32,3 +36,30 @@ const processBootstrap: WorkerBootstrap = { }; void processBootstrap; + +const legacyLifecycle: ILifecycle = { + async acknowledgeStart(_workerBootstrap: WorkerBootstrap) {}, + acknowledgeStop() {}, + async sendHeartbeat() {}, + get edgeFunctionName() { + return undefined; + }, + get queueName() { + return 'legacy-queue'; + }, + get isCreated() { + return true; + }, + get isRunning() { + return false; + }, + get isStopping() { + return false; + }, + get isStopped() { + return false; + }, + transitionToStopping() {}, +}; + +void legacyLifecycle; diff --git a/pkgs/edge-worker/tests/unit/Queries.test.ts b/pkgs/edge-worker/tests/unit/Queries.test.ts index 97d08953d..e84920083 100644 --- a/pkgs/edge-worker/tests/unit/Queries.test.ts +++ b/pkgs/edge-worker/tests/unit/Queries.test.ts @@ -26,11 +26,21 @@ Deno.test('Queries.trackWorkerFunction - calls correct SQL function', async () = await queries.trackWorkerFunction('my-edge-function'); assertEquals(calls.length, 1); - assertEquals(calls[0].values, ['my-edge-function', 'http']); + assertEquals(calls[0].values, ['my-edge-function']); // Check that query references the correct function assertEquals(calls[0].query.includes('pgflow.track_worker_function'), true); }); +Deno.test('Queries.trackWorkerFunction - keeps explicit HTTP mode compatible with old databases', async () => { + const { mockSql, calls } = createMockSql(); + const queries = new Queries(mockSql); + + await queries.trackWorkerFunction('http-worker', 'http'); + + assertEquals(calls.length, 1); + assertEquals(calls[0].values, ['http-worker']); +}); + Deno.test('Queries.trackWorkerFunction - passes explicit process start mode', async () => { const { mockSql, calls } = createMockSql(); const queries = new Queries(mockSql); @@ -48,7 +58,7 @@ Deno.test('Queries.trackWorkerFunction - handles special characters in function await queries.trackWorkerFunction('my_function-with-special_chars'); assertEquals(calls.length, 1); - assertEquals(calls[0].values, ['my_function-with-special_chars', 'http']); + assertEquals(calls[0].values, ['my_function-with-special_chars']); }); Deno.test('Queries.markWorkerStopped - calls correct SQL function', async () => { diff --git a/pkgs/edge-worker/tests/unit/platform/ProcessPlatformAdapter.test.ts b/pkgs/edge-worker/tests/unit/platform/ProcessPlatformAdapter.test.ts index 6e1dc442f..7fdb295b8 100644 --- a/pkgs/edge-worker/tests/unit/platform/ProcessPlatformAdapter.test.ts +++ b/pkgs/edge-worker/tests/unit/platform/ProcessPlatformAdapter.test.ts @@ -37,7 +37,13 @@ type WorkerStub = { }; function createDeps(env: Record = {}) { - const handlers = new Map void | Promise>(); + type SignalHandler = () => void | Promise; + const handlers = new Map(); + const offSignal = createSpy<[ProcessSignal, SignalHandler], void>( + (signal, handler) => { + if (handlers.get(signal) === handler) handlers.delete(signal); + } + ); const exit = createSpy<[number], never>((code) => { throw new Error(`exit:${code}`); }); @@ -46,12 +52,13 @@ function createDeps(env: Record = {}) { onSignal: (signal, handler) => { handlers.set(signal, handler); }, + offSignal, exit: exit as unknown as (code: number) => never, setExitCode: createSpy<[number], void>(() => undefined), randomUUID: createSpy<[], string>(() => '00000000-0000-4000-8000-000000000001'), }; - return { deps, handlers, exit }; + return { deps, handlers, exit, offSignal }; } function createSqlStub(events?: string[]): SqlStub { @@ -110,6 +117,18 @@ Deno.test('ProcessPlatformAdapter throws when no database source is available', ); }); +Deno.test('ProcessPlatformAdapter prefers DATABASE_URL when both database variables are set', async () => { + const { deps } = createDeps(validEnv({ + DATABASE_URL: 'postgresql://process:5432/database', + EDGE_WORKER_DB_URL: 'postgresql://edge:5432/database', + })); + const adapter = new ProcessPlatformAdapter(undefined, deps); + + assertEquals(adapter.connectionString, 'postgresql://process:5432/database'); + + await adapter.stopWorker(); +}); + Deno.test('ProcessPlatformAdapter starts immediately with process start mode and generated worker id', async () => { const { deps } = createDeps(validEnv({ WORKER_NAME: 'emails' })); const sql = createSqlStub(); @@ -252,7 +271,7 @@ Deno.test('ProcessPlatformAdapter deprecation drains and exits zero', async () = assertEquals(exit.calls, [[0]]); }); -Deno.test('ProcessPlatformAdapter signal handlers registered after startup readiness', async () => { +Deno.test('ProcessPlatformAdapter handles a first signal during startup after readiness', async () => { let resolveStartup = () => {}; const { deps, handlers } = createDeps(validEnv()); const sql = createSqlStub(); @@ -264,10 +283,61 @@ Deno.test('ProcessPlatformAdapter signal handlers registered after startup readi const startupPromise = adapter.startWorker(() => worker as never); - assertEquals(handlers.size, 0, 'No signal handlers should be registered during startup'); + assertEquals(handlers.size, 3, 'Signal handlers must exist during startup'); + + const shutdownPromise = handlers.get('SIGTERM')?.(); + await Promise.resolve(); + assertEquals(worker.stop.calls.length, 0, 'Worker must not stop while startup is pending'); resolveStartup(); await startupPromise; + await assertRejects(async () => await shutdownPromise, Error, 'exit:0'); + + assertEquals(worker.stop.calls.length, 1); + assertEquals(handlers.size, 0); +}); + +Deno.test('ProcessPlatformAdapter startup failure cleans owned resources once', async () => { + const { deps, handlers, offSignal } = createDeps(validEnv()); + const worker = createWorkerStub(); + worker.startOnlyOnce.implementation = () => Promise.reject(new Error('startup failed')); + const adapter = new ProcessPlatformAdapter(undefined, deps); + const end = createSpy<[], Promise>(() => Promise.resolve()); + (adapter.sql as unknown as { end: typeof end }).end = end; - assertEquals(handlers.size, 3, 'All signal handlers should be registered after startup'); + await assertRejects( + () => adapter.startWorker(() => worker as never), + Error, + 'startup failed' + ); + + assertEquals(end.calls.length, 1); + assertEquals(offSignal.calls.length, 3); + assertEquals(handlers.size, 0); + + await assertRejects(() => adapter.stopWorker(), Error, 'startup failed'); + assertEquals(end.calls.length, 1); + assertEquals(offSignal.calls.length, 3); +}); + +Deno.test('ProcessPlatformAdapter shutdown cleans owned resources once', async () => { + const { deps, handlers, offSignal } = createDeps(validEnv()); + const worker = createWorkerStub(); + const adapter = new ProcessPlatformAdapter(undefined, deps); + const end = createSpy<[], Promise>(() => Promise.resolve()); + (adapter.sql as unknown as { end: typeof end }).end = end; + const markWorkerStopped = createSpy<[string], Promise>(() => Promise.resolve()); + (adapter as unknown as { + queries: { markWorkerStopped: typeof markWorkerStopped }; + }).queries.markWorkerStopped = markWorkerStopped; + + await adapter.startWorker(() => worker as never); + await adapter.stopWorker(); + await adapter.stopWorker(); + + assertEquals(worker.stop.calls.length, 1); + assertEquals(markWorkerStopped.calls.length, 1); + assertEquals(end.calls.length, 1); + assertEquals(offSignal.calls.length, 3); + assertEquals(handlers.size, 0); }); diff --git a/pkgs/edge-worker/tests/unit/platform/SupabasePlatformAdapter.test.ts b/pkgs/edge-worker/tests/unit/platform/SupabasePlatformAdapter.test.ts index 83daa33a6..24d03366d 100644 --- a/pkgs/edge-worker/tests/unit/platform/SupabasePlatformAdapter.test.ts +++ b/pkgs/edge-worker/tests/unit/platform/SupabasePlatformAdapter.test.ts @@ -209,6 +209,30 @@ Deno.test({ }, }); +Deno.test({ + name: 'ignores DATABASE_URL when EDGE_WORKER_DB_URL is also set', + sanitizeResources: false, + fn: () => { + const deps = createMockDeps({ + getEnv: () => ({ + SUPABASE_URL: 'https://abc123.supabase.co', + SUPABASE_ANON_KEY: 'test-anon-key', + SUPABASE_SERVICE_ROLE_KEY: 'test-service-key', + SB_EXECUTION_ID: 'test-exec-id', + DATABASE_URL: 'postgresql://unrelated:5432/database', + EDGE_WORKER_DB_URL: 'postgresql://edge-worker:5432/database', + }), + }); + + const adapter = new SupabasePlatformAdapter(undefined, deps); + + assertEquals( + adapter.connectionString, + 'postgresql://edge-worker:5432/database' + ); + }, +}); + // ============================================================ // Local Environment Detection Tests // ============================================================ @@ -361,6 +385,78 @@ Deno.test({ }, }); +Deno.test({ + name: 'HTTP startup returns a controlled 500 and retries with a fresh worker', + sanitizeResources: false, + fn: async () => { + let serveHandler: ((req: Request) => Response | Promise) | null = null; + let rejectStartup = (_error: Error) => {}; + let createCount = 0; + + const deps = createMockDeps({ + serve: (h) => { + serveHandler = h; + }, + }); + const sql = (() => Promise.resolve([])) as unknown as { + end: () => Promise; + }; + sql.end = () => Promise.resolve(); + + const adapter = new SupabasePlatformAdapter({ sql: sql as never }, deps); + await adapter.startWorker(() => { + createCount++; + return { + startOnlyOnce: () => createCount === 1 + ? new Promise((_resolve, reject) => { + rejectStartup = reject; + }) + : Promise.resolve(), + stop: () => Promise.resolve(), + get isDeprecated() { + return false; + }, + get isStopped() { + return false; + }, + } as unknown as Worker; + }); + + const handler = serveHandler as unknown as (req: Request) => Response | Promise; + const request = () => new Request('http://localhost/functions/v1/my-worker', { + headers: { authorization: 'Bearer test-service-key' }, + }); + + let responseSettled = false; + const responsePromise = Promise.resolve(handler(request())); + responsePromise.then(() => { + responseSettled = true; + }); + await Promise.resolve(); + + assertEquals(responseSettled, false, 'Response must wait for worker readiness'); + + rejectStartup(new Error('sensitive database details')); + const failedResponse = await responsePromise; + const failedBody = await failedResponse.text(); + + assertEquals(failedResponse.status, 500); + assertEquals(failedBody.includes('sensitive database details'), false); + assertEquals(JSON.parse(failedBody), { + error: 'Internal Server Error', + message: 'Internal Server Error', + }); + + const retryResponse = await handler(request()) as Response; + const retryBody = await retryResponse.json(); + + assertEquals(retryResponse.status, 200); + assertEquals(retryBody.status, 'started'); + assertEquals(createCount, 2); + assertEquals(retryBody.workerId === 'test-exec-id', false); + }, +}); + Deno.test({ name: 'HTTP startup uses a fresh worker id when replacing deprecated worker', sanitizeResources: false, diff --git a/pkgs/edge-worker/tests/unit/platform/connectionPriority.test.ts b/pkgs/edge-worker/tests/unit/platform/connectionPriority.test.ts index c88be00e6..bf6b23c1f 100644 --- a/pkgs/edge-worker/tests/unit/platform/connectionPriority.test.ts +++ b/pkgs/edge-worker/tests/unit/platform/connectionPriority.test.ts @@ -58,6 +58,20 @@ Deno.test('connection priority - production uses EDGE_WORKER_DB_URL', () => { assertEquals(result, 'postgresql://prod:5432/db'); }); +Deno.test('connection priority - DATABASE_URL is opt-in for process workers', () => { + const env = { + SUPABASE_URL: 'https://abc123.supabase.co', + DATABASE_URL: 'postgresql://process:5432/db', + EDGE_WORKER_DB_URL: 'postgresql://edge:5432/db', + }; + + assertEquals(resolveConnectionString(env), 'postgresql://edge:5432/db'); + assertEquals( + resolveConnectionString(env, { allowDatabaseUrl: true }), + 'postgresql://process:5432/db' + ); +}); + Deno.test('connection priority - production config.connectionString overrides env var', () => { const env = { SUPABASE_URL: 'https://abc123.supabase.co', diff --git a/pkgs/website/src/content/docs/deploy/database-connection.mdx b/pkgs/website/src/content/docs/deploy/database-connection.mdx index 735c661f9..388a8b24c 100644 --- a/pkgs/website/src/content/docs/deploy/database-connection.mdx +++ b/pkgs/website/src/content/docs/deploy/database-connection.mdx @@ -11,7 +11,9 @@ pgflow auto-detects your database connection in local development and requires m ## Connection Priority -When an Edge Worker starts, pgflow resolves the database connection using this priority chain: +### Supabase Edge Functions + +Supabase Edge Functions use this priority chain: | Priority | Source | Use Case | | -------- | ------------------------- | ---------------------------------------------- | @@ -20,12 +22,26 @@ When an Edge Worker starts, pgflow resolves the database connection using this p | 3 | `EDGE_WORKER_DB_URL` | Production environment variable | | 4 | Local fallback | Supabase local development (automatic) | +`DATABASE_URL` does not affect Supabase Edge Functions. This prevents an unrelated application database URL from redirecting an Edge Worker. + +### Node and Bun Process Workers + +Process workers use this priority chain: + +| Priority | Source | +| -------- | ------------------------- | +| 1 | `config.sql` | +| 2 | `config.connectionString` | +| 3 | `DATABASE_URL` | +| 4 | `EDGE_WORKER_DB_URL` | +| 5 | Local fallback | + If none of these are available, the worker throws an error: `"No database connection available"`. ## Local Development diff --git a/pkgs/website/src/content/docs/deploy/node-bun-process-workers.mdx b/pkgs/website/src/content/docs/deploy/node-bun-process-workers.mdx index 4748ef1aa..2906dfb7e 100644 --- a/pkgs/website/src/content/docs/deploy/node-bun-process-workers.mdx +++ b/pkgs/website/src/content/docs/deploy/node-bun-process-workers.mdx @@ -67,7 +67,15 @@ Set these variables in your process host: | `EDGE_WORKER_DB_URL` | Yes, unless `DATABASE_URL` is set | Alternative PostgreSQL connection string name shared with Supabase Edge Functions. | | `WORKER_NAME` | No | Worker function name recorded in pgflow. Defaults to `pgflow-worker`. | -Use `DATABASE_URL` for process hosts when possible. `EDGE_WORKER_DB_URL` remains supported for deployments that share configuration with Supabase Edge Functions. +Use `DATABASE_URL` for process hosts when possible. `EDGE_WORKER_DB_URL` remains supported for deployments that share configuration with Supabase Edge Functions. If both are set, `DATABASE_URL` takes priority. + + ## Run The Worker