From dcb0fa4c767eb7c60c34572b0e89348804c7e0b7 Mon Sep 17 00:00:00 2001 From: yml2213 Date: Sun, 30 Aug 2026 13:15:47 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96=E8=AE=A2=E5=8D=95=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E4=B8=8E=E6=95=B0=E6=8D=AE=E5=BA=93=E5=85=B3=E9=97=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/repositories/order-item-repo.ts | 93 +++++++------- apps/backend/src/repositories/order-repo.ts | 39 +++++- .../src/services/order/order-service.ts | 118 ++++++++++-------- apps/backend/src/startup/shutdown.test.ts | 66 ++++++++++ apps/backend/src/startup/shutdown.ts | 40 ++++-- 5 files changed, 243 insertions(+), 113 deletions(-) create mode 100644 apps/backend/src/startup/shutdown.test.ts diff --git a/apps/backend/src/repositories/order-item-repo.ts b/apps/backend/src/repositories/order-item-repo.ts index 5d91fb10..bfe4a2d8 100644 --- a/apps/backend/src/repositories/order-item-repo.ts +++ b/apps/backend/src/repositories/order-item-repo.ts @@ -4,7 +4,7 @@ import { query, withTransaction } from '../db/client.js' import type { OrderItemReplaceInput } from '../types/repository/inputs.js' import type { OrderItemRow } from '../types/repository/rows.js' -type QueryExecutor = (text: string, params?: unknown[]) => Promise> +export type QueryExecutor = (text: string, params?: unknown[]) => Promise> type OrderItemSyncPlan = { updates: Array<{ @@ -25,14 +25,22 @@ export async function replaceOrderItems( orderId: number | string, items: OrderItemReplaceInput[], ): Promise { - return withTransaction(async (client: PoolClient) => { - const executor: QueryExecutor = client.query.bind(client) - const existingItems = await listOrderItemsByOrderIdWithExecutor(executor, orderId) - const plan = resolveOrderItemSyncPlan(existingItems, items) + return withTransaction((client: PoolClient) => + replaceOrderItemsWithExecutor(client.query.bind(client), orderId, items), + ) +} - for (const update of plan.updates) { - await executor( - ` +export async function replaceOrderItemsWithExecutor( + executor: QueryExecutor, + orderId: number | string, + items: OrderItemReplaceInput[], +): Promise { + const existingItems = await listOrderItemsByOrderIdWithExecutor(executor, orderId) + const plan = resolveOrderItemSyncPlan(existingItems, items) + + for (const update of plan.updates) { + await executor( + ` UPDATE order_items SET sku_code = $1, @@ -43,21 +51,21 @@ export async function replaceOrderItems( updated_at = $6 WHERE id = $7 `, - [ - update.item.skuCode, - update.item.skuName, - update.item.quantity, - update.item.specJson || '{}', - update.item.itemSnapshotJson || update.item.specJson || '{}', - update.item.updatedAt, - update.orderItemId, - ], - ) - } + [ + update.item.skuCode, + update.item.skuName, + update.item.quantity, + update.item.specJson || '{}', + update.item.itemSnapshotJson || update.item.specJson || '{}', + update.item.updatedAt, + update.orderItemId, + ], + ) + } - for (const create of plan.creates) { - await executor( - ` + for (const create of plan.creates) { + await executor( + ` INSERT INTO order_items ( order_id, sku_code, @@ -69,29 +77,28 @@ export async function replaceOrderItems( updated_at ) VALUES ($1, $2, $3, $4, $5::jsonb, $6::jsonb, $7, $8) `, - [ - Number(orderId), - create.skuCode, - create.skuName, - create.quantity, - create.specJson || '{}', - create.itemSnapshotJson || create.specJson || '{}', - create.createdAt, - create.updatedAt, - ], - ) + [ + Number(orderId), + create.skuCode, + create.skuName, + create.quantity, + create.specJson || '{}', + create.itemSnapshotJson || create.specJson || '{}', + create.createdAt, + create.updatedAt, + ], + ) + } + + if (plan.deletes.length > 0) { + const deletableIds = await listDeletableOrderItemIdsWithExecutor(executor, plan.deletes) + + if (deletableIds.length > 0) { + await executor('DELETE FROM order_items WHERE id = ANY($1::bigint[])', [deletableIds]) } + } - if (plan.deletes.length > 0) { - const deletableIds = await listDeletableOrderItemIdsWithExecutor(executor, plan.deletes) - - if (deletableIds.length > 0) { - await executor('DELETE FROM order_items WHERE id = ANY($1::bigint[])', [deletableIds]) - } - } - - return listOrderItemsByOrderIdWithExecutor(executor, orderId) - }) + return listOrderItemsByOrderIdWithExecutor(executor, orderId) } export function resolveOrderItemSyncPlan( diff --git a/apps/backend/src/repositories/order-repo.ts b/apps/backend/src/repositories/order-repo.ts index 74ef499c..3ea342be 100644 --- a/apps/backend/src/repositories/order-repo.ts +++ b/apps/backend/src/repositories/order-repo.ts @@ -1,3 +1,5 @@ +import type { QueryResult, QueryResultRow } from 'pg' + import { query } from '../db/client.js' import type { OrderCreateInput, @@ -6,6 +8,11 @@ import type { } from '../types/repository/inputs.js' import type { OrderListQueryResult, OrderListRow, OrderRow } from '../types/repository/rows.js' +export type OrderQueryExecutor = ( + text: string, + params?: unknown[], +) => Promise> + type OrderPlatformLookupInput = { provider?: string platform: string @@ -37,7 +44,18 @@ export async function findLatestOrderByPlatformOrderId({ platform, platformOrderId, }: Omit): Promise { - const result = await query( + return findLatestOrderByPlatformOrderIdWithExecutor(query, { + provider, + platform, + platformOrderId, + }) +} + +export async function findLatestOrderByPlatformOrderIdWithExecutor( + executor: OrderQueryExecutor, + { provider = '91kaquan', platform, platformOrderId }: Omit, +): Promise { + const result = await executor( ` SELECT * FROM orders @@ -69,7 +87,14 @@ export async function findLatestOrderByAnyPlatformOrderId( } export async function createOrder(input: OrderCreateInput): Promise { - const result = await query( + return createOrderWithExecutor(query, input) +} + +export async function createOrderWithExecutor( + executor: OrderQueryExecutor, + input: OrderCreateInput, +): Promise { + const result = await executor( ` INSERT INTO orders ( provider, @@ -118,7 +143,15 @@ export async function updateOrder( orderId: number | string, input: OrderUpdateInput, ): Promise { - const result = await query( + return updateOrderWithExecutor(query, orderId, input) +} + +export async function updateOrderWithExecutor( + executor: OrderQueryExecutor, + orderId: number | string, + input: OrderUpdateInput, +): Promise { + const result = await executor( ` UPDATE orders SET diff --git a/apps/backend/src/services/order/order-service.ts b/apps/backend/src/services/order/order-service.ts index 41bfb229..5e12dccd 100644 --- a/apps/backend/src/services/order/order-service.ts +++ b/apps/backend/src/services/order/order-service.ts @@ -1,9 +1,10 @@ import { - createOrder, - findLatestOrderByPlatformOrderId, - updateOrder, + createOrderWithExecutor, + findLatestOrderByPlatformOrderIdWithExecutor, + updateOrderWithExecutor, } from '../../repositories/order-repo.js' -import { replaceOrderItems } from '../../repositories/order-item-repo.js' +import { replaceOrderItemsWithExecutor } from '../../repositories/order-item-repo.js' +import { withTransaction } from '../../db/client.js' import { syncDeliveryTasksForOrder } from './delivery-task-service.js' import { resolveOrderItemForFulfillment } from '../fulfillment/product-resolution-service.js' import { @@ -98,11 +99,6 @@ export async function upsertOrderFromSource( { sourceLabel = 'source' }: UpsertOrderSourceOptions = {}, ): Promise { const now = nowIso() - const existing = await findLatestOrderByPlatformOrderId({ - provider: event.provider, - platform: event.platform, - platformOrderId: event.platformOrderId, - }) logIntegration('[order-service]', `开始处理 ${sourceLabel} 订单 upsert`, { provider: event.provider, @@ -110,7 +106,6 @@ export async function upsertOrderFromSource( shopId: event.shopId, shopName: event.shopName, platformOrderId: event.platformOrderId, - existingOrderId: existing?.id || null, }) const resolvedItems = await Promise.all( @@ -125,56 +120,71 @@ export async function upsertOrderFromSource( const resolvedFulfillmentItems = resolvedItems as FulfillmentOrderItem[] const configuredItems = resolvedFulfillmentItems.filter((item) => item.isConfigured) - const basePayload = { - provider: event.provider, - platform: event.platform, - shopId: String(existing?.shop_id || event.shopId || '').trim(), - shopName: String(existing?.shop_name || event.shopName || '').trim(), - platformOrderId: event.platformOrderId, - orderStatus: event.orderStatus, - payStatus: event.payStatus, - buyerId: event.buyerId, - buyerName: event.buyerName, - receiverContact: event.receiverContact, - totalAmount: event.totalAmount, - currency: event.currency, - rawPayloadJson: JSON.stringify(event.rawPayload), - ...(event.paidAt !== undefined ? { paidAt: event.paidAt } : {}), - } - const mergedPayload = mergeSourceOrderState(existing, basePayload) + const persisted = await withTransaction(async (client) => { + const executor = client.query.bind(client) + const lockKey = [event.provider, event.platform, event.platformOrderId].join('|') + await executor('SELECT pg_advisory_xact_lock(hashtextextended($1, 0))', [lockKey]) - const order = existing - ? await updateOrder(existing.id, { - ...basePayload, - ...mergedPayload, - updatedAt: now, + const existing = await findLatestOrderByPlatformOrderIdWithExecutor(executor, { + provider: event.provider, + platform: event.platform, + platformOrderId: event.platformOrderId, + }) + const basePayload = { + provider: event.provider, + platform: event.platform, + shopId: String(existing?.shop_id || event.shopId || '').trim(), + shopName: String(existing?.shop_name || event.shopName || '').trim(), + platformOrderId: event.platformOrderId, + orderStatus: event.orderStatus, + payStatus: event.payStatus, + buyerId: event.buyerId, + buyerName: event.buyerName, + receiverContact: event.receiverContact, + totalAmount: event.totalAmount, + currency: event.currency, + rawPayloadJson: JSON.stringify(event.rawPayload), + ...(event.paidAt !== undefined ? { paidAt: event.paidAt } : {}), + } + const mergedPayload = mergeSourceOrderState(existing, basePayload) + const order = existing + ? await updateOrderWithExecutor(executor, existing.id, { + ...basePayload, + ...mergedPayload, + updatedAt: now, + }) + : await createOrderWithExecutor(executor, { + ...basePayload, + ...mergedPayload, + createdAt: now, + updatedAt: now, + }) + + if (!order) { + throw createHttpError('订单写入失败', { + statusCode: 500, + errorCode: 'order_write_failed', }) - : await createOrder({ - ...basePayload, - ...mergedPayload, + } + + const orderItems = await replaceOrderItemsWithExecutor( + executor, + order.id, + resolvedFulfillmentItems.map((item) => ({ + skuCode: item.skuCode, + skuName: item.skuName, + quantity: item.quantity, + specJson: JSON.stringify(item.spec || {}), + itemSnapshotJson: JSON.stringify(item.snapshot || item.spec || {}), createdAt: now, updatedAt: now, - }) + })), + ) - if (!order) { - throw createHttpError('订单写入失败', { - statusCode: 500, - errorCode: 'order_write_failed', - }) - } + return { order, orderItems } + }) - const orderItems = await replaceOrderItems( - order.id, - resolvedFulfillmentItems.map((item) => ({ - skuCode: item.skuCode, - skuName: item.skuName, - quantity: item.quantity, - specJson: JSON.stringify(item.spec || {}), - itemSnapshotJson: JSON.stringify(item.snapshot || item.spec || {}), - createdAt: now, - updatedAt: now, - })), - ) + const { order, orderItems } = persisted const configuredOrderItems = orderItems.filter(isOrderItemConfiguredForFulfillment) const readiness = diff --git a/apps/backend/src/startup/shutdown.test.ts b/apps/backend/src/startup/shutdown.test.ts new file mode 100644 index 00000000..e350079a --- /dev/null +++ b/apps/backend/src/startup/shutdown.test.ts @@ -0,0 +1,66 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import type { Server } from 'node:http' + +import { createShutdownController } from './shutdown.js' +import type { StartupState } from './state.js' + +function createStartupState(): StartupState { + return { + phase: 'ready', + core: { + running: false, + ready: true, + attemptCount: 1, + lastAttemptAt: '', + lastError: '', + readyAt: '', + }, + process: { + lastUnhandledRejection: null, + lastUncaughtException: null, + }, + } +} + +test('shutdown closes the HTTP server and database pool', async () => { + let serverClosed = false + let databaseClosed = false + const server = { + listening: true, + close(callback: (error?: Error) => void) { + serverClosed = true + callback() + }, + } as unknown as Server + const controller = createShutdownController(server, createStartupState(), { + closeDatabase: async () => { + databaseClosed = true + }, + }) + + await controller.shutdown('test') + + assert.equal(serverClosed, true) + assert.equal(databaseClosed, true) + assert.equal(controller.isShutdownStarted(), true) +}) + +test('shutdown closes the database pool when the server is already stopped', async () => { + let databaseClosed = false + const server = { + listening: false, + close() { + throw new Error('close should not be called') + }, + } as unknown as Server + const controller = createShutdownController(server, createStartupState(), { + closeDatabase: async () => { + databaseClosed = true + }, + }) + + await controller.shutdown('test') + + assert.equal(databaseClosed, true) +}) diff --git a/apps/backend/src/startup/shutdown.ts b/apps/backend/src/startup/shutdown.ts index 07254a02..3755a73e 100644 --- a/apps/backend/src/startup/shutdown.ts +++ b/apps/backend/src/startup/shutdown.ts @@ -2,10 +2,19 @@ import type { Server } from 'node:http' import { stopKuaishouIndustrySendCallbackRetryWorker } from '../services/platforms/kuaishou-industry/send-code-service.js' import { stopScheduledJobs } from '../services/scheduler/scheduler-service.js' +import { closeDb } from '../db/client.js' import { logError, logInfo } from '../utils/logger.js' import type { StartupState } from './state.js' -export function createShutdownController(server: Server, startupState: StartupState) { +type ShutdownOptions = { + closeDatabase?: () => Promise +} + +export function createShutdownController( + server: Server, + startupState: StartupState, + { closeDatabase = closeDb }: ShutdownOptions = {}, +) { let shutdownStarted = false async function shutdown(signal: string) { @@ -20,20 +29,25 @@ export function createShutdownController(server: Server, startupState: StartupSt stopScheduledJobs() stopKuaishouIndustrySendCallbackRetryWorker() - if (!server.listening) { - return + if (server.listening) { + await new Promise((resolve) => { + server.close((error) => { + if (error) { + logError('[shutdown]', 'failed to close HTTP server', error) + process.exitCode = 1 + } + + resolve() + }) + }) } - await new Promise((resolve) => { - server.close((error) => { - if (error) { - logError('[shutdown]', 'failed to close HTTP server', error) - process.exitCode = 1 - } - - resolve() - }) - }) + try { + await closeDatabase() + } catch (error) { + logError('[shutdown]', 'failed to close database pool', error) + process.exitCode = 1 + } } return {