mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
refactor(cloud): share PostgreSQL schema startup between services
This commit is contained in:
@@ -3,11 +3,13 @@ WORKDIR /app
|
||||
RUN corepack enable
|
||||
COPY package.json pnpm-lock.yaml pnpm-workspace.yaml tsconfig.base.json ./
|
||||
COPY packages/relay-contract/package.json packages/relay-contract/package.json
|
||||
COPY packages/postgres-schema/package.json packages/postgres-schema/package.json
|
||||
COPY apps/relay/package.json apps/relay/package.json
|
||||
RUN pnpm install --frozen-lockfile
|
||||
COPY packages/relay-contract packages/relay-contract
|
||||
COPY apps/relay apps/relay
|
||||
RUN pnpm --filter @orca-cloud/relay-contract build && pnpm --filter @orca-cloud/relay build
|
||||
COPY packages/postgres-schema packages/postgres-schema
|
||||
RUN pnpm --filter @orca-cloud/postgres-schema build && pnpm --filter @orca-cloud/relay-contract build && pnpm --filter @orca-cloud/relay build
|
||||
|
||||
FROM node:24-alpine AS runtime
|
||||
ENV NODE_ENV=production
|
||||
@@ -16,8 +18,10 @@ WORKDIR /app
|
||||
RUN corepack enable
|
||||
COPY package.json pnpm-lock.yaml pnpm-workspace.yaml ./
|
||||
COPY packages/relay-contract/package.json packages/relay-contract/package.json
|
||||
COPY packages/postgres-schema/package.json packages/postgres-schema/package.json
|
||||
COPY apps/relay/package.json apps/relay/package.json
|
||||
COPY --from=build /app/packages/relay-contract/dist packages/relay-contract/dist
|
||||
COPY --from=build /app/packages/postgres-schema/dist packages/postgres-schema/dist
|
||||
COPY --from=build /app/apps/relay/dist apps/relay/dist
|
||||
RUN pnpm install --prod --frozen-lockfile --filter @orca-cloud/relay...
|
||||
USER node
|
||||
|
||||
@@ -9,13 +9,14 @@
|
||||
"clean": "node -e \"require('fs').rmSync('dist', { recursive: true, force: true })\"",
|
||||
"dev": "tsx watch src/index.ts",
|
||||
"lint": "tsc -p tsconfig.json --noEmit",
|
||||
"pretest": "pnpm --filter @orca-cloud/relay-contract build",
|
||||
"pretest": "pnpm --filter @orca-cloud/postgres-schema build && pnpm --filter @orca-cloud/relay-contract build",
|
||||
"start": "node dist/index.js",
|
||||
"test": "vitest run",
|
||||
"typecheck": "tsc -p tsconfig.json --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@hono/node-server": "^1.19.17",
|
||||
"@orca-cloud/postgres-schema": "workspace:*",
|
||||
"@orca-cloud/relay-contract": "workspace:*",
|
||||
"hono": "^4.13.7",
|
||||
"jose": "^6.1.3",
|
||||
|
||||
@@ -1,120 +1 @@
|
||||
const RETRYABLE_SCHEMA_CODES = new Set(['55P03', '57014'])
|
||||
const DEFAULT_RETRY_DEADLINE_MS = 30_000
|
||||
const RETRY_BASE_DELAY_MS = 250
|
||||
const RETRY_MAX_DELAY_MS = 2_000
|
||||
|
||||
type SchemaStartupOptions = {
|
||||
now?: () => number
|
||||
random?: () => number
|
||||
retryDeadlineMs?: number
|
||||
wait?: (delayMs: number) => Promise<void>
|
||||
}
|
||||
|
||||
function retryDelayMs(attempt: number, random: () => number): number {
|
||||
const ceiling = Math.min(
|
||||
RETRY_BASE_DELAY_MS * 2 ** (attempt - 1),
|
||||
RETRY_MAX_DELAY_MS
|
||||
)
|
||||
return Math.ceil(ceiling * (0.5 + random() * 0.5))
|
||||
}
|
||||
|
||||
function wait(delayMs: number): Promise<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, delayMs))
|
||||
}
|
||||
|
||||
const CREATE_TABLE_IF_NOT_EXISTS = /^\s*CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\b/i
|
||||
const CREATE_INDEX_IF_NOT_EXISTS = /^\s*CREATE\s+(?:UNIQUE\s+)?INDEX\s+IF\s+NOT\s+EXISTS\b/i
|
||||
|
||||
// `IF NOT EXISTS` only checks the name before the catalog inserts, so the loser of a concurrent
|
||||
// CREATE can fail on the catalog unique index (23505) or, when the winner has already committed by
|
||||
// the time the loser reaches TypeCreate/heap_create_with_catalog, on the name check those routines
|
||||
// repeat (42710 duplicate type, 42P07 duplicate relation). Each is a no-op on the next attempt.
|
||||
function concurrentCreateCollision(
|
||||
value: { code?: unknown; constraint?: unknown },
|
||||
statement: string
|
||||
): boolean {
|
||||
if (CREATE_TABLE_IF_NOT_EXISTS.test(statement)) {
|
||||
return (
|
||||
(value.code === '23505' && value.constraint === 'pg_type_typname_nsp_index') ||
|
||||
value.code === '42710' ||
|
||||
value.code === '42P07'
|
||||
)
|
||||
}
|
||||
if (CREATE_INDEX_IF_NOT_EXISTS.test(statement)) {
|
||||
return (
|
||||
(value.code === '23505' && value.constraint === 'pg_class_relname_nsp_index') ||
|
||||
value.code === '42P07'
|
||||
)
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
const ALTER_TABLE_ADD_CONSTRAINT =
|
||||
/^\s*ALTER\s+TABLE\s+\S+\s+ADD\s+CONSTRAINT\b/i
|
||||
|
||||
// Postgres has no `ADD CONSTRAINT IF NOT EXISTS`, so a re-run and a concurrent
|
||||
// startup both land on 42710 once the constraint exists. Unlike a CREATE race
|
||||
// this is terminal, not transient: retrying only repeats it, so the statement
|
||||
// counts as applied.
|
||||
function constraintAlreadyApplied(error: unknown, statement: string): boolean {
|
||||
return (
|
||||
ALTER_TABLE_ADD_CONSTRAINT.test(statement) &&
|
||||
(error as { code?: unknown }).code === '42710'
|
||||
)
|
||||
}
|
||||
|
||||
function retryableSchemaError(error: unknown, statement: string): boolean {
|
||||
const value = error as { code?: unknown; constraint?: unknown }
|
||||
return (
|
||||
RETRYABLE_SCHEMA_CODES.has(String(value.code)) || concurrentCreateCollision(value, statement)
|
||||
)
|
||||
}
|
||||
|
||||
export async function applyPostgresSchema(
|
||||
statements: string[],
|
||||
query: (statement: string) => Promise<unknown>,
|
||||
options: SchemaStartupOptions = {}
|
||||
): Promise<void> {
|
||||
const now = options.now ?? Date.now
|
||||
const random = options.random ?? Math.random
|
||||
const pause = options.wait ?? wait
|
||||
const deadlineAt = now() + (options.retryDeadlineMs ?? DEFAULT_RETRY_DEADLINE_MS)
|
||||
|
||||
for (const statement of statements) {
|
||||
let attempt = 1
|
||||
while (true) {
|
||||
try {
|
||||
await query(statement)
|
||||
break
|
||||
} catch (error) {
|
||||
if (constraintAlreadyApplied(error, statement)) break
|
||||
const code = String((error as { code?: unknown }).code)
|
||||
const remainingMs = deadlineAt - now()
|
||||
const retryable = retryableSchemaError(error, statement)
|
||||
if (!retryable || remainingMs <= 0) {
|
||||
if (retryable) {
|
||||
console.warn(
|
||||
JSON.stringify({
|
||||
event: 'orca_relay_postgres_schema_retry_exhausted',
|
||||
code,
|
||||
attempts: attempt
|
||||
})
|
||||
)
|
||||
}
|
||||
throw error
|
||||
}
|
||||
const delayMs = Math.min(remainingMs, retryDelayMs(attempt, random))
|
||||
console.warn(
|
||||
JSON.stringify({
|
||||
event: 'orca_relay_postgres_schema_retry',
|
||||
code,
|
||||
attempt,
|
||||
delayMs
|
||||
})
|
||||
)
|
||||
await pause(delayMs)
|
||||
attempt += 1
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
export { applyPostgresSchema } from '@orca-cloud/postgres-schema'
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
{
|
||||
"name": "@orca-cloud/postgres-schema",
|
||||
"version": "0.0.0",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
"types": "dist/index.d.ts",
|
||||
"scripts": {
|
||||
"build": "tsc -p tsconfig.build.json",
|
||||
"clean": "node -e \"require('fs').rmSync('dist', { recursive: true, force: true })\"",
|
||||
"lint": "tsc -p tsconfig.json --noEmit",
|
||||
"test": "pnpm build",
|
||||
"typecheck": "tsc -p tsconfig.json --noEmit"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^24.10.0",
|
||||
"typescript": "^5.9.3",
|
||||
"vitest": "^4.0.8"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,113 @@
|
||||
const RETRYABLE_SCHEMA_CODES = new Set(['55P03', '57014'])
|
||||
const DEFAULT_RETRY_DEADLINE_MS = 30_000
|
||||
const RETRY_BASE_DELAY_MS = 250
|
||||
const RETRY_MAX_DELAY_MS = 2_000
|
||||
|
||||
type SchemaStartupOptions = {
|
||||
eventPrefix?: string
|
||||
now?: () => number
|
||||
random?: () => number
|
||||
retryDeadlineMs?: number
|
||||
wait?: (delayMs: number) => Promise<void>
|
||||
}
|
||||
|
||||
function retryDelayMs(attempt: number, random: () => number): number {
|
||||
const ceiling = Math.min(RETRY_BASE_DELAY_MS * 2 ** (attempt - 1), RETRY_MAX_DELAY_MS)
|
||||
return Math.ceil(ceiling * (0.5 + random() * 0.5))
|
||||
}
|
||||
|
||||
function wait(delayMs: number): Promise<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, delayMs))
|
||||
}
|
||||
|
||||
const CREATE_TABLE_IF_NOT_EXISTS = /^\s*CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\b/i
|
||||
const CREATE_INDEX_IF_NOT_EXISTS = /^\s*CREATE\s+(?:UNIQUE\s+)?INDEX\s+IF\s+NOT\s+EXISTS\b/i
|
||||
|
||||
// `IF NOT EXISTS` only checks the name before the catalog inserts, so the loser of a concurrent
|
||||
// CREATE can fail on the catalog unique index (23505) or, when the winner has already committed by
|
||||
// the time the loser reaches TypeCreate/heap_create_with_catalog, on the name check those routines
|
||||
// repeat (42710 duplicate type, 42P07 duplicate relation). Each is a no-op on the next attempt.
|
||||
function concurrentCreateCollision(
|
||||
value: { code?: unknown; constraint?: unknown },
|
||||
statement: string
|
||||
): boolean {
|
||||
if (CREATE_TABLE_IF_NOT_EXISTS.test(statement)) {
|
||||
return (
|
||||
(value.code === '23505' && value.constraint === 'pg_type_typname_nsp_index') ||
|
||||
value.code === '42710' ||
|
||||
value.code === '42P07'
|
||||
)
|
||||
}
|
||||
if (CREATE_INDEX_IF_NOT_EXISTS.test(statement)) {
|
||||
return (
|
||||
(value.code === '23505' && value.constraint === 'pg_class_relname_nsp_index') ||
|
||||
value.code === '42P07'
|
||||
)
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
const ALTER_TABLE_ADD_CONSTRAINT = /^\s*ALTER\s+TABLE\s+\S+\s+ADD\s+CONSTRAINT\b/i
|
||||
|
||||
function constraintAlreadyApplied(error: unknown, statement: string): boolean {
|
||||
return (
|
||||
ALTER_TABLE_ADD_CONSTRAINT.test(statement) &&
|
||||
(error as { code?: unknown }).code === '42710'
|
||||
)
|
||||
}
|
||||
|
||||
function retryableSchemaError(error: unknown, statement: string): boolean {
|
||||
const value = error as { code?: unknown; constraint?: unknown }
|
||||
return (
|
||||
RETRYABLE_SCHEMA_CODES.has(String(value.code)) || concurrentCreateCollision(value, statement)
|
||||
)
|
||||
}
|
||||
|
||||
export async function applyPostgresSchema(
|
||||
statements: string[],
|
||||
query: (statement: string) => Promise<unknown>,
|
||||
options: SchemaStartupOptions = {}
|
||||
): Promise<void> {
|
||||
const now = options.now ?? Date.now
|
||||
const random = options.random ?? Math.random
|
||||
const pause = options.wait ?? wait
|
||||
const deadlineAt = now() + (options.retryDeadlineMs ?? DEFAULT_RETRY_DEADLINE_MS)
|
||||
|
||||
for (const statement of statements) {
|
||||
let attempt = 1
|
||||
while (true) {
|
||||
try {
|
||||
await query(statement)
|
||||
break
|
||||
} catch (error) {
|
||||
if (constraintAlreadyApplied(error, statement)) break
|
||||
const code = String((error as { code?: unknown }).code)
|
||||
const remainingMs = deadlineAt - now()
|
||||
const retryable = retryableSchemaError(error, statement)
|
||||
if (!retryable || remainingMs <= 0) {
|
||||
if (retryable) {
|
||||
console.warn(
|
||||
JSON.stringify({
|
||||
event: `${options.eventPrefix ?? 'orca_relay_postgres_schema'}_retry_exhausted`,
|
||||
code,
|
||||
attempts: attempt
|
||||
})
|
||||
)
|
||||
}
|
||||
throw error
|
||||
}
|
||||
const delayMs = Math.min(remainingMs, retryDelayMs(attempt, random))
|
||||
console.warn(
|
||||
JSON.stringify({
|
||||
event: `${options.eventPrefix ?? 'orca_relay_postgres_schema'}_retry`,
|
||||
code,
|
||||
attempt,
|
||||
delayMs
|
||||
})
|
||||
)
|
||||
await pause(delayMs)
|
||||
attempt += 1
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
{
|
||||
"extends": "./tsconfig.json",
|
||||
"compilerOptions": {
|
||||
"declaration": true,
|
||||
"emitDeclarationOnly": false,
|
||||
"noEmit": false,
|
||||
"outDir": "dist",
|
||||
"rootDir": "src"
|
||||
},
|
||||
"exclude": ["src/**/*.test.ts"]
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
{
|
||||
"extends": "../../tsconfig.base.json",
|
||||
"compilerOptions": { "noEmit": true },
|
||||
"include": ["src/**/*.ts"]
|
||||
}
|
||||
Generated
+16
@@ -23,6 +23,9 @@ importers:
|
||||
|
||||
apps/relay:
|
||||
dependencies:
|
||||
'@orca-cloud/postgres-schema':
|
||||
specifier: workspace:*
|
||||
version: link:../../packages/postgres-schema
|
||||
'@hono/node-server':
|
||||
specifier: ^1.19.17
|
||||
version: 1.19.17(hono@4.13.7)
|
||||
@@ -117,6 +120,19 @@ importers:
|
||||
specifier: ^4.0.8
|
||||
version: 4.1.9(@types/node@24.13.2)(vite@8.0.16(@types/node@24.13.2)(esbuild@0.28.1)(tsx@4.22.4))
|
||||
|
||||
packages/postgres-schema:
|
||||
devDependencies:
|
||||
'@types/node':
|
||||
specifier: ^24.10.0
|
||||
version: 24.13.2
|
||||
typescript:
|
||||
specifier: ^5.9.3
|
||||
version: 5.9.3
|
||||
vitest:
|
||||
specifier: ^4.0.8
|
||||
version: 4.1.9(@types/node@24.13.2)(vite@8.0.16(@types/node@24.13.2)(esbuild@0.28.1)(tsx@4.22.4))
|
||||
|
||||
|
||||
packages/relay-contract:
|
||||
dependencies:
|
||||
zod:
|
||||
|
||||
Reference in New Issue
Block a user