开始多平台 多店铺整改

This commit is contained in:
yml
2026-04-08 22:21:11 +08:00
parent d5b10e3f7f
commit a2b09c89e8
40 changed files with 1750 additions and 685 deletions
+59 -4
View File
@@ -7,13 +7,17 @@ import { fileURLToPath } from 'node:url'
const require = createRequire(import.meta.url)
const CURRENT_DIR = path.dirname(fileURLToPath(import.meta.url))
export const PROJECT_ROOT = path.resolve(CURRENT_DIR, '../..')
const WORKSPACE_ROOT = path.resolve(PROJECT_ROOT, '../..')
const CONFIG_ROOT = path.join(PROJECT_ROOT, 'config')
const defaultConfig = loadConfig(path.join(CONFIG_ROOT, 'default.cjs'))
const localConfig = loadConfig(path.join(CONFIG_ROOT, 'local.cjs'))
const mergedConfig = deepMerge(defaultConfig, localConfig)
loadEnvFiles([
path.join(WORKSPACE_ROOT, '.env'),
path.join(PROJECT_ROOT, '.env'),
])
export const runtimeConfig = applyEnvOverrides(mergedConfig)
const defaultConfig = loadConfig(path.join(CONFIG_ROOT, 'default.cjs'))
export const runtimeConfig = applyEnvOverrides(defaultConfig)
function loadConfig(configPath) {
if (!fs.existsSync(configPath)) {
@@ -24,6 +28,57 @@ function loadConfig(configPath) {
return isPlainObject(loaded) ? loaded : {}
}
function loadEnvFiles(filePaths) {
for (const filePath of filePaths) {
loadEnvFile(filePath)
}
}
function loadEnvFile(filePath) {
if (!fs.existsSync(filePath)) {
return
}
const rawText = fs.readFileSync(filePath, 'utf8')
const lines = rawText.split(/\r?\n/)
for (const rawLine of lines) {
const line = rawLine.trim()
if (!line || line.startsWith('#')) {
continue
}
const separatorIndex = line.indexOf('=')
if (separatorIndex <= 0) {
continue
}
const key = line.slice(0, separatorIndex).trim()
if (!key || key in process.env) {
continue
}
const value = parseEnvValue(line.slice(separatorIndex + 1))
process.env[key] = value
}
}
function parseEnvValue(rawValue) {
const value = String(rawValue || '').trim()
if (!value) {
return ''
}
const quote = value[0]
if ((quote === '"' || quote === "'") && value.endsWith(quote)) {
return value.slice(1, -1)
}
return value
}
function applyEnvOverrides(baseConfig) {
const nextConfig = deepMerge(baseConfig, {})
@@ -0,0 +1,94 @@
PRAGMA foreign_keys = OFF;
CREATE TABLE IF NOT EXISTS orders__new (
id INTEGER PRIMARY KEY AUTOINCREMENT,
provider TEXT NOT NULL DEFAULT 'agiso',
platform TEXT NOT NULL,
shop_id TEXT NOT NULL DEFAULT '',
shop_name TEXT NOT NULL DEFAULT '',
platform_order_id TEXT NOT NULL,
order_status TEXT NOT NULL DEFAULT 'created',
pay_status TEXT NOT NULL DEFAULT 'unpaid',
buyer_id TEXT NOT NULL DEFAULT '',
buyer_name TEXT NOT NULL DEFAULT '',
receiver_contact TEXT NOT NULL DEFAULT '',
total_amount INTEGER NOT NULL DEFAULT 0,
currency TEXT NOT NULL DEFAULT 'CNY',
raw_payload_json TEXT NOT NULL DEFAULT '{}',
paid_at TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
INSERT INTO orders__new (
id,
provider,
platform,
shop_id,
shop_name,
platform_order_id,
order_status,
pay_status,
buyer_id,
buyer_name,
receiver_contact,
total_amount,
currency,
raw_payload_json,
paid_at,
created_at,
updated_at
)
SELECT
id,
'agiso',
platform,
'',
'',
platform_order_id,
order_status,
pay_status,
buyer_id,
buyer_name,
receiver_contact,
total_amount,
currency,
raw_payload_json,
paid_at,
created_at,
updated_at
FROM orders;
DROP TABLE orders;
ALTER TABLE orders__new RENAME TO orders;
CREATE UNIQUE INDEX IF NOT EXISTS idx_orders_provider_platform_shop_order
ON orders(provider, platform, shop_id, platform_order_id);
ALTER TABLE webhook_events ADD COLUMN provider TEXT NOT NULL DEFAULT 'agiso';
ALTER TABLE webhook_events ADD COLUMN shop_id TEXT NOT NULL DEFAULT '';
ALTER TABLE webhook_events ADD COLUMN shop_name TEXT NOT NULL DEFAULT '';
UPDATE webhook_events
SET provider = CASE
WHEN trim(provider) = '' THEN 'agiso'
ELSE provider
END;
CREATE INDEX IF NOT EXISTS idx_webhook_events_provider_platform_shop
ON webhook_events(provider, platform, shop_id);
ALTER TABLE message_deliveries ADD COLUMN provider TEXT NOT NULL DEFAULT 'agiso';
ALTER TABLE message_deliveries ADD COLUMN shop_id TEXT NOT NULL DEFAULT '';
ALTER TABLE message_deliveries ADD COLUMN shop_name TEXT NOT NULL DEFAULT '';
UPDATE message_deliveries
SET provider = CASE
WHEN trim(provider) = '' THEN 'agiso'
ELSE provider
END;
CREATE INDEX IF NOT EXISTS idx_message_deliveries_provider_platform_shop
ON message_deliveries(provider, platform, shop_id);
PRAGMA foreign_keys = ON;
+14 -13
View File
@@ -9,6 +9,7 @@ import webhooksRouter from './routes/webhooks.js'
import { ensureAdminUsersBootstrapped } from './services/admin-auth-service.js'
import { closeLocalOcrWorker, warmupLocalOcrWorker } from './services/ocr.js'
import { closeAllTencentBrowserSessions, warmupTencentBrowser } from './services/session.js'
import { logError, logInfo, logWarn } from './utils/logger.js'
import tencentRouter from './routes/tencent.js'
const app = express()
@@ -47,7 +48,7 @@ app.use('/api/v1/claim', claimsRouter)
app.use('/api/v1/admin', adminRouter)
const server = app.listen(port, () => {
console.log(`order-site-backend listening on http://127.0.0.1:${port}`)
logInfo('[startup]', `order-site-backend listening on http://127.0.0.1:${port}`)
void bootstrapBrowser()
void bootstrapOcr()
})
@@ -57,34 +58,34 @@ async function bootstrapBrowser() {
const result = await warmupTencentBrowser()
if (!result.warmed) {
console.log('[startup] browser prewarm skipped')
logInfo('[startup]', 'browser prewarm skipped')
return
}
console.log('[startup] browser prewarm ready')
logInfo('[startup]', 'browser prewarm ready')
} catch (error) {
if (isMissingPlaywrightBrowserError(error)) {
const installCommand = resolveBrowserInstallCommand()
console.warn(`[startup] browser prewarm skipped: ${formatErrorMessage(error)}`)
console.warn(`[startup] 请先安装浏览器依赖: ${installCommand}`)
logWarn('[startup]', `browser prewarm skipped: ${formatErrorMessage(error)}`)
logWarn('[startup]', `请先安装浏览器依赖: ${installCommand}`)
if (String(runtimeConfig.browser.chromePath || '').trim()) {
console.warn('[startup] 当前设置了 CHROME_PATH,也请确认该路径指向的浏览器可执行文件真实存在')
logWarn('[startup]', '当前设置了 CHROME_PATH,也请确认该路径指向的浏览器可执行文件真实存在')
}
return
}
console.error('[startup] browser prewarm failed:', error)
logError('[startup]', 'browser prewarm failed', error)
}
}
async function bootstrapOcr() {
try {
await warmupLocalOcrWorker()
console.log('[startup] local OCR worker ready')
logInfo('[startup]', 'local OCR ready')
} catch (error) {
console.warn(`[startup] local OCR worker skipped: ${formatStartupError(error)}`)
logWarn('[startup]', `local OCR skipped: ${formatStartupError(error)}`)
}
}
@@ -131,18 +132,18 @@ async function shutdown(signal) {
}
shutdownStarted = true
console.log(`[shutdown] received ${signal}, closing browser sessions and HTTP server`)
logInfo('[shutdown]', `received ${signal}, closing browser sessions and HTTP server`)
try {
await closeAllTencentBrowserSessions({ markClosed: true })
} catch (error) {
console.error('[shutdown] failed to close browser sessions:', error)
logError('[shutdown]', 'failed to close browser sessions', error)
}
try {
await closeLocalOcrWorker()
} catch (error) {
console.error('[shutdown] failed to close local OCR worker:', error)
logError('[shutdown]', 'failed to close local OCR worker', error)
}
if (!server.listening) {
@@ -152,7 +153,7 @@ async function shutdown(signal) {
await new Promise((resolve) => {
server.close((error) => {
if (error) {
console.error('[shutdown] failed to close HTTP server:', error)
logError('[shutdown]', 'failed to close HTTP server', error)
process.exitCode = 1
}
@@ -3,7 +3,10 @@ import { getDb } from '../db/client.js'
export function createMessageDelivery(input) {
const result = getDb().prepare(`
INSERT INTO message_deliveries (
provider,
platform,
shop_id,
shop_name,
channel,
order_id,
task_id,
@@ -21,9 +24,12 @@ export function createMessageDelivery(input) {
sent_at,
created_at,
updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`).run(
input.provider,
input.platform,
input.shopId,
input.shopName,
input.channel,
input.orderId,
input.taskId,
+18 -4
View File
@@ -1,18 +1,21 @@
import { getDb } from '../db/client.js'
export function findOrderByPlatformOrderId(platform, platformOrderId) {
export function findOrderByPlatformOrderId({ provider = 'agiso', platform, shopId = '', platformOrderId }) {
return getDb().prepare(`
SELECT *
FROM orders
WHERE platform = ? AND platform_order_id = ?
WHERE provider = ? AND platform = ? AND shop_id = ? AND platform_order_id = ?
LIMIT 1
`).get(platform, platformOrderId) || null
`).get(provider, platform, shopId, platformOrderId) || null
}
export function createOrder(input) {
const result = getDb().prepare(`
INSERT INTO orders (
provider,
platform,
shop_id,
shop_name,
platform_order_id,
order_status,
pay_status,
@@ -25,9 +28,12 @@ export function createOrder(input) {
paid_at,
created_at,
updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`).run(
input.provider,
input.platform,
input.shopId,
input.shopName,
input.platformOrderId,
input.orderStatus,
input.payStatus,
@@ -49,6 +55,10 @@ export function updateOrder(orderId, input) {
getDb().prepare(`
UPDATE orders
SET
provider = ?,
platform = ?,
shop_id = ?,
shop_name = ?,
order_status = ?,
pay_status = ?,
buyer_id = ?,
@@ -61,6 +71,10 @@ export function updateOrder(orderId, input) {
updated_at = ?
WHERE id = ?
`).run(
input.provider,
input.platform,
input.shopId,
input.shopName,
input.orderStatus,
input.payStatus,
input.buyerId,
@@ -3,7 +3,10 @@ import { getDb } from '../db/client.js'
export function createWebhookEvent(input) {
const result = getDb().prepare(`
INSERT INTO webhook_events (
provider,
platform,
shop_id,
shop_name,
event_type,
event_key,
signature_valid,
@@ -14,9 +17,12 @@ export function createWebhookEvent(input) {
process_error,
related_order_id,
created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`).run(
input.provider,
input.platform,
input.shopId,
input.shopName,
input.eventType,
input.eventKey,
input.signatureValid ? 1 : 0,
@@ -60,6 +66,7 @@ export function getWebhookEventById(eventId) {
export function listWebhookEvents({
page = 1,
pageSize = 20,
provider = '',
platform = '',
processed = '',
relatedOrderId = '',
@@ -70,6 +77,11 @@ export function listWebhookEvents({
const filters = []
const params = []
if (provider) {
filters.push('provider = ?')
params.push(provider)
}
if (platform) {
filters.push('platform = ?')
params.push(platform)
+32
View File
@@ -1,19 +1,51 @@
import { Router } from 'express'
import { processAgisoTradeWebhook } from '../services/webhook-service.js'
import { createRequestId, logWebhook } from '../utils/logger.js'
import { buildNotFoundPayload, buildSuccessPayload, sendRouteError } from '../utils/http.js'
const router = Router()
router.post('/agiso/trade', async (req, res) => {
const requestId = createRequestId('wh')
const startedAt = Date.now()
logWebhook('[webhooks/agiso/trade]', '收到 Agiso webhook 请求', {
requestId,
method: req.method,
originalUrl: req.originalUrl,
ip: req.ip,
headers: req.headers,
query: req.query,
body: req.body,
})
try {
const data = await processAgisoTradeWebhook({
requestId,
headers: req.headers,
query: req.query,
body: req.body,
})
logWebhook('[webhooks/agiso/trade]', 'Agiso webhook 处理完成', {
requestId,
durationMs: Date.now() - startedAt,
result: data,
})
res.json(buildSuccessPayload(data, 'success'))
} catch (error) {
logWebhook(
'[webhooks/agiso/trade]',
'Agiso webhook 处理失败',
{
requestId,
durationMs: Date.now() - startedAt,
error,
},
{ level: 'error' },
)
sendRouteError(res, error, '处理 Agiso 订单通知失败', '[webhooks/agiso/trade]')
}
})
+21 -1
View File
@@ -70,7 +70,10 @@ export function getAdminOrderDetail(orderId) {
return {
order: {
orderId: order.id,
provider: order.provider || 'agiso',
platform: order.platform,
shopId: order.shop_id || '',
shopName: order.shop_name || '',
platformOrderId: order.platform_order_id,
orderStatus: order.order_status,
payStatus: order.pay_status,
@@ -96,7 +99,10 @@ export function getAdminOrderDetail(orderId) {
tasks: tasks.map(mapAdminTaskSummary),
webhookEvents: webhookEvents.map((event) => ({
eventId: event.id,
provider: event.provider || 'agiso',
platform: event.platform,
shopId: event.shop_id || '',
shopName: event.shop_name || '',
eventType: event.event_type,
eventKey: event.event_key,
signatureValid: Boolean(event.signature_valid),
@@ -154,6 +160,10 @@ export function getAdminTaskDetail(taskId) {
order: order
? {
orderId: order.id,
provider: order.provider || 'agiso',
platform: order.platform,
shopId: order.shop_id || '',
shopName: order.shop_name || '',
platformOrderId: order.platform_order_id,
payStatus: order.pay_status,
orderStatus: order.order_status,
@@ -337,6 +347,7 @@ export function getAdminWebhookEvents(query = {}) {
const { items, total } = listWebhookEvents({
page,
pageSize,
provider: String(query.provider || '').trim(),
platform: String(query.platform || '').trim(),
processed: String(query.processed || '').trim(),
relatedOrderId: String(query.relatedOrderId || '').trim(),
@@ -347,7 +358,10 @@ export function getAdminWebhookEvents(query = {}) {
return {
items: items.map((item) => ({
eventId: item.id,
provider: item.provider || 'agiso',
platform: item.platform,
shopId: item.shop_id || '',
shopName: item.shop_name || '',
eventType: item.event_type,
eventKey: item.event_key,
signatureValid: Boolean(item.signature_valid),
@@ -372,7 +386,10 @@ export function getAdminWebhookEventDetail(eventId) {
return {
eventId: event.id,
provider: event.provider || 'agiso',
platform: event.platform,
shopId: event.shop_id || '',
shopName: event.shop_name || '',
eventType: event.event_type,
eventKey: event.event_key,
signatureValid: Boolean(event.signature_valid),
@@ -396,7 +413,7 @@ export async function replayAdminWebhookEvent(eventId) {
})
}
if (event.platform !== 'agiso') {
if (String(event.provider || event.platform || '').trim() !== 'agiso') {
throw createHttpError('当前只支持重放 agiso webhook', {
statusCode: 409,
errorCode: 'admin_webhook_replay_not_supported',
@@ -686,7 +703,10 @@ function mapAdminOrderListItem(item) {
return {
orderId: item.id,
provider: item.provider || 'agiso',
platform: item.platform,
shopId: item.shop_id || '',
shopName: item.shop_name || '',
platformOrderId: item.platform_order_id,
orderStatus: item.order_status,
payStatus: item.pay_status,
+4 -1
View File
@@ -44,7 +44,10 @@ export async function ensureAgisoClaimMessageDeliveredForTask({ order, task, cla
})
const createdAt = nowIso()
const delivery = createMessageDelivery({
platform: 'agiso',
provider: 'agiso',
platform: String(order.platform || '').trim() || 'unknown',
shopId: String(order.shop_id || '').trim(),
shopName: String(order.shop_name || '').trim(),
channel: AGISO_MESSAGE_CHANNEL,
orderId: order.id,
taskId: task.id,
+126 -142
View File
@@ -1,193 +1,177 @@
import { spawn } from 'node:child_process'
import readline from 'node:readline'
import fs from 'node:fs'
import path from 'node:path'
import process from 'node:process'
import { runtimeConfig } from '../config/runtime.js'
const DEFAULT_TIMEOUT_MS = 60_000
let workerPromise = null
let requestSequence = 0
const pendingRequests = new Map()
let stderrBuffer = []
export async function recognizeTencentCaptcha(payload) {
return callOcrWorker('recognize', payload)
return callOcrRequest('recognize', payload)
}
export async function batchRecognizeTencentCaptcha(payload) {
return callOcrWorker('batch', payload, { timeoutMs: 180_000 })
return callOcrRequest('batch', payload, { timeoutMs: 180_000 })
}
export async function warmupLocalOcrWorker() {
await ensureOcrWorker()
await runOcrCommand('healthcheck', { timeoutMs: 30_000 })
}
export async function closeLocalOcrWorker() {
const worker = await workerPromise?.catch(() => null)
if (!worker) {
workerPromise = null
return
}
worker.child.kill()
workerPromise = null
// OCR 改为按次调用 Python 子进程,这里不再维护常驻 worker。
}
async function callOcrWorker(action, payload, { timeoutMs = DEFAULT_TIMEOUT_MS } = {}) {
const worker = await ensureOcrWorker()
const requestId = `ocr-${Date.now()}-${++requestSequence}`
async function callOcrRequest(action, payload, { timeoutMs = DEFAULT_TIMEOUT_MS } = {}) {
const response = await runOcrCommand('request', {
timeoutMs,
input: JSON.stringify({
action,
payload: payload || {},
}),
})
if (!response || typeof response !== 'object' || Array.isArray(response)) {
throw new Error('本地 OCR 返回了无效响应')
}
return response
}
async function runOcrCommand(command, { timeoutMs = DEFAULT_TIMEOUT_MS, input = '' } = {}) {
const projectRoot = resolveOcrProjectRoot()
const workerCommand = resolveWorkerCommand(projectRoot)
const env = {
...process.env,
PYTHONUNBUFFERED: '1',
PYTHONPATH: buildPythonPath(path.join(projectRoot, 'src')),
}
return await new Promise((resolve, reject) => {
const timer = setTimeout(() => {
pendingRequests.delete(requestId)
worker.child.kill()
reject(new Error('本地 OCR 识别超时,请确认 uv 环境和 OCR worker 依赖正常'))
}, timeoutMs)
pendingRequests.set(requestId, {
resolve: (value) => {
clearTimeout(timer)
resolve(value)
},
reject: (error) => {
clearTimeout(timer)
reject(error)
},
})
worker.child.stdin.write(
`${JSON.stringify({ id: requestId, action, payload: payload || {} })}\n`,
'utf8',
)
})
}
async function ensureOcrWorker() {
if (!workerPromise) {
workerPromise = createOcrWorker().catch((error) => {
workerPromise = null
throw error
})
}
return workerPromise
}
function createOcrWorker() {
return new Promise((resolve, reject) => {
const projectRoot = resolveOcrProjectRoot()
stderrBuffer = []
const child = spawn('uv', ['run', 'ocr-worker', 'worker'], {
const child = spawn(workerCommand.command, [...workerCommand.args, command], {
cwd: projectRoot,
stdio: ['pipe', 'pipe', 'pipe'],
env: process.env,
env,
})
child.stdin.setDefaultEncoding('utf8')
const stdoutReader = readline.createInterface({ input: child.stdout })
const worker = { child, stdoutReader }
let stdout = ''
let stderr = ''
let settled = false
let timedOut = false
const readyTimer = setTimeout(() => {
reject(new Error(`本地 OCR worker 启动超时,请检查 ${projectRoot} 下是否已执行 uv sync`))
const timer = setTimeout(() => {
timedOut = true
child.kill()
}, 15_000)
}, timeoutMs)
let ready = false
stdoutReader.on('line', (line) => {
const message = tryParseJson(line)
if (!message || typeof message !== 'object') {
return
}
if (message.type === 'ready') {
if (!ready) {
ready = true
clearTimeout(readyTimer)
resolve(worker)
}
return
}
if (message.type !== 'response') {
return
}
const requestId = String(message.id || '')
const pending = pendingRequests.get(requestId)
if (!pending) {
return
}
pendingRequests.delete(requestId)
pending.resolve(message.payload)
child.stdout.on('data', (chunk) => {
stdout += String(chunk || '')
})
child.stderr.on('data', (chunk) => {
const text = String(chunk || '').trim()
if (!text) {
return
}
stderrBuffer.push(text)
if (stderrBuffer.length > 20) {
stderrBuffer = stderrBuffer.slice(-20)
}
stderr += String(chunk || '')
})
child.on('error', (error) => {
clearTimeout(readyTimer)
const workerError = buildWorkerError(error)
failAllPendingRequests(workerError)
if (!ready) {
reject(workerError)
if (settled) {
return
}
workerPromise = null
settled = true
clearTimeout(timer)
reject(buildWorkerError(error))
})
child.on('exit', (code, signal) => {
clearTimeout(readyTimer)
const workerError = new Error(
buildWorkerExitMessage({
code,
signal,
projectRoot,
stderrText: stderrBuffer.join('\n'),
}),
)
failAllPendingRequests(workerError)
if (!ready) {
reject(workerError)
child.on('close', (code, signal) => {
if (settled) {
return
}
workerPromise = null
settled = true
clearTimeout(timer)
if (timedOut) {
reject(new Error(`本地 OCR 执行超时,请检查 ${projectRoot} 的 Python 依赖是否已经准备完成`))
return
}
if (signal || code !== 0) {
reject(
new Error(
buildWorkerExitMessage({
code,
signal,
projectRoot,
stderrText: stderr.trim(),
}),
),
)
return
}
if (command === 'healthcheck') {
resolve(null)
return
}
const parsed = tryParseJson(stdout.trim())
if (!parsed || typeof parsed !== 'object') {
reject(new Error(`本地 OCR 返回无法解析的结果${stderr.trim() ? `\n${stderr.trim()}` : ''}`))
return
}
resolve(parsed)
})
if (input) {
child.stdin.end(`${input}\n`, 'utf8')
return
}
child.stdin.end()
})
}
function failAllPendingRequests(error) {
for (const [requestId, pending] of pendingRequests.entries()) {
pendingRequests.delete(requestId)
pending.reject(error)
}
}
function resolveOcrProjectRoot() {
return String(runtimeConfig.ocr.projectRoot || '').trim()
}
function resolveWorkerCommand(projectRoot) {
const venvPython = path.join(projectRoot, '.venv', 'bin', 'python')
if (fs.existsSync(venvPython)) {
return {
command: venvPython,
args: ['-u', '-m', 'ocr_worker.cli'],
}
}
return {
command: 'python3',
args: ['-u', '-m', 'ocr_worker.cli'],
}
}
function buildPythonPath(srcPath) {
const current = String(process.env.PYTHONPATH || '').trim()
return current ? `${srcPath}${path.delimiter}${current}` : srcPath
}
function buildWorkerError(error) {
if (error instanceof Error && error.message.includes('spawn uv ENOENT')) {
return new Error('未找到 uv 命令,请先安装 uv,并确认它在 PATH 里')
if (
error instanceof Error &&
(error.message.includes('spawn python3 ENOENT') ||
error.message.includes('spawn ') && error.message.includes('python'))
) {
const projectRoot = resolveOcrProjectRoot()
return new Error(
`无法启动本地 OCR。请检查 Python 是否可用,并确认 OCR_PROJECT_ROOT 指向的目录存在:${projectRoot}`,
)
}
return new Error(
`本地 OCR worker 启动失败: ${error instanceof Error ? error.message : String(error || '未知错误')}`,
`本地 OCR 启动失败: ${error instanceof Error ? error.message : String(error || '未知错误')}`,
)
}
@@ -198,10 +182,10 @@ function buildWorkerExitMessage({ code, signal, projectRoot, stderrText }) {
? `exit code ${code}`
: 'unknown reason'
const installHint = `先执行: cd ${projectRoot} && uv sync`
const installHint = `检查 ${projectRoot} 的 Python 依赖是否已安装;本地源码运行请先执行 uv syncDocker 请重建 backend 镜像`
const stderrHint = stderrText ? `\n${stderrText}` : ''
return `本地 OCR worker 已退出(${reason})。${installHint}${stderrHint}`
return `本地 OCR 已退出(${reason})。${installHint}${stderrHint}`
}
function tryParseJson(text) {
+33 -3
View File
@@ -5,12 +5,31 @@ import { buildClaimUrl } from './claim-service.js'
import { syncDeliveryTasksForOrder } from './delivery-task-service.js'
import { ensureAgisoClaimMessageDeliveredForTask } from './message-service.js'
import { nowIso } from '../utils/time.js'
import { logWebhook } from '../utils/logger.js'
export async function upsertOrderFromWebhook(event) {
const now = nowIso()
const existing = findOrderByPlatformOrderId(event.platform, event.platformOrderId)
const basePayload = {
const existing = findOrderByPlatformOrderId({
provider: event.provider,
platform: event.platform,
shopId: event.shopId,
platformOrderId: event.platformOrderId,
})
logWebhook('[order-service]', '开始处理 webhook 订单 upsert', {
provider: event.provider,
platform: event.platform,
shopId: event.shopId,
shopName: event.shopName,
platformOrderId: event.platformOrderId,
existingOrderId: existing?.id || null,
})
const basePayload = {
provider: event.provider,
platform: event.platform,
shopId: event.shopId,
shopName: event.shopName,
platformOrderId: event.platformOrderId,
orderStatus: event.orderStatus,
payStatus: event.payStatus,
@@ -50,8 +69,19 @@ export async function upsertOrderFromWebhook(event) {
const tasks = syncDeliveryTasksForOrder(order, orderItems)
const messageDeliveries = []
logWebhook('[order-service]', 'Webhook 订单 upsert 完成', {
orderId: order.id,
provider: order.provider,
platform: order.platform,
shopId: order.shop_id,
shopName: order.shop_name,
platformOrderId: order.platform_order_id,
orderItemCount: orderItems.length,
taskCount: tasks.length,
})
for (const task of tasks) {
if (event.platform !== 'agiso' || String(task.task_status || '') !== 'link_generated' || !task.claim_token_id) {
if (event.provider !== 'agiso' || String(task.task_status || '') !== 'link_generated' || !task.claim_token_id) {
continue
}
@@ -1,4 +1,5 @@
import { runtimeConfig } from '../config/runtime.js'
import { logDebug } from '../utils/logger.js'
const SESSION_DEBUG_ENABLED = Boolean(runtimeConfig.session.debug)
@@ -7,12 +8,7 @@ export function logBrowserSessionDebug(scope, detail) {
return
}
if (typeof detail === 'undefined') {
console.log(`[browser/session][debug] ${scope}`)
return
}
console.log(`[browser/session][debug] ${scope}:`, detail)
logDebug('[browser/session]', scope, detail)
}
export async function waitForFrame(page, predicate, timeoutMs = 15_000) {
+387 -26
View File
@@ -4,29 +4,59 @@ import { runtimeConfig } from '../config/runtime.js'
import { createWebhookEvent, updateWebhookEvent } from '../repositories/webhook-event-repo.js'
import { upsertOrderFromWebhook } from './order-service.js'
import { createHttpError } from '../utils/http.js'
import { logWebhook } from '../utils/logger.js'
import { nowIso } from '../utils/time.js'
export async function processAgisoTradeWebhook(requestLike) {
const parsed = parseAgisoTradeRequest(requestLike)
const webhookEvent = createWebhookEvent(buildWebhookEventInput(requestLike, parsed))
return executeAgisoTradeWebhook(parsed, webhookEvent.id)
logWebhook('[webhook-service/agiso]', 'Webhook 已解析并写入 webhook_events', {
requestId: requestLike.requestId || '',
webhookEventId: webhookEvent.id,
provider: parsed.provider,
platform: parsed.platform,
shopId: parsed.shopId,
shopName: parsed.shopName,
eventKey: parsed.eventKey,
eventType: parsed.eventType,
platformOrderId: parsed.platformOrderId,
signatureValid: parsed.signatureValid,
})
return executeAgisoTradeWebhook(parsed, webhookEvent.id, { requestId: requestLike.requestId || '' })
}
export async function replayAgisoTradeWebhookEvent(webhookEvent) {
const requestLike = {
requestId: `replay-${webhookEvent.id}`,
headers: safeParseJson(webhookEvent.headers_json),
query: safeParseJson(webhookEvent.query_json),
body: safeParseJson(webhookEvent.body_json),
}
const parsed = parseAgisoTradeRequest(requestLike)
return executeAgisoTradeWebhook(parsed, webhookEvent.id)
logWebhook('[webhook-service/agiso]', '开始重放 webhook 事件', {
requestId: requestLike.requestId,
webhookEventId: webhookEvent.id,
provider: parsed.provider,
platform: parsed.platform,
shopId: parsed.shopId,
shopName: parsed.shopName,
eventKey: parsed.eventKey,
eventType: parsed.eventType,
platformOrderId: parsed.platformOrderId,
})
return executeAgisoTradeWebhook(parsed, webhookEvent.id, { requestId: requestLike.requestId })
}
function buildWebhookEventInput(requestLike, parsed) {
return {
platform: 'agiso',
provider: parsed.provider,
platform: parsed.platform,
shopId: parsed.shopId,
shopName: parsed.shopName,
eventType: parsed.eventType,
eventKey: parsed.eventKey,
signatureValid: parsed.signatureValid,
@@ -40,7 +70,7 @@ function buildWebhookEventInput(requestLike, parsed) {
}
}
async function executeAgisoTradeWebhook(parsed, webhookEventId) {
async function executeAgisoTradeWebhook(parsed, webhookEventId, { requestId = '' } = {}) {
try {
if (!parsed.signatureValid) {
throw createHttpError('验签失败', {
@@ -63,8 +93,28 @@ async function executeAgisoTradeWebhook(parsed, webhookEventId) {
related_order_id: result.order.id,
})
logWebhook('[webhook-service/agiso]', 'Webhook 业务处理成功', {
requestId,
webhookEventId,
provider: parsed.provider,
platform: parsed.platform,
shopId: parsed.shopId,
shopName: parsed.shopName,
eventKey: parsed.eventKey,
eventType: parsed.eventType,
platformOrderId: parsed.platformOrderId,
orderId: result.order.id,
taskCount: result.tasks.length,
messageDeliveryCount: Array.isArray(result.messageDeliveries) ? result.messageDeliveries.length : 0,
})
return {
accepted: true,
provider: parsed.provider,
platform: parsed.platform,
shopId: parsed.shopId,
shopName: parsed.shopName,
platformRaw: parsed.platformRaw,
eventType: parsed.eventType,
platformOrderId: parsed.platformOrderId,
orderId: result.order.id,
@@ -82,6 +132,25 @@ async function executeAgisoTradeWebhook(parsed, webhookEventId) {
process_error: error instanceof Error ? error.message : String(error || ''),
related_order_id: null,
})
logWebhook(
'[webhook-service/agiso]',
'Webhook 业务处理失败',
{
requestId,
webhookEventId,
provider: parsed.provider,
platform: parsed.platform,
shopId: parsed.shopId,
shopName: parsed.shopName,
eventKey: parsed.eventKey,
eventType: parsed.eventType,
platformOrderId: parsed.platformOrderId,
signatureValid: parsed.signatureValid,
error,
},
{ level: 'error' },
)
throw error
}
}
@@ -93,34 +162,83 @@ function parseAgisoTradeRequest(requestLike) {
const payload = rawJson ? parsePayloadJson(rawJson) : body
const timestamp = String(query.timestamp || '').trim()
const sign = String(query.sign || '').trim().toLowerCase()
const eventType = resolveEventType(query.aopic)
const provider = 'agiso'
const platformRaw = pickFirstNonEmpty([
query.fromPlatform,
query.from_platform,
body.fromPlatform,
body.from_platform,
payload.fromPlatform,
payload.FromPlatform,
payload.platform,
payload.Platform,
])
const platform = resolveBusinessPlatform(platformRaw)
const shop = resolveShop(payload)
const eventType = resolveEventType(query.aopic, payload)
const signatureValid = verifyAgisoSignature({ rawJson, timestamp, sign })
const itemSources = extractOrderItemSources(payload)
const firstItem = itemSources[0] || {}
const platformOrderId = pickFirstNonEmpty([
payload.biz_order_id,
payload.Tid,
payload.tid,
payload.Oid,
payload.oid,
payload.order_id,
payload.orderId,
firstItem.Oid,
firstItem.oid,
])
return {
platform: 'agiso',
provider,
platform,
platformRaw,
eventType,
eventKey: `${platformOrderId || 'unknown'}:${eventType}:${timestamp || 'na'}`,
eventKey: [provider, platform || 'unknown', shop.shopId || 'unknown', platformOrderId || 'unknown', eventType, timestamp || 'na'].join(':'),
signatureValid,
platformOrderId,
shopId: shop.shopId,
shopName: shop.shopName,
orderStatus: resolveOrderStatus(eventType, payload),
payStatus: resolvePayStatus(eventType, payload),
buyerId: pickFirstNonEmpty([payload.buyer_id, payload.buyerId, payload.openid]),
buyerName: pickFirstNonEmpty([payload.buyer_name, payload.buyerName, payload.nick]),
buyerId: pickFirstNonEmpty([
payload.buyer_id,
payload.buyerId,
payload.BuyerId,
payload.openid,
payload.BuyerOpenUid,
payload.buyer_open_uid,
payload.buyerOpenUid,
]),
buyerName: pickFirstNonEmpty([
payload.buyer_name,
payload.buyerName,
payload.BuyerName,
payload.nick,
payload.BuyerNick,
payload.buyer_nick,
payload.buyerNick,
]),
receiverContact: pickFirstNonEmpty([
payload.receiver_contact,
payload.receiverContact,
payload.mobile,
payload.phone,
payload.receiver_mobile,
payload.receiverMobile,
]),
totalAmount: normalizeInteger(
pickFirstNonEmpty([payload.total_fee, payload.totalFee, payload.pay_fee, payload.payFee]),
pickFirstNonEmpty([
payload.total_fee,
payload.totalFee,
payload.TotalFee,
payload.pay_fee,
payload.payFee,
payload.Payment,
payload.payment,
]),
),
currency: pickFirstNonEmpty([payload.currency, 'CNY']) || 'CNY',
paidAt: resolvePaidAt(eventType, payload),
@@ -159,13 +277,17 @@ function verifyAgisoSignature({ rawJson, timestamp, sign }) {
return legacyDigest === sign
}
function resolveEventType(aopic) {
function resolveEventType(aopic, payload) {
const normalized = String(aopic || '').trim()
if (normalized === '1') {
return 'payment_success'
}
if (normalized === '21') {
return 'trade_create'
}
if (normalized === '32') {
return 'trade_create'
}
@@ -178,10 +300,29 @@ function resolveEventType(aopic) {
return 'trade_closed'
}
const status = resolveProviderStatus(payload)
if (status === 'closed') {
return 'trade_closed'
}
return 'trade_event'
}
function resolveOrderStatus(eventType, payload) {
const providerStatus = resolveProviderStatus(payload)
if (providerStatus === 'paid') {
return 'paid'
}
if (providerStatus === 'closed') {
return 'closed'
}
if (providerStatus === 'refunded') {
return 'refunded'
}
if (eventType === 'payment_success') {
return 'paid'
}
@@ -210,6 +351,20 @@ function resolveOrderStatus(eventType, payload) {
}
function resolvePayStatus(eventType, payload) {
const providerStatus = resolveProviderStatus(payload)
if (providerStatus === 'paid') {
return 'paid'
}
if (providerStatus === 'closed') {
return 'unpaid'
}
if (providerStatus === 'refunded') {
return 'refunded'
}
if (eventType === 'payment_success') {
return 'paid'
}
@@ -234,6 +389,24 @@ function resolvePayStatus(eventType, payload) {
}
function resolvePaidAt(eventType, payload) {
const providerPaidAt = normalizeProviderDateTime(
pickFirstNonEmpty([
payload.paid_at,
payload.paidAt,
payload.pay_time,
payload.payTime,
payload.PayTime,
]),
)
if (providerPaidAt) {
return providerPaidAt
}
if (resolveProviderStatus(payload) === 'paid') {
return nowIso()
}
if (eventType === 'payment_success') {
return nowIso()
}
@@ -266,13 +439,17 @@ function parsePayloadJson(rawJson) {
}
function normalizeOrderItems(payload) {
const items = Array.isArray(payload.items) && payload.items.length > 0
? payload.items
: [payload]
const items = extractOrderItemSources(payload)
return items.map((item) => {
const source = isPlainObject(item) ? item : {}
const rawSkuKey = pickFirstNonEmpty([
const rawSkuCandidates = [
source.OuterSkuId,
source.outerSkuId,
source.outer_sku_id,
source.OuterIid,
source.outerIid,
source.outer_iid,
source.sku_code,
source.skuCode,
source.goods_sku,
@@ -280,13 +457,18 @@ function normalizeOrderItems(payload) {
source.itemId,
source.goods_id,
source.goodsId,
source.NumIid,
source.num_iid,
])
const skuCode = resolveSkuCode(rawSkuKey)
payload.item_id,
payload.itemId,
]
const skuCode = resolveSkuCode(rawSkuCandidates)
return {
skuCode,
skuName: pickFirstNonEmpty([
source.Title,
source.title,
source.sku_name,
source.skuName,
source.goods_name,
@@ -295,25 +477,204 @@ function normalizeOrderItems(payload) {
]),
quantity: Math.max(
1,
normalizeInteger(pickFirstNonEmpty([source.quantity, source.num, source.buy_amount])) || 1,
normalizeInteger(
pickFirstNonEmpty([
source.quantity,
source.num,
source.Num,
source.buy_amount,
payload.quantity,
payload.num,
payload.Num,
]),
) || 1,
),
spec: source,
}
})
}
function resolveSkuCode(rawKey) {
const mappings = runtimeConfig.orders.skuMappings
function extractOrderItemSources(payload) {
const candidates = [
payload.items,
payload.Items,
payload.orders,
payload.Orders,
payload.order_list,
payload.OrderList,
]
if (rawKey && mappings && typeof mappings === 'object') {
const direct = mappings[String(rawKey)]
if (typeof direct === 'string' && direct.trim()) {
return direct.trim()
for (const current of candidates) {
if (Array.isArray(current) && current.length > 0) {
return current
}
}
if (typeof rawKey === 'string' && rawKey.trim()) {
return rawKey.trim()
return [payload]
}
function resolveBusinessPlatform(rawPlatform) {
const normalized = normalizePlatformKey(rawPlatform)
if (!normalized) {
return 'unknown'
}
const aliases = {
xianyu: 'xianyu',
idlefish: 'xianyu',
idle: 'xianyu',
taobaoidle: 'xianyu',
taobaoxianyu: 'xianyu',
taobao: 'taobao',
tb: 'taobao',
pdd: 'pdd',
pinduoduo: 'pdd',
douyin: 'douyin',
weidian: 'weidian',
jd: 'jd',
jingdong: 'jd',
}
return aliases[normalized] || normalized
}
function normalizePlatformKey(value) {
return String(value || '')
.trim()
.toLowerCase()
.replace(/[\s_-]+/g, '')
}
function resolveShop(payload) {
const shopName = pickFirstNonEmpty([
payload.shop_name,
payload.shopName,
payload.ShopName,
payload.seller_name,
payload.sellerName,
payload.seller_nick,
payload.sellerNick,
payload.SellerNick,
])
return {
shopId: pickFirstNonEmpty([
payload.shop_id,
payload.shopId,
payload.ShopId,
payload.seller_id,
payload.sellerId,
payload.SellerId,
payload.PlatformUserId,
payload.platformUserId,
shopName,
]),
shopName,
}
}
function resolveProviderStatus(payload) {
const rawStatus = normalizeInteger(
pickFirstNonEmpty([payload.order_status, payload.orderStatus, payload.status]),
)
if (rawStatus === 3 || rawStatus === 4) {
return 'paid'
}
if (rawStatus === 6) {
return 'closed'
}
if (rawStatus === 5) {
return 'refunded'
}
const statusText = pickFirstNonEmpty([
payload.Status,
payload.status,
payload.order_status_text,
payload.orderStatusText,
]).toUpperCase()
if (!statusText) {
return pickFirstNonEmpty([payload.PayTime, payload.payTime, payload.pay_time]) ? 'paid' : 'unknown'
}
if (['WAIT_SELLER_SEND_GOODS', 'WAIT_BUYER_CONFIRM_GOODS', 'TRADE_FINISHED', 'SUCCESS', 'PAID'].includes(statusText)) {
return 'paid'
}
if (['TRADE_CLOSED', 'TRADE_CLOSED_BY_TAOBAO', 'CLOSED'].includes(statusText)) {
return 'closed'
}
if (['REFUND_SUCCESS', 'TRADE_REFUND', 'REFUNDED'].includes(statusText)) {
return 'refunded'
}
if (['WAIT_BUYER_PAY', 'CREATED', 'NEW'].includes(statusText)) {
return 'unpaid'
}
return pickFirstNonEmpty([payload.PayTime, payload.payTime, payload.pay_time]) ? 'paid' : 'unknown'
}
function normalizeProviderDateTime(value) {
const normalized = String(value || '').trim()
if (!normalized) {
return null
}
const isoText = normalized.replace(' ', 'T')
const utcMatch = isoText.match(/^(\d{4})-(\d{2})-(\d{2})T(\d{2}):(\d{2}):(\d{2})$/)
if (utcMatch) {
const [, year, month, day, hours, minutes, seconds] = utcMatch
return new Date(
Date.UTC(
Number(year),
Number(month) - 1,
Number(day),
Number(hours) - 8,
Number(minutes),
Number(seconds),
),
).toISOString()
}
const date = new Date(normalized)
return Number.isNaN(date.getTime()) ? null : date.toISOString()
}
function resolveSkuCode(rawKey) {
const mappings = runtimeConfig.orders.skuMappings
const candidates = Array.isArray(rawKey) ? rawKey : [rawKey]
if (mappings && typeof mappings === 'object') {
for (const candidate of candidates) {
const normalizedCandidate = String(candidate || '').trim()
if (!normalizedCandidate) {
continue
}
const direct = mappings[normalizedCandidate]
if (typeof direct === 'string' && direct.trim()) {
return direct.trim()
}
}
}
for (const candidate of candidates) {
if (typeof candidate === 'string' && candidate.trim()) {
return candidate.trim()
}
if (typeof candidate === 'number' && Number.isFinite(candidate)) {
return String(candidate)
}
}
return ''
+27 -9
View File
@@ -1,3 +1,7 @@
import { logError, logWarn } from './logger.js'
const RESPONSE_TIME_ZONE = 'Asia/Shanghai'
export function buildSuccessPayload(data, msg = 'ok') {
return {
code: 0,
@@ -21,8 +25,12 @@ export function buildErrorPayload(error, fallbackMessage) {
export function sendRouteError(res, error, fallbackMessage, scope) {
const classification = classifyRouteError(error)
const logger = classification.statusCode >= 500 ? console.error : console.warn
logger(`${scope} ${fallbackMessage}:`, error)
const logger = classification.statusCode >= 500 ? logError : logWarn
logger(scope, fallbackMessage, {
statusCode: classification.statusCode,
errorCode: classification.errorCode,
error,
})
res.status(classification.statusCode).json(buildErrorPayload(error, fallbackMessage))
}
@@ -150,12 +158,22 @@ function formatResponseDateTime(value) {
return String(value || '').replace('T', ' ').replace('Z', '')
}
const year = date.getFullYear()
const month = String(date.getMonth() + 1).padStart(2, '0')
const day = String(date.getDate()).padStart(2, '0')
const hours = String(date.getHours()).padStart(2, '0')
const minutes = String(date.getMinutes()).padStart(2, '0')
const seconds = String(date.getSeconds()).padStart(2, '0')
const parts = new Intl.DateTimeFormat('zh-CN', {
timeZone: RESPONSE_TIME_ZONE,
hour12: false,
year: 'numeric',
month: '2-digit',
day: '2-digit',
hour: '2-digit',
minute: '2-digit',
second: '2-digit',
}).formatToParts(date)
return `${year}-${month}-${day} ${hours}:${minutes}:${seconds}`
const tokens = Object.fromEntries(
parts
.filter((item) => item.type !== 'literal')
.map((item) => [item.type, item.value]),
)
return `${tokens.year}-${tokens.month}-${tokens.day} ${tokens.hour}:${tokens.minute}:${tokens.second}`
}
+134
View File
@@ -0,0 +1,134 @@
import { existsSync } from 'node:fs'
import fs from 'node:fs/promises'
import path from 'node:path'
import process from 'node:process'
import { PROJECT_ROOT, runtimeConfig } from '../config/runtime.js'
const DATA_ROOT = resolveDataRoot()
const LOG_DIR = path.join(DATA_ROOT, 'logs')
const APP_LOG_FILE = path.join(LOG_DIR, 'app.log')
const WEBHOOK_LOG_FILE = path.join(LOG_DIR, 'webhook.log')
let writeQueue = Promise.resolve()
export function createRequestId(prefix = 'req') {
return `${prefix}-${Date.now().toString(36)}-${Math.random().toString(36).slice(2, 8)}`
}
export function logDebug(scope, message, detail) {
return writeLog('debug', scope, message, detail)
}
export function logInfo(scope, message, detail) {
return writeLog('info', scope, message, detail)
}
export function logWarn(scope, message, detail) {
return writeLog('warn', scope, message, detail)
}
export function logError(scope, message, detail) {
return writeLog('error', scope, message, detail)
}
export function logWebhook(scope, message, detail, { level = 'info' } = {}) {
return writeLog(level, scope, message, detail, { channel: 'webhook' })
}
function writeLog(level, scope, message, detail, { channel = 'app' } = {}) {
const entry = {
time: new Date().toISOString(),
level,
scope: String(scope || 'app').trim() || 'app',
message: String(message || '').trim() || '-',
pid: process.pid,
detail: normalizeLogValue(detail),
}
writeConsole(entry)
enqueueFileWrite(entry, channel === 'webhook' ? WEBHOOK_LOG_FILE : APP_LOG_FILE)
return entry
}
function enqueueFileWrite(entry, filePath) {
const line = JSON.stringify(entry)
writeQueue = writeQueue
.then(async () => {
await fs.mkdir(LOG_DIR, { recursive: true })
await fs.appendFile(filePath, `${line}\n`, 'utf8')
})
.catch((error) => {
const reason = error instanceof Error ? error.message : String(error || '未知错误')
console.error('[logger] failed to write log file:', reason)
})
}
function writeConsole(entry) {
const prefix = `[${entry.time}] [${entry.level}] [${entry.scope}] ${entry.message}`
const logger = resolveConsoleMethod(entry.level)
if (typeof entry.detail === 'undefined') {
logger(prefix)
return
}
logger(prefix, entry.detail)
}
function resolveConsoleMethod(level) {
if (level === 'error') {
return console.error
}
if (level === 'warn') {
return console.warn
}
return console.log
}
function resolveDataRoot() {
const configuredDatabasePath = String(runtimeConfig.database?.filePath || '').trim()
if (configuredDatabasePath) {
const configuredDataRoot = path.dirname(configuredDatabasePath)
if (existsSync(configuredDataRoot)) {
return configuredDataRoot
}
}
return path.resolve(PROJECT_ROOT, 'data')
}
function normalizeLogValue(value) {
if (typeof value === 'undefined') {
return undefined
}
try {
return JSON.parse(
JSON.stringify(value, (_key, current) => {
if (current instanceof Error) {
return {
name: current.name,
message: current.message,
stack: current.stack,
statusCode: current.statusCode,
errorCode: current.errorCode,
}
}
if (typeof current === 'bigint') {
return String(current)
}
return current
}),
)
} catch {
return String(value)
}
}