From 464428a81d0e2e87210cc4b6bcbef4ae1da53361 Mon Sep 17 00:00:00 2001 From: imeepos Date: Sun, 5 Jul 2026 23:22:16 -0700 Subject: [PATCH] feat: implement database functions migration and testing --- scripts/migrations/database-functions.ts | 52 +++++++++++++++++++ scripts/migrations/migrate-functions.ts | 7 +-- server/src/index.ts | 23 +++++++- .../unit/db/migrate-functions-script.test.ts | 41 +++++++++++++++ 4 files changed, 116 insertions(+), 7 deletions(-) create mode 100644 scripts/migrations/database-functions.ts create mode 100644 server/tests/unit/db/migrate-functions-script.test.ts diff --git a/scripts/migrations/database-functions.ts b/scripts/migrations/database-functions.ts new file mode 100644 index 000000000..3f158b176 --- /dev/null +++ b/scripts/migrations/database-functions.ts @@ -0,0 +1,52 @@ +import { readFileSync } from "node:fs"; +import { join } from "node:path"; +import postgres from "postgres"; + +export const SQL_FUNCTION_FILES = [ + "deductFromRollovers.sql", + "deductFromMainBalance.sql", + "unwindFromLockReceipt.sql", + "getTotalBalance.sql", + "deductFromAdditionalBalance.sql", + "getAvailableOverageFromSpendLimit.sql", + "performDeduction.sql", + "syncBalances.sql", + "syncBalancesV2.sql", + "resetCusEnts.sql", +] as const; + +const SQL_DIR = join( + import.meta.dir, + "..", + "..", + "server", + "src", + "internal", + "balances", + "utils", + "sql", +); + +export const initializeDatabaseFunctions = async ( + databaseUrl = process.env.DATABASE_URL, +) => { + if (!databaseUrl) { + throw new Error("DATABASE_URL is required to initialize database functions"); + } + + const sql = postgres(databaseUrl, { + max: 1, + prepare: false, + connect_timeout: 30, + }); + + try { + for (const file of SQL_FUNCTION_FILES) { + const body = readFileSync(join(SQL_DIR, file), "utf8"); + await sql.unsafe(body); + console.log(`loaded ${file}`); + } + } finally { + await sql.end({ timeout: 5 }); + } +}; diff --git a/scripts/migrations/migrate-functions.ts b/scripts/migrations/migrate-functions.ts index 302d6575b..fe15ff40a 100644 --- a/scripts/migrations/migrate-functions.ts +++ b/scripts/migrations/migrate-functions.ts @@ -1,5 +1,6 @@ import { loadLocalEnv } from "@server/utils/envUtils"; import inquirer from "inquirer"; +import { initializeDatabaseFunctions } from "./database-functions.ts"; loadLocalEnv(); @@ -15,10 +16,6 @@ if (!process.env.DATABASE_URL?.includes("us-east-2")) { } export const migrateFunctions = async () => { - // Dynamic import to ensure env is loaded first - const { initializeDatabaseFunctions } = await import( - "@server/db/initializeDatabaseFunctions" - ); const databaseUrl = process.env.DATABASE_URL; if (databaseUrl?.includes("us-east-2")) { @@ -38,7 +35,7 @@ export const migrateFunctions = async () => { } } - await initializeDatabaseFunctions(); + await initializeDatabaseFunctions(databaseUrl); }; await migrateFunctions(); diff --git a/server/src/index.ts b/server/src/index.ts index 270ddd617..588e9e46f 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -130,6 +130,23 @@ const getWebhookTraceObject = (event: unknown) => { : undefined; }; +const toHex = async (value: string): Promise => { + const bytes = new TextEncoder().encode(value); + const digest = await crypto.subtle.digest("SHA-256", bytes); + return Array.from(new Uint8Array(digest)) + .slice(0, 8) + .map((byte) => byte.toString(16).padStart(2, "0")) + .join(""); +}; + +const describeRawBodyForTrace = async (rawBody: string) => ({ + raw_body_length: rawBody.length, + raw_body_sha256_prefix: await toHex(rawBody), + raw_body_tail_char_codes: Array.from(rawBody.slice(-8)).map((char) => + char.charCodeAt(0), + ), +}); + const appendStripeWebhookIngressTrace = async ({ kv, key, @@ -163,6 +180,7 @@ const recordStripeWebhookIngressTrace = async ( cf_ray: request.headers.get("cf-ray"), user_agent: request.headers.get("user-agent"), }; + let rawBody = ""; try { if (traceKv) { @@ -177,7 +195,7 @@ const recordStripeWebhookIngressTrace = async ( }); } - const rawBody = await request.clone().text(); + rawBody = await request.clone().text(); const event = JSON.parse(rawBody) as { id?: string; type?: string; @@ -194,7 +212,7 @@ const recordStripeWebhookIngressTrace = async ( created_at: receivedAt, data: { ...baseData, - raw_body_length: rawBody.length, + ...(await describeRawBodyForTrace(rawBody)), }, }; @@ -214,6 +232,7 @@ const recordStripeWebhookIngressTrace = async ( data: { ...baseData, error: error instanceof Error ? error.message : String(error), + ...(await describeRawBodyForTrace(rawBody)), }, }; diff --git a/server/tests/unit/db/migrate-functions-script.test.ts b/server/tests/unit/db/migrate-functions-script.test.ts new file mode 100644 index 000000000..fa943faec --- /dev/null +++ b/server/tests/unit/db/migrate-functions-script.test.ts @@ -0,0 +1,41 @@ +import { existsSync, readFileSync } from "node:fs"; +import { join } from "node:path"; +import { describe, expect, test } from "bun:test"; + +const serverRoot = join(import.meta.dir, "../../.."); +const projectRoot = join(serverRoot, ".."); + +const readProjectSource = (relativePath: string) => + readFileSync(join(projectRoot, relativePath), "utf8"); + +describe("database function migration script", () => { + test("uses a Node script module instead of a removed Worker source module", () => { + const source = readProjectSource("scripts/migrations/migrate-functions.ts"); + + expect(source).not.toContain("@server/db/initializeDatabaseFunctions"); + expect( + existsSync(join(projectRoot, "scripts/migrations/database-functions.ts")), + ).toBe(true); + }); + + test("installs the deduction entrypoint and helper functions in dependency order", () => { + const source = readProjectSource("scripts/migrations/database-functions.ts"); + const expectedFiles = [ + "deductFromRollovers.sql", + "deductFromMainBalance.sql", + "unwindFromLockReceipt.sql", + "getTotalBalance.sql", + "deductFromAdditionalBalance.sql", + "getAvailableOverageFromSpendLimit.sql", + "performDeduction.sql", + "syncBalances.sql", + "syncBalancesV2.sql", + "resetCusEnts.sql", + ]; + + const positions = expectedFiles.map((file) => source.indexOf(file)); + + expect(positions.every((position) => position >= 0)).toBe(true); + expect(positions).toEqual([...positions].sort((a, b) => a - b)); + }); +});