import crypto from 'node:crypto' import { query } from '../db/client.js' import type { KuaishouIndustryVoucherUpdatePatch, KuaishouIndustryVoucherUpsertInput, } from '../types/repository/inputs.js' import type { KuaishouIndustryVoucherRow } from '../types/repository/rows.js' type VoucherPatchColumn = { column: string value: unknown cast?: string } export type KuaishouIndustryVoucherAdminListQuery = { oid?: string voucherCode?: string taskId?: number | string sellerId?: string status?: string page?: number | string pageSize?: number | string } const VOUCHER_CODE_PREFIX = 'KSV' const VOUCHER_CODE_RANDOM_BYTES = 10 const VOUCHER_CODE_MAX_GENERATE_ATTEMPTS = 5 const CROCKFORD_BASE32_ALPHABET = '0123456789ABCDEFGHJKMNPQRSTVWXYZ' export async function upsertKuaishouIndustryVoucher( input: KuaishouIndustryVoucherUpsertInput, ): Promise { for (let attempt = 1; attempt <= VOUCHER_CODE_MAX_GENERATE_ATTEMPTS; attempt += 1) { try { return await insertKuaishouIndustryVoucher(input, generateKuaishouIndustryVoucherCode()) } catch (error) { if (attempt < VOUCHER_CODE_MAX_GENERATE_ATTEMPTS && isVoucherCodeUniqueViolation(error)) { continue } throw error } } return null } export function generateKuaishouIndustryVoucherCode( randomBytes: (size: number) => Uint8Array = crypto.randomBytes, ): string { return `${VOUCHER_CODE_PREFIX}${encodeCrockfordBase32(randomBytes(VOUCHER_CODE_RANDOM_BYTES))}` } async function insertKuaishouIndustryVoucher( input: KuaishouIndustryVoucherUpsertInput, voucherCode: string, ): Promise { const result = await query( ` INSERT INTO kuaishou_industry_vouchers ( voucher_code, oid, order_id, task_id, unit_index, seller_id, token, eticket_type, status, valid_start_time, valid_end_time, consume_serial_num, consume_details_json, consumed_at, destroyed_at, send_callback_status, send_callback_attempt_count, send_callback_last_error, send_callback_response_json, send_callback_sent_at, raw_payload_json, created_at, updated_at ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13::jsonb, $14, $15, $16, $17, $18, $19::jsonb, $20, $21::jsonb, $22, $23 ) ON CONFLICT (oid, unit_index) DO UPDATE SET seller_id = CASE WHEN EXCLUDED.seller_id <> '' THEN EXCLUDED.seller_id ELSE kuaishou_industry_vouchers.seller_id END, token = CASE WHEN EXCLUDED.token <> '' THEN EXCLUDED.token ELSE kuaishou_industry_vouchers.token END, eticket_type = CASE WHEN EXCLUDED.eticket_type <> '' THEN EXCLUDED.eticket_type ELSE kuaishou_industry_vouchers.eticket_type END, order_id = COALESCE(kuaishou_industry_vouchers.order_id, EXCLUDED.order_id), task_id = COALESCE(kuaishou_industry_vouchers.task_id, EXCLUDED.task_id), valid_start_time = EXCLUDED.valid_start_time, valid_end_time = EXCLUDED.valid_end_time, send_callback_status = CASE WHEN kuaishou_industry_vouchers.send_callback_status = 'success' THEN kuaishou_industry_vouchers.send_callback_status ELSE EXCLUDED.send_callback_status END, send_callback_attempt_count = CASE WHEN kuaishou_industry_vouchers.send_callback_status = 'success' THEN kuaishou_industry_vouchers.send_callback_attempt_count ELSE EXCLUDED.send_callback_attempt_count END, send_callback_last_error = CASE WHEN kuaishou_industry_vouchers.send_callback_status = 'success' THEN kuaishou_industry_vouchers.send_callback_last_error ELSE EXCLUDED.send_callback_last_error END, send_callback_response_json = CASE WHEN kuaishou_industry_vouchers.send_callback_status = 'success' THEN kuaishou_industry_vouchers.send_callback_response_json ELSE EXCLUDED.send_callback_response_json END, send_callback_sent_at = CASE WHEN kuaishou_industry_vouchers.send_callback_status = 'success' THEN kuaishou_industry_vouchers.send_callback_sent_at ELSE EXCLUDED.send_callback_sent_at END, raw_payload_json = EXCLUDED.raw_payload_json, updated_at = EXCLUDED.updated_at RETURNING * `, [ voucherCode, input.oid, normalizeNullableId(input.orderId), normalizeNullableId(input.taskId), input.unitIndex, input.sellerId || '', input.token || '', input.eticketType || '', input.status || 'UNUSED', input.validStartTime || 0, input.validEndTime || 0, input.consumeSerialNum || '', stringifyJson(input.consumeDetailsJson ?? []), input.consumedAt || null, input.destroyedAt || null, input.sendCallbackStatus || 'success', normalizeNonNegativeInteger(input.sendCallbackAttemptCount, 0), input.sendCallbackLastError || '', stringifyJson(input.sendCallbackResponseJson ?? {}), input.sendCallbackSentAt || null, stringifyJson(input.rawPayloadJson ?? {}), input.createdAt, input.updatedAt, ], ) return result.rows[0] || null } export async function listKuaishouIndustryVouchersByOid( oid: string, ): Promise { const result = await query( ` SELECT * FROM kuaishou_industry_vouchers WHERE oid = $1 ORDER BY unit_index ASC, id ASC `, [oid], ) return result.rows } export type KuaishouIndustryVoucherWorkOrderSyncSource = { oid: string raw_payload_json: string | Record updated_at: string } export async function listKuaishouIndustryVoucherWorkOrderSyncSources( limit = 100, ): Promise { const result = await query( ` SELECT oid, raw_payload_json, updated_at FROM ( SELECT DISTINCT ON (oid) oid, raw_payload_json, updated_at, id FROM kuaishou_industry_vouchers WHERE raw_payload_json ? 'body' ORDER BY oid, updated_at DESC, id DESC ) AS latest_source ORDER BY updated_at DESC, id DESC LIMIT $1 `, [normalizePositiveLimit(limit, 100)], ) return result.rows } export async function listKuaishouIndustryVouchersForSendCallbackRetry( limit = 100, maxAttempts = 8, ): Promise { const result = await query( ` SELECT * FROM kuaishou_industry_vouchers WHERE send_callback_status <> 'success' AND COALESCE(send_callback_attempt_count, 0) < $2 ORDER BY COALESCE(send_callback_sent_at, created_at) ASC, id ASC LIMIT $1 `, [normalizePositiveLimit(limit, 100), Math.max(1, Math.trunc(Number(maxAttempts) || 8))], ) return result.rows } export async function listKuaishouIndustryVouchersByTaskId( taskId: number | string, ): Promise { const result = await query( ` SELECT * FROM kuaishou_industry_vouchers WHERE task_id = $1 ORDER BY unit_index ASC, id ASC `, [Number(taskId)], ) return result.rows } export type KuaishouIndustryVoucherAdminListRow = KuaishouIndustryVoucherRow & { sku_code?: string | null sku_name?: string | null task_no?: string | null } export async function listKuaishouIndustryVouchersForAdmin( input: KuaishouIndustryVoucherAdminListQuery = {}, ): Promise<{ items: KuaishouIndustryVoucherAdminListRow[] total: number page: number pageSize: number }> { const page = normalizePositiveLimit(input.page, 1) const pageSize = Math.min(normalizePositiveLimit(input.pageSize, 50), 200) const where: string[] = [] const params: unknown[] = [] const oid = String(input.oid || '').trim() if (oid) { params.push(oid) where.push(`v.oid = $${params.length}`) } const voucherCode = String(input.voucherCode || '').trim() if (voucherCode) { params.push(Array.from(new Set([voucherCode, voucherCode.toUpperCase()]))) where.push(`v.voucher_code = ANY($${params.length}::text[])`) } const taskId = Number(input.taskId || 0) if (Number.isFinite(taskId) && taskId > 0) { params.push(Math.trunc(taskId)) where.push(`v.task_id = $${params.length}`) } const sellerId = String(input.sellerId || '').trim() if (sellerId) { params.push(sellerId) where.push(`v.seller_id = $${params.length}`) } const status = String(input.status || '') .trim() .toUpperCase() if (status) { params.push(status) where.push(`v.status = $${params.length}`) } const whereSql = where.length > 0 ? `WHERE ${where.join(' AND ')}` : '' const totalResult = await query<{ total: string | number }>( ` SELECT COUNT(*) AS total FROM kuaishou_industry_vouchers v ${whereSql} `, params, ) const listParams = [...params, pageSize, (page - 1) * pageSize] // 关联任务/订单商品,便于运营后台展示「买的是什么」 const listResult = await query( ` SELECT v.*, COALESCE(oi_task.sku_code, oi_order.sku_code, '') AS sku_code, COALESCE(oi_task.sku_name, oi_order.sku_name, '') AS sku_name, COALESCE(ft.task_no, '') AS task_no FROM kuaishou_industry_vouchers v LEFT JOIN fulfillment_tasks ft ON ft.id = v.task_id LEFT JOIN order_items oi_task ON oi_task.id = ft.order_item_id LEFT JOIN LATERAL ( SELECT oi.sku_code, oi.sku_name FROM order_items oi WHERE v.order_id IS NOT NULL AND oi.order_id = v.order_id ORDER BY oi.id ASC LIMIT 1 ) oi_order ON true ${whereSql} ORDER BY v.updated_at DESC, v.id DESC LIMIT $${params.length + 1} OFFSET $${params.length + 2} `, listParams, ) return { items: listResult.rows, total: Number(totalResult.rows[0]?.total || 0) || 0, page, pageSize, } } export async function findKuaishouIndustryVoucherByCode( voucherCode: string, oid = '', ): Promise { const normalizedCode = String(voucherCode || '').trim() const normalizedOid = String(oid || '').trim() if (!normalizedCode) { return null } const candidateCodes = Array.from(new Set([normalizedCode, normalizedCode.toUpperCase()])) const params: unknown[] = [candidateCodes] const oidFilter = normalizedOid ? 'AND oid = $2' : '' if (normalizedOid) { params.push(normalizedOid) } const result = await query( ` SELECT * FROM kuaishou_industry_vouchers WHERE voucher_code = ANY($1::text[]) ${oidFilter} LIMIT 1 `, params, ) return result.rows[0] || null } export async function updateKuaishouIndustryVoucherByCode( voucherCode: string, patch: KuaishouIndustryVoucherUpdatePatch, ): Promise { const normalizedCode = String(voucherCode || '').trim() if (!normalizedCode) { return null } const columns = normalizeVoucherPatchColumns(patch) if (columns.length === 0) { return findKuaishouIndustryVoucherByCode(normalizedCode) } const assignments = columns .map((item, index) => `${item.column} = $${index + 1}${item.cast || ''}`) .join(', ') const params = columns.map((item) => item.value) params.push(normalizedCode) const result = await query( ` UPDATE kuaishou_industry_vouchers SET ${assignments} WHERE voucher_code = $${params.length} RETURNING * `, params, ) return result.rows[0] || null } function normalizeVoucherPatchColumns( patch: KuaishouIndustryVoucherUpdatePatch, ): VoucherPatchColumn[] { const columns: VoucherPatchColumn[] = [] if (patch.orderId !== undefined) { columns.push({ column: 'order_id', value: normalizeNullableId(patch.orderId) }) } if (patch.taskId !== undefined) { columns.push({ column: 'task_id', value: normalizeNullableId(patch.taskId) }) } if (patch.token !== undefined) { columns.push({ column: 'token', value: patch.token || '' }) } if (patch.sellerId !== undefined) { columns.push({ column: 'seller_id', value: patch.sellerId || '' }) } if (patch.eticketType !== undefined) { columns.push({ column: 'eticket_type', value: patch.eticketType || '' }) } if (patch.status !== undefined) { columns.push({ column: 'status', value: patch.status || 'UNUSED' }) } if (patch.validStartTime !== undefined) { columns.push({ column: 'valid_start_time', value: patch.validStartTime || 0 }) } if (patch.validEndTime !== undefined) { columns.push({ column: 'valid_end_time', value: patch.validEndTime || 0 }) } if (patch.consumeSerialNum !== undefined) { columns.push({ column: 'consume_serial_num', value: patch.consumeSerialNum || '' }) } if (patch.consumeDetailsJson !== undefined) { columns.push({ column: 'consume_details_json', value: stringifyJson(patch.consumeDetailsJson ?? []), cast: '::jsonb', }) } if (patch.consumedAt !== undefined) { columns.push({ column: 'consumed_at', value: patch.consumedAt || null }) } if (patch.destroyedAt !== undefined) { columns.push({ column: 'destroyed_at', value: patch.destroyedAt || null }) } if (patch.sendCallbackStatus !== undefined) { columns.push({ column: 'send_callback_status', value: patch.sendCallbackStatus || 'success' }) } if (patch.sendCallbackAttemptCount !== undefined) { columns.push({ column: 'send_callback_attempt_count', value: normalizeNonNegativeInteger(patch.sendCallbackAttemptCount, 0), }) } if (patch.sendCallbackLastError !== undefined) { columns.push({ column: 'send_callback_last_error', value: patch.sendCallbackLastError || '' }) } if (patch.sendCallbackResponseJson !== undefined) { columns.push({ column: 'send_callback_response_json', value: stringifyJson(patch.sendCallbackResponseJson ?? {}), cast: '::jsonb', }) } if (patch.sendCallbackSentAt !== undefined) { columns.push({ column: 'send_callback_sent_at', value: patch.sendCallbackSentAt || null }) } if (patch.rawPayloadJson !== undefined) { columns.push({ column: 'raw_payload_json', value: stringifyJson(patch.rawPayloadJson ?? {}), cast: '::jsonb', }) } if (patch.updatedAt !== undefined) { columns.push({ column: 'updated_at', value: patch.updatedAt }) } return columns } function encodeCrockfordBase32(bytes: Uint8Array): string { let bits = 0 let value = 0 let output = '' for (const byte of bytes) { value = (value << 8) | byte bits += 8 while (bits >= 5) { output += CROCKFORD_BASE32_ALPHABET[(value >>> (bits - 5)) & 31] || '' bits -= 5 } } if (bits > 0) { output += CROCKFORD_BASE32_ALPHABET[(value << (5 - bits)) & 31] || '' } return output } function isVoucherCodeUniqueViolation(error: unknown): boolean { const current = error && typeof error === 'object' ? (error as { code?: unknown; constraint?: unknown; detail?: unknown }) : {} if (String(current.code || '') !== '23505') { return false } const marker = `${String(current.constraint || '')} ${String(current.detail || '')}` return marker.includes('voucher_code') } function normalizeNullableId(value: unknown): number | null { if (value === null || value === undefined || value === '') { return null } const parsed = Number(value) return Number.isFinite(parsed) && parsed > 0 ? parsed : null } function normalizeNonNegativeInteger(value: unknown, fallback: number) { const parsed = Number(value) return Number.isInteger(parsed) && parsed >= 0 ? parsed : fallback } function normalizePositiveLimit(value: unknown, fallback: number): number { const parsed = Number(value) return Number.isFinite(parsed) && parsed > 0 ? Math.trunc(parsed) : fallback } function stringifyJson(value: unknown): string { if (typeof value === 'string') { return value } return JSON.stringify(value ?? {}) }