优化订单履约分流与后台筛选

This commit is contained in:
yml2213
2026-09-02 11:59:38 +08:00
parent 50c526a855
commit de63daf1ac
19 changed files with 1269 additions and 162 deletions
@@ -0,0 +1,29 @@
CREATE EXTENSION IF NOT EXISTS pg_trgm;
CREATE INDEX IF NOT EXISTS idx_orders_admin_search_text_trgm
ON orders USING GIN ((
platform_order_id || ' ' || shop_id || ' ' || shop_name || ' ' ||
buyer_id || ' ' || buyer_name || ' ' || receiver_contact
) gin_trgm_ops);
CREATE INDEX IF NOT EXISTS idx_order_items_admin_search_text_trgm
ON order_items USING GIN ((
sku_code || ' ' || sku_name || ' ' ||
COALESCE(item_snapshot_json->'kuaishouSendCode'->>'itemTitle', '') || ' ' ||
COALESCE(item_snapshot_json->'kuaishouSendCode'->>'skuNick', '')
) gin_trgm_ops);
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_status_order
ON fulfillment_tasks(task_status, order_id);
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_executor_order
ON fulfillment_tasks(executor_key, order_id);
COMMENT ON INDEX idx_orders_admin_search_text_trgm IS
'支持后台订单号、店铺、买家及联系方式综合模糊搜索';
COMMENT ON INDEX idx_order_items_admin_search_text_trgm IS
'支持后台订单商品名、SKU、快手商品标题和 SKU 标题综合模糊搜索';
COMMENT ON INDEX idx_fulfillment_tasks_status_order IS
'支持后台按任务状态筛选订单';
COMMENT ON INDEX idx_fulfillment_tasks_executor_order IS
'支持后台按履约渠道筛选订单';
@@ -0,0 +1,7 @@
CREATE EXTENSION IF NOT EXISTS pg_trgm;
CREATE INDEX IF NOT EXISTS idx_fulfillment_tasks_admin_search_text_trgm
ON fulfillment_tasks USING GIN ((task_no || ' ' || platform_order_id) gin_trgm_ops);
COMMENT ON INDEX idx_fulfillment_tasks_admin_search_text_trgm IS
'支持后台任务号和平台订单号综合模糊搜索';
@@ -21,6 +21,22 @@ export async function listOrderItemsByOrderId(orderId: number | string): Promise
return listOrderItemsByOrderIdWithExecutor(query, orderId)
}
export async function listOrderItemsByOrderIds(orderIds: number[]): Promise<OrderItemRow[]> {
const normalizedIds = orderIds.map(Number).filter((value) => value > 0)
if (normalizedIds.length === 0) return []
const result = await query<OrderItemRow>(
`
SELECT *
FROM order_items
WHERE order_id = ANY($1::bigint[])
ORDER BY order_id ASC, id ASC
`,
[normalizedIds],
)
return result.rows
}
export async function replaceOrderItems(
orderId: number | string,
items: OrderItemReplaceInput[],
@@ -0,0 +1,73 @@
import assert from 'node:assert/strict'
import test from 'node:test'
import { listOrdersWithExecutor, type OrderQueryExecutor } from './order-repo.js'
test('订单列表综合筛选下推到总数和分页查询', async () => {
const calls: Array<{ sql: string; params: unknown[] }> = []
const executor = createExecutor(calls, 3)
const result = await listOrdersWithExecutor(executor, {
page: 2,
pageSize: 10,
keyword: '皮肤店',
provider: '91kaquan',
executorKey: 'kuaishou_ct_assisted',
payStatus: 'paid',
taskStatus: 'manual_review',
dateFrom: '2026-09-01T00:00:00.000Z',
dateTo: '2026-09-02T23:59:59.999Z',
})
assert.equal(result.total, 3)
assert.equal(calls.length, 2)
assert.match(calls[0]?.sql || '', /o\.platform_order_id.*ILIKE \$1/s)
assert.match(calls[0]?.sql || '', /oi\.sku_code.*oi\.sku_name.*ILIKE \$1/s)
assert.match(calls[0]?.sql || '', /kuaishouSendCode'.*itemTitle.*kuaishouSendCode'.*skuNick/s)
assert.match(calls[0]?.sql || '', /o\.provider = \$2/)
assert.match(calls[0]?.sql || '', /o\.pay_status = \$3/)
assert.match(calls[0]?.sql || '', /ft\.task_status = \$4/)
assert.match(calls[0]?.sql || '', /ft\.executor_key = \$5/)
assert.match(calls[0]?.sql || '', /o\.created_at >= \$6/)
assert.match(calls[0]?.sql || '', /o\.created_at <= \$7/)
assert.deepEqual(calls[0]?.params, [
'%皮肤店%',
'91kaquan',
'paid',
'manual_review',
'kuaishou_ct_assisted',
'2026-09-01T00:00:00.000Z',
'2026-09-02T23:59:59.999Z',
])
assert.match(calls[1]?.sql || '', /LIMIT \$8 OFFSET \$9/)
assert.deepEqual(calls[1]?.params, [...(calls[0]?.params || []), 10, 10])
})
test('订单列表保留订单号和精确 SKU 旧筛选参数兼容', async () => {
const calls: Array<{ sql: string; params: unknown[] }> = []
await listOrdersWithExecutor(createExecutor(calls, 0), {
platformOrderId: '26245',
skuCode: 'SKU-1',
})
assert.match(calls[0]?.sql || '', /o\.platform_order_id ILIKE \$1/)
assert.match(calls[0]?.sql || '', /oi\.sku_code = \$2/)
assert.deepEqual(calls[0]?.params, ['%26245%', 'SKU-1'])
})
function createExecutor(
calls: Array<{ sql: string; params: unknown[] }>,
total: number,
): OrderQueryExecutor {
return (async (sql: string, params: unknown[] = []) => {
calls.push({ sql, params: [...params] })
return {
command: 'SELECT',
rowCount: calls.length === 1 ? 1 : 0,
oid: 0,
fields: [],
rows: calls.length === 1 ? [{ total }] : [],
}
}) as OrderQueryExecutor
}
+63 -12
View File
@@ -202,18 +202,57 @@ export async function getOrderById(orderId: number | string): Promise<OrderRow |
}
export async function listOrders({
page = 1,
pageSize = 20,
platformOrderId = '',
payStatus = '',
skuCode = '',
dateFrom = '',
dateTo = '',
...input
}: OrderListQueryInput = {}): Promise<OrderListQueryResult> {
return listOrdersWithExecutor(query, input)
}
export async function listOrdersWithExecutor(
executor: OrderQueryExecutor,
{
page = 1,
pageSize = 20,
keyword = '',
provider = '',
executorKey = '',
platformOrderId = '',
payStatus = '',
taskStatus = '',
skuCode = '',
dateFrom = '',
dateTo = '',
}: OrderListQueryInput = {},
): Promise<OrderListQueryResult> {
const offset = (page - 1) * pageSize
const filters: string[] = []
const params: unknown[] = []
if (keyword) {
params.push(`%${keyword}%`)
const keywordParam = `$${params.length}`
filters.push(`(
(
o.platform_order_id || ' ' || o.shop_id || ' ' || o.shop_name || ' ' ||
o.buyer_id || ' ' || o.buyer_name || ' ' || o.receiver_contact
) ILIKE ${keywordParam}
OR EXISTS (
SELECT 1
FROM order_items oi
WHERE oi.order_id = o.id
AND (
oi.sku_code || ' ' || oi.sku_name || ' ' ||
COALESCE(oi.item_snapshot_json->'kuaishouSendCode'->>'itemTitle', '') || ' ' ||
COALESCE(oi.item_snapshot_json->'kuaishouSendCode'->>'skuNick', '')
) ILIKE ${keywordParam}
)
)`)
}
if (provider) {
params.push(provider)
filters.push(`o.provider = $${params.length}`)
}
if (platformOrderId) {
params.push(`%${platformOrderId}%`)
filters.push(`o.platform_order_id ILIKE $${params.length}`)
@@ -224,6 +263,20 @@ export async function listOrders({
filters.push(`o.pay_status = $${params.length}`)
}
if (taskStatus) {
params.push(taskStatus)
filters.push(
`EXISTS (SELECT 1 FROM fulfillment_tasks ft WHERE ft.order_id = o.id AND ft.task_status = $${params.length})`,
)
}
if (executorKey) {
params.push(executorKey)
filters.push(
`EXISTS (SELECT 1 FROM fulfillment_tasks ft WHERE ft.order_id = o.id AND ft.executor_key = $${params.length})`,
)
}
if (skuCode) {
params.push(skuCode)
filters.push(
@@ -242,18 +295,16 @@ export async function listOrders({
}
const whereClause = filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : ''
const totalResult = await query<{ [column: string]: unknown; total: number }>(
const totalResult = await executor<{ [column: string]: unknown; total: number }>(
`SELECT COUNT(*)::int AS total FROM orders o ${whereClause}`,
params,
)
params.push(pageSize)
params.push(offset)
const itemsResult = await query<OrderListRow>(
const itemsResult = await executor<OrderListRow>(
`
SELECT
o.*,
(SELECT COUNT(*)::int FROM fulfillment_tasks ft WHERE ft.order_id = o.id) AS task_count
SELECT o.*
FROM orders o
${whereClause}
ORDER BY o.id DESC
@@ -0,0 +1,48 @@
import assert from 'node:assert/strict'
import test from 'node:test'
import { buildTaskListFilters } from './task-repo.js'
test('任务列表综合关键词覆盖任务订单商品与角色字段', () => {
const result = buildTaskListFilters({
keyword: '指挥官密钥',
executorKey: 'kuaishou_ct_assisted',
status: 'manual_review',
dateFrom: '2026-09-01T00:00:00.000Z',
dateTo: '2026-09-02T23:59:59.999Z',
})
assert.match(result.whereClause, /ft\.task_no.*ft\.platform_order_id.*ILIKE \$1/s)
assert.match(result.whereClause, /oi\.sku_code ILIKE \$1/)
assert.match(result.whereClause, /oi\.sku_name ILIKE \$1/)
assert.match(result.whereClause, /kuaishouSendCode'.*itemTitle.*ILIKE \$1/s)
assert.match(result.whereClause, /kuaishouSendCode'.*skuNick.*ILIKE \$1/s)
assert.match(result.whereClause, /kcts\.role_id.*ctx\.role_id.*ILIKE \$1/s)
assert.match(result.whereClause, /kcts\.role_name.*ctx\.role_name.*ctx\.nickname.*ILIKE \$1/s)
assert.match(result.whereClause, /ft\.executor_key = \$2/)
assert.match(result.whereClause, /ft\.task_status = \$3/)
assert.match(result.whereClause, /ft\.created_at >= \$4/)
assert.match(result.whereClause, /ft\.created_at <= \$5/)
assert.deepEqual(result.params, [
'%指挥官密钥%',
'kuaishou_ct_assisted',
'manual_review',
'2026-09-01T00:00:00.000Z',
'2026-09-02T23:59:59.999Z',
])
})
test('任务列表保留旧版分字段筛选参数', () => {
const result = buildTaskListFilters({
platformOrderId: '26245',
taskNo: 'DT-1',
skuCode: 'SKU-1',
roleId: '10001',
})
assert.match(result.whereClause, /ft\.platform_order_id ILIKE \$1/)
assert.match(result.whereClause, /ft\.task_no ILIKE \$2/)
assert.match(result.whereClause, /oi\.sku_code = \$3/)
assert.match(result.whereClause, /kcts\.role_id.*ctx\.role_id.*ILIKE \$4/s)
assert.deepEqual(result.params, ['%26245%', '%DT-1%', 'SKU-1', '%10001%'])
})
+72 -28
View File
@@ -79,6 +79,19 @@ export async function listTasksByOrderId(orderId: number | string): Promise<Task
return result.rows
}
export async function listTasksByOrderIds(orderIds: number[]): Promise<TaskRow[]> {
const normalizedIds = orderIds.map(Number).filter((value) => value > 0)
if (normalizedIds.length === 0) return []
const result = await query<TaskRow>(
`${buildTaskSelect()}
WHERE ft.order_id = ANY($1::bigint[])
ORDER BY ft.order_id ASC, ft.id ASC`,
[normalizedIds],
)
return result.rows
}
export async function createTask(input: TaskCreateInput): Promise<TaskRow | null> {
return withTransaction(async (client) => {
const taskResult = await client.query<{ id: number }>(
@@ -418,6 +431,43 @@ export async function findAffiliateDashTaskByOrder(input: {
export async function listTasks({
page = 1,
pageSize = 20,
...filters
}: TaskListQueryInput = {}): Promise<TaskListQueryResult> {
const offset = (page - 1) * pageSize
const { whereClause, params } = buildTaskListFilters(filters)
const totalResult = await query<{ [column: string]: unknown; total: number }>(
`
SELECT COUNT(*)::int AS total
FROM fulfillment_tasks ft
LEFT JOIN order_items oi ON oi.id = ft.order_item_id
LEFT JOIN task_runtime_contexts ctx ON ctx.task_id = ft.id
LEFT JOIN kuaishou_cloud_task_states kcts ON kcts.task_id = ft.id
${whereClause}
`,
params,
)
params.push(pageSize)
params.push(offset)
const itemsResult = await query<TaskRow>(
`${buildTaskSelect('oi.sku_code, oi.sku_name, oi.quantity')}
LEFT JOIN order_items oi ON oi.id = ft.order_item_id
${whereClause}
ORDER BY ft.id DESC
LIMIT $${params.length - 1} OFFSET $${params.length}`,
params,
)
return {
items: itemsResult.rows,
total: Number(totalResult.rows[0]?.total || 0),
}
}
export function buildTaskListFilters({
keyword = '',
executorKey = '',
status = '',
platformOrderId = '',
taskNo = '',
@@ -425,11 +475,29 @@ export async function listTasks({
roleId = '',
dateFrom = '',
dateTo = '',
}: TaskListQueryInput = {}): Promise<TaskListQueryResult> {
const offset = (page - 1) * pageSize
}: TaskListQueryInput = {}): { whereClause: string; params: unknown[] } {
const filters: string[] = []
const params: unknown[] = []
if (keyword) {
params.push(`%${keyword}%`)
const keywordParam = `$${params.length}`
filters.push(`(
(ft.task_no || ' ' || ft.platform_order_id) ILIKE ${keywordParam}
OR oi.sku_code ILIKE ${keywordParam}
OR oi.sku_name ILIKE ${keywordParam}
OR COALESCE(oi.item_snapshot_json->'kuaishouSendCode'->>'itemTitle', '') ILIKE ${keywordParam}
OR COALESCE(oi.item_snapshot_json->'kuaishouSendCode'->>'skuNick', '') ILIKE ${keywordParam}
OR COALESCE(NULLIF(kcts.role_id, ''), ctx.role_id, '') ILIKE ${keywordParam}
OR COALESCE(NULLIF(kcts.role_name, ''), NULLIF(ctx.role_name, ''), ctx.nickname, '') ILIKE ${keywordParam}
)`)
}
if (executorKey) {
params.push(executorKey)
filters.push(`ft.executor_key = $${params.length}`)
}
if (status) {
params.push(status)
filters.push(`ft.task_status = $${params.length}`)
@@ -465,33 +533,9 @@ export async function listTasks({
filters.push(`ft.created_at <= $${params.length}`)
}
const whereClause = filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : ''
const totalResult = await query<{ [column: string]: unknown; total: number }>(
`
SELECT COUNT(*)::int AS total
FROM fulfillment_tasks ft
LEFT JOIN order_items oi ON oi.id = ft.order_item_id
LEFT JOIN task_runtime_contexts ctx ON ctx.task_id = ft.id
LEFT JOIN kuaishou_cloud_task_states kcts ON kcts.task_id = ft.id
${whereClause}
`,
params,
)
params.push(pageSize)
params.push(offset)
const itemsResult = await query<TaskRow>(
`${buildTaskSelect('oi.sku_code, oi.sku_name, oi.quantity')}
LEFT JOIN order_items oi ON oi.id = ft.order_item_id
${whereClause}
ORDER BY ft.id DESC
LIMIT $${params.length - 1} OFFSET $${params.length}`,
params,
)
return {
items: itemsResult.rows,
total: Number(totalResult.rows[0]?.total || 0),
whereClause: filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : '',
params,
}
}
@@ -1,17 +1,14 @@
import { listOrderItemsByOrderId } from '../../repositories/order-item-repo.js'
import { listTasksByOrderId } from '../../repositories/task-repo.js'
import { formatFenToAmount, normalizeFen } from '../../utils/money.js'
import { resolveDisplayShopName, resolveOrderItemTitle } from './admin-read-shared-helpers.js'
import { buildOrderFulfillmentProgress } from './admin-task-read-helpers.js'
import type { AdminOrderListItem } from '../../types/admin/read-models.js'
import type { OrderItemRow, OrderListRow } from '../../types/repository/rows.js'
import type { OrderItemRow, OrderListRow, TaskRow } from '../../types/repository/rows.js'
export async function mapAdminOrderListItem(item: OrderListRow): Promise<AdminOrderListItem> {
const [tasks, orderItems] = await Promise.all([
listTasksByOrderId(item.id),
listOrderItemsByOrderId(item.id),
])
export function mapAdminOrderListItem(
item: OrderListRow,
{ tasks = [], orderItems = [] }: { tasks?: TaskRow[]; orderItems?: OrderItemRow[] } = {},
): AdminOrderListItem {
const fulfillmentProgress = buildOrderFulfillmentProgress(tasks)
const itemSummary = summarizeOrderItems(orderItems)
@@ -39,7 +36,7 @@ export async function mapAdminOrderListItem(item: OrderListRow): Promise<AdminOr
0,
),
itemSummary,
taskCount: Number(item.task_count || tasks.length || 0),
taskCount: tasks.length,
resourceStatus: fulfillmentProgress.resourceStatus,
customerStatus: fulfillmentProgress.customerStatus,
preparedTaskCount: fulfillmentProgress.preparedTaskCount,
@@ -1,9 +1,17 @@
import { buildClaimUrl } from '../claim/claim-service.js'
import { buildClaimIdentityAdminSummary } from '../claim/claim-identity.js'
import { getClaimTokenById } from '../../repositories/claim-token-repo.js'
import { listOrderItemsByOrderId } from '../../repositories/order-item-repo.js'
import {
listOrderItemsByOrderId,
listOrderItemsByOrderIds,
} from '../../repositories/order-item-repo.js'
import { getOrderById, listOrders } from '../../repositories/order-repo.js'
import { getTaskById, listTasks, listTasksByOrderId } from '../../repositories/task-repo.js'
import {
getTaskById,
listTasks,
listTasksByOrderId,
listTasksByOrderIds,
} from '../../repositories/task-repo.js'
import { listTaskEventsByTaskId } from '../../repositories/task-event-repo.js'
import type { JsonRecord } from '../../types/json.js'
import {
@@ -76,19 +84,46 @@ export async function getAdminOrders(
const { items, total } = await listOrders({
page,
pageSize,
keyword: String(query.keyword || '').trim(),
provider: String(query.provider || '').trim(),
executorKey: String(query.executorKey || '').trim(),
platformOrderId: String(query.platformOrderId || '').trim(),
payStatus: String(query.payStatus || '').trim(),
taskStatus: String(query.taskStatus || '').trim(),
skuCode: String(query.skuCode || '').trim(),
dateFrom: normalizeDateQuery(query.dateFrom),
dateTo: normalizeDateQuery(query.dateTo, true),
})
const orderIds = items.map((item) => item.id)
const [orderItems, tasks] = await Promise.all([
listOrderItemsByOrderIds(orderIds),
listTasksByOrderIds(orderIds),
])
const orderItemsByOrderId = groupRowsByOrderId(orderItems)
const tasksByOrderId = groupRowsByOrderId(tasks)
return {
items: await Promise.all(items.map((item) => mapAdminOrderListItem(item))),
items: items.map((item) =>
mapAdminOrderListItem(item, {
orderItems: orderItemsByOrderId.get(item.id) || [],
tasks: tasksByOrderId.get(item.id) || [],
}),
),
pagination: { page, pageSize, total },
}
}
function groupRowsByOrderId<T extends { order_id: number }>(rows: T[]): Map<number, T[]> {
const grouped = new Map<number, T[]>()
for (const row of rows) {
const group = grouped.get(row.order_id) || []
group.push(row)
grouped.set(row.order_id, group)
}
return grouped
}
export async function getAdminOrderDetail(
orderId: AdminEntityIdInput,
session: AdminViewerSessionInput | null = null,
@@ -196,6 +231,8 @@ export async function getAdminTasks(
const { items, total } = await listTasks({
page,
pageSize,
keyword: String(query.keyword || '').trim(),
executorKey: String(query.executorKey || '').trim(),
status: String(query.status || '').trim(),
platformOrderId: String(query.platformOrderId || '').trim(),
taskNo: String(query.taskNo || '').trim(),
@@ -6,7 +6,9 @@ import type { OrderItemRow, OrderRow, TaskRow } from '../../types/repository/row
import type { JsonObject } from '../../types/json.js'
type CreatedTaskInput = {
orderId?: number
orderItemId?: number
unitIndex?: number
contextJson?: string
taskStatus?: string
executorKey?: string
@@ -59,9 +61,9 @@ function stubOrderItem(
function stubCreatedTask(id: number, input: CreatedTaskInput): TaskRow {
return {
id,
order_id: 0,
order_id: Number(input.orderId || 0),
order_item_id: Number(input.orderItemId || 0),
unit_index: 1,
unit_index: Number(input.unitIndex || 1),
platform_order_id: '',
profile_id: 0,
task_no: 'DT-test',
@@ -145,6 +147,7 @@ test('syncDeliveryTasksForOrder 多账号候选时不提前固定 cloudtentacles
},
nowIso: () => '2026-05-29T09:05:46.658Z',
randomId: () => 'DT-test',
createTaskEvent: async () => null,
},
)
@@ -194,6 +197,7 @@ test('syncDeliveryTasksForOrder 单账号候选时保留固定 cloudtentacles
},
nowIso: () => '2026-05-29T09:05:46.658Z',
randomId: () => 'DT-test',
createTaskEvent: async () => null,
},
)
@@ -203,3 +207,229 @@ test('syncDeliveryTasksForOrder 单账号候选时保留固定 cloudtentacles
assert.deepEqual(binding.cloudSourceKeys, ['account-a'])
assert.equal(binding.resolvedSourceKey, 'account-a')
})
test('syncDeliveryTasksForOrder 已有部分任务时按订单项补建缺失任务', async () => {
const existingTask = stubCreatedTask(201, {
orderId: 3,
orderItemId: 20,
unitIndex: 1,
executorKey: 'kuaishou_ct_assisted',
})
const createdInputs: CreatedTaskInput[] = []
const taskEvents: Array<{ taskId: number; eventType: string }> = []
const tasks = await syncDeliveryTasksForOrderWithDeps(
stubOrder({ id: 3, platform_order_id: 'partial-order' }),
[
stubOrderItem({
id: 20,
order_id: 3,
sku_code: 'SKU-existing',
item_snapshot_json: buildCloudSnapshot(74, '商品 A'),
}),
stubOrderItem({
id: 21,
order_id: 3,
sku_code: 'SKU-missing',
item_snapshot_json: buildCloudSnapshot(75, '商品 B'),
}),
],
{
listTasksByOrderId: async () => [existingTask],
getFulfillmentProfileByKey: async () => buildCloudProfile(),
createTask: async (input) => {
createdInputs.push(input as CreatedTaskInput)
return stubCreatedTask(202, input as CreatedTaskInput)
},
createTaskEvent: async (taskId, eventType) => {
taskEvents.push({ taskId: Number(taskId), eventType })
return null
},
logInfo: () => null,
logWarn: () => null,
nowIso: () => '2026-05-29T09:05:46.658Z',
randomId: () => 'DT-partial',
},
)
assert.deepEqual(
createdInputs.map((input) => [input.orderItemId, input.unitIndex]),
[[21, 1]],
)
assert.deepEqual(
tasks.map((task) => task.id),
[201, 202],
)
assert.deepEqual(taskEvents, [{ taskId: 202, eventType: 'fulfillment_task_planned' }])
})
test('syncDeliveryTasksForOrder 数量增加时只补建缺失单元', async () => {
const existingTask = stubCreatedTask(301, {
orderId: 4,
orderItemId: 30,
unitIndex: 1,
executorKey: 'kuaishou_ct_assisted',
})
const createdInputs: CreatedTaskInput[] = []
const tasks = await syncDeliveryTasksForOrderWithDeps(
stubOrder({ id: 4, platform_order_id: 'quantity-increased' }),
[
stubOrderItem({
id: 30,
order_id: 4,
sku_code: 'SKU-quantity',
quantity: 2,
item_snapshot_json: buildCloudSnapshot(74, '商品 A'),
}),
],
{
listTasksByOrderId: async () => [existingTask],
getFulfillmentProfileByKey: async () => buildCloudProfile(),
createTask: async (input) => {
createdInputs.push(input as CreatedTaskInput)
return stubCreatedTask(302, input as CreatedTaskInput)
},
createTaskEvent: async () => null,
logInfo: () => null,
logWarn: () => null,
nowIso: () => '2026-05-29T09:05:46.658Z',
randomId: () => 'DT-quantity',
},
)
assert.deepEqual(
createdInputs.map((input) => input.unitIndex),
[2],
)
assert.deepEqual(
tasks.map((task) => task.unit_index),
[1, 2],
)
})
test('syncDeliveryTasksForOrder 规划失败时记录可定位的路由诊断', async () => {
const warnings: Array<{ message: string; detail: Record<string, unknown> }> = []
const tasks = await syncDeliveryTasksForOrderWithDeps(
stubOrder({ id: 5, platform_order_id: 'planning-failed' }),
[
stubOrderItem({
id: 40,
order_id: 5,
sku_code: 'SKU-broken',
item_snapshot_json: {
isConfigured: true,
fulfillmentRoute: {
selectedExecutorKey: 'kuaishou_ct_assisted',
reason: '规则指定云履约',
},
cloudtentacles: {
cloudSkuId: 74,
cloudSkuName: '商品 A',
cloudSourceKeys: [],
},
} as JsonObject,
}),
],
{
listTasksByOrderId: async () => [],
getFulfillmentProfileByKey: async () => buildCloudProfile(),
createTaskEvent: async () => null,
logInfo: () => null,
logWarn: (_scope, message, detail) => {
warnings.push({
message: String(message),
detail: detail as Record<string, unknown>,
})
return null
},
},
)
assert.deepEqual(tasks, [])
assert.equal(warnings.length, 1)
assert.equal(warnings[0]?.message, '订单项已配置履约路由但无法生成任务计划')
assert.equal(warnings[0]?.detail.orderItemId, 40)
assert.deepEqual(warnings[0]?.detail.missingUnitIndexes, [1])
assert.deepEqual(warnings[0]?.detail.planningDiagnostic, {
snapshotIsConfigured: true,
selectedExecutorKey: 'kuaishou_ct_assisted',
routeReason: '规则指定云履约',
candidateAvailability: [],
sourceItemTitle: '',
sourceSkuNick: '',
cloudSkuId: 74,
cloudSkuName: '商品 A',
cloudSourceKeys: [],
feifeiProductCode: '',
affiliateDashSku: '',
})
})
test('syncDeliveryTasksForOrder 并发创建冲突时复用已落库任务', async () => {
const concurrentTask = stubCreatedTask(401, {
orderId: 6,
orderItemId: 50,
unitIndex: 1,
executorKey: 'kuaishou_ct_assisted',
})
let listCount = 0
let eventCount = 0
const tasks = await syncDeliveryTasksForOrderWithDeps(
stubOrder({ id: 6, platform_order_id: 'concurrent-order' }),
[
stubOrderItem({
id: 50,
order_id: 6,
sku_code: 'SKU-concurrent',
item_snapshot_json: buildCloudSnapshot(74, '商品 A'),
}),
],
{
listTasksByOrderId: async () => (listCount++ === 0 ? [] : [concurrentTask]),
getFulfillmentProfileByKey: async () => buildCloudProfile(),
createTask: async () => {
throw new Error('duplicate key value violates unique constraint')
},
createTaskEvent: async () => {
eventCount += 1
return null
},
logInfo: () => null,
logWarn: () => null,
nowIso: () => '2026-05-29T09:05:46.658Z',
randomId: () => 'DT-concurrent',
},
)
assert.deepEqual(
tasks.map((task) => task.id),
[401],
)
assert.equal(eventCount, 0)
})
function buildCloudSnapshot(cloudSkuId: number, cloudSkuName: string): JsonObject {
return {
isConfigured: true,
fulfillmentRoute: {
selectedExecutorKey: 'kuaishou_ct_assisted',
},
cloudtentacles: {
cloudSourceKeys: ['account-a'],
cloudSkuId,
cloudSkuName,
},
}
}
function buildCloudProfile() {
return {
id: 9,
profile_key: 'kuaishou_ct_assisted',
profile_name: 'kuaishou-lewan 履约',
executor_key: 'kuaishou_ct_assisted',
}
}
@@ -1,4 +1,5 @@
import { createTask, listTasksByOrderId, updateTask } from '../../repositories/task-repo.js'
import { createTaskEvent } from '../../repositories/task-event-repo.js'
import { getFulfillmentProfileByKey } from '../../repositories/fulfillment-profile-repo.js'
import { createTaskClaimToken } from '../claim/claim-service.js'
import { notifyTaskAutoManualReview } from '../notification/domain-notifications.js'
@@ -10,7 +11,9 @@ import { preparePaidFulfillmentTask } from '../fulfillment/executors/registry.js
import type { FulfillmentPrepareDeps } from '../fulfillment/executors/types.js'
import { nowIso } from '../../utils/time.js'
import { randomId } from '../../utils/random.js'
import { logInfo, logWarn } from '../../utils/logger.js'
import { TASK_STATUS, resolveInitialPaidTaskStatus } from '../../domain/task-status.js'
import type { JsonObject } from '../../types/json.js'
import type { OrderItemRow, OrderRow, TaskRow } from '../../types/repository/rows.js'
type ClaimTokenLike = {
@@ -37,6 +40,9 @@ type DeliveryTaskDeps = {
}) => Promise<unknown> | unknown
nowIso?: () => string
randomId?: (prefix?: string) => string
createTaskEvent?: typeof createTaskEvent
logInfo?: typeof logInfo
logWarn?: typeof logWarn
}
export async function syncDeliveryTasksForOrder(
@@ -60,6 +66,9 @@ export async function syncDeliveryTasksForOrderWithDeps(
notifyTaskAutoManualReview: notifyManualReview = notifyTaskAutoManualReview,
nowIso: getNowIso = nowIso,
randomId: createRandomId = randomId,
createTaskEvent: appendTaskEvent = createTaskEvent,
logInfo: writeInfoLog = logInfo,
logWarn: writeWarnLog = logWarn,
} = deps
const runtimeDeps: FulfillmentPrepareDeps = {
@@ -70,90 +79,216 @@ export async function syncDeliveryTasksForOrderWithDeps(
}
const existingTasks = await listTasks(order.id)
const tasks: DeliveryTaskRow[] = [...existingTasks]
const existingTaskKeys = new Set(existingTasks.map(buildTaskUnitKey))
if (existingTasks.length > 0) {
if (order.pay_status !== 'paid') {
return existingTasks
}
return preparePaidTasks(existingTasks, runtimeDeps)
}
const tasks: DeliveryTaskRow[] = []
writeInfoLog('[delivery-task-service]', '开始核对订单履约任务完整性', {
...buildOrderLogContext(order),
orderItemCount: orderItems.length,
existingTaskCount: existingTasks.length,
})
for (const item of orderItems) {
const plan = await planFulfillmentTaskForOrderItem({
order,
item,
getProfileByKey,
})
const quantity = Math.max(1, Number(item.quantity || 1))
const missingUnitIndexes = Array.from({ length: quantity }, (_, index) => index + 1).filter(
(unitIndex) => !existingTaskKeys.has(buildOrderItemUnitKey(item.id, unitIndex)),
)
if (!plan) {
if (missingUnitIndexes.length === 0) {
continue
}
const quantity = Math.max(1, Number(item.quantity || 1))
let plan
try {
plan = await planFulfillmentTaskForOrderItem({
order,
item,
getProfileByKey,
})
} catch (error) {
writeWarnLog('[delivery-task-service]', '订单项履约任务规划执行异常', {
...buildOrderLogContext(order),
...buildOrderItemLogContext(item),
missingUnitIndexes,
planningDiagnostic: buildPlanningDiagnostic(item),
error: error instanceof Error ? error.message : String(error),
})
throw error
}
for (let index = 0; index < quantity; index += 1) {
if (!plan) {
writeWarnLog('[delivery-task-service]', '订单项已配置履约路由但无法生成任务计划', {
...buildOrderLogContext(order),
...buildOrderItemLogContext(item),
missingUnitIndexes,
planningDiagnostic: buildPlanningDiagnostic(item),
})
continue
}
for (const unitIndex of missingUnitIndexes) {
const createdAt = getNowIso()
const initialStatus =
order.pay_status === 'paid'
? resolvePaidTaskStatus(plan.profile)
: TASK_STATUS.PENDING_PAYMENT
const task = await createDeliveryTask({
orderId: order.id,
orderItemId: item.id,
unitIndex: index + 1,
provider: order.provider,
platform: order.platform,
shopId: order.shop_id,
shopName: order.shop_name,
platformOrderId: order.platform_order_id,
taskNo: createRandomId('DT'),
profileId: plan.profileId,
executorKey: plan.executorKey,
taskStatus: initialStatus,
deliveryStatus: 'pending',
resultCode: '',
resultMessage: '',
claimToken: '',
claimExpiresAt: null,
automationMode: plan.autoDispatch ? 'automatic' : 'manual',
requiresClaim: plan.requiresClaim,
userActionStatus: plan.requiresClaim ? 'pending_claim' : 'not_required',
attemptCount: 0,
lastError: '',
contextJson: JSON.stringify(plan.context),
createdAt,
updatedAt: createdAt,
})
let createdFresh = true
let task: TaskRow | null
try {
task = await createDeliveryTask({
orderId: order.id,
orderItemId: item.id,
unitIndex,
provider: order.provider,
platform: order.platform,
shopId: order.shop_id,
shopName: order.shop_name,
platformOrderId: order.platform_order_id,
taskNo: createRandomId('DT'),
profileId: plan.profileId,
executorKey: plan.executorKey,
taskStatus: initialStatus,
deliveryStatus: 'pending',
resultCode: '',
resultMessage: '',
claimToken: '',
claimExpiresAt: null,
automationMode: plan.autoDispatch ? 'automatic' : 'manual',
requiresClaim: plan.requiresClaim,
userActionStatus: plan.requiresClaim ? 'pending_claim' : 'not_required',
attemptCount: 0,
lastError: '',
contextJson: JSON.stringify(plan.context),
createdAt,
updatedAt: createdAt,
})
} catch (error) {
const concurrentTasks = await listTasks(order.id)
task =
concurrentTasks.find(
(candidate) =>
buildTaskUnitKey(candidate) === buildOrderItemUnitKey(item.id, unitIndex),
) || null
if (!task) {
writeWarnLog('[delivery-task-service]', '履约任务创建执行异常', {
...buildOrderLogContext(order),
...buildOrderItemLogContext(item),
unitIndex,
profileId: plan.profileId,
profileKey: plan.profileKey,
executorKey: plan.executorKey,
error: error instanceof Error ? error.message : String(error),
})
throw error
}
createdFresh = false
writeInfoLog('[delivery-task-service]', '并发同步已创建相同履约任务,本次直接复用', {
...buildOrderLogContext(order),
...buildOrderItemLogContext(item),
unitIndex,
taskId: task.id,
taskNo: task.task_no,
})
}
if (task) {
tasks.push({
const deliveryTask = {
...task,
skuCode: item.sku_code,
skuName: item.sku_name,
}
tasks.push(deliveryTask)
existingTaskKeys.add(buildTaskUnitKey(task))
if (createdFresh) {
await appendTaskEvent(
task.id,
'fulfillment_task_planned',
{
source: 'order_task_reconciliation',
orderId: order.id,
platformOrderId: order.platform_order_id,
orderItemId: item.id,
unitIndex,
skuCode: item.sku_code,
skuName: item.sku_name,
profileId: plan.profileId,
profileKey: plan.profileKey,
executorKey: plan.executorKey,
automationMode: plan.autoDispatch ? 'automatic' : 'manual',
initialStatus,
},
createdAt,
)
}
} else {
writeWarnLog('[delivery-task-service]', '履约任务创建未返回记录', {
...buildOrderLogContext(order),
...buildOrderItemLogContext(item),
unitIndex,
profileId: plan.profileId,
profileKey: plan.profileKey,
executorKey: plan.executorKey,
})
}
}
}
writeInfoLog('[delivery-task-service]', '订单履约任务完整性核对完成', {
...buildOrderLogContext(order),
existingTaskCount: existingTasks.length,
createdTaskCount: tasks.length - existingTasks.length,
totalTaskCount: tasks.length,
expectedTaskCount: orderItems.reduce(
(total, item) => total + Math.max(1, Number(item.quantity || 1)),
0,
),
})
if (order.pay_status !== 'paid') {
return tasks
}
const preparedTasks = await Promise.all(
tasks.map((task) => preparePaidFulfillmentTask(task, runtimeDeps)),
tasks.map((task) => preparePaidTaskWithDiagnostics(order, task, runtimeDeps, writeWarnLog)),
)
return preparedTasks.filter(isTaskRow)
}
async function preparePaidTasks(tasks: TaskRow[], deps: FulfillmentPrepareDeps) {
const preparedTasks = await Promise.all(
tasks.map((task) => preparePaidFulfillmentTask(task, deps)),
)
return preparedTasks.filter(isTaskRow)
async function preparePaidTaskWithDiagnostics(
order: OrderRow,
task: TaskRow,
deps: FulfillmentPrepareDeps,
writeWarnLog: typeof logWarn,
) {
try {
const prepared = await preparePaidFulfillmentTask(task, deps)
if (!prepared) {
writeWarnLog('[delivery-task-service]', '已支付履约任务准备未返回任务记录', {
...buildOrderLogContext(order),
taskId: task.id,
taskNo: task.task_no,
orderItemId: task.order_item_id,
unitIndex: task.unit_index,
executorKey: task.executor_key,
taskStatus: task.task_status,
})
}
return prepared
} catch (error) {
writeWarnLog('[delivery-task-service]', '已支付履约任务准备执行异常', {
...buildOrderLogContext(order),
taskId: task.id,
taskNo: task.task_no,
orderItemId: task.order_item_id,
unitIndex: task.unit_index,
executorKey: task.executor_key,
taskStatus: task.task_status,
error: error instanceof Error ? error.message : String(error),
})
throw error
}
}
function resolvePaidTaskStatus(profile: FulfillmentBindingLike | null | undefined): string {
@@ -163,3 +298,71 @@ function resolvePaidTaskStatus(profile: FulfillmentBindingLike | null | undefine
function isTaskRow(task: TaskRow | DeliveryTaskRow | null | undefined): task is TaskRow {
return Boolean(task && Number(task.id || 0) > 0)
}
function buildTaskUnitKey(task: Pick<TaskRow, 'order_item_id' | 'unit_index'>) {
return buildOrderItemUnitKey(task.order_item_id, task.unit_index)
}
function buildOrderItemUnitKey(orderItemId: number | string, unitIndex: number | string) {
return `${Number(orderItemId)}:${Math.max(1, Number(unitIndex || 1))}`
}
function buildOrderLogContext(order: OrderRow) {
return {
orderId: order.id,
provider: order.provider,
platform: order.platform,
shopId: order.shop_id,
platformOrderId: order.platform_order_id,
payStatus: order.pay_status,
}
}
function buildOrderItemLogContext(item: OrderItemRow) {
return {
orderItemId: item.id,
skuCode: item.sku_code,
skuName: item.sku_name,
quantity: item.quantity,
}
}
function buildPlanningDiagnostic(item: OrderItemRow) {
const snapshot = parseJsonObject(item.item_snapshot_json)
const route = parseJsonObject(snapshot.fulfillmentRoute)
const cloudtentacles = parseJsonObject(snapshot.cloudtentacles)
const kuaishouFeifei = parseJsonObject(snapshot.kuaishouFeifei)
const affiliateDash = parseJsonObject(snapshot.affiliateDash)
const kuaishouSendCode = parseJsonObject(snapshot.kuaishouSendCode)
return {
snapshotIsConfigured: Boolean(snapshot.isConfigured),
selectedExecutorKey: String(route.selectedExecutorKey || '').trim(),
routeReason: String(route.reason || '').trim(),
candidateAvailability: route.candidates || route.candidateAvailability || [],
sourceItemTitle: String(kuaishouSendCode.itemTitle || '').trim(),
sourceSkuNick: String(kuaishouSendCode.skuNick || '').trim(),
cloudSkuId: Number(cloudtentacles.cloudSkuId || 0) || 0,
cloudSkuName: String(cloudtentacles.cloudSkuName || '').trim(),
cloudSourceKeys: Array.isArray(cloudtentacles.cloudSourceKeys)
? cloudtentacles.cloudSourceKeys
: [],
feifeiProductCode: String(kuaishouFeifei.productCode || '').trim(),
affiliateDashSku: String(affiliateDash.sku || '').trim(),
}
}
function parseJsonObject(value: unknown): JsonObject {
if (value && typeof value === 'object' && !Array.isArray(value)) {
return value as JsonObject
}
try {
const parsed = JSON.parse(String(value || '{}'))
return parsed && typeof parsed === 'object' && !Array.isArray(parsed)
? (parsed as JsonObject)
: {}
} catch {
return {}
}
}
@@ -17,7 +17,7 @@ import {
publishWorkOrderRealtimeChange,
} from '../realtime/realtime-event-service.js'
import { nowIso } from '../../utils/time.js'
import { logIntegration } from '../../utils/logger.js'
import { logInfo, logIntegration, logWarn } from '../../utils/logger.js'
import { createHttpError } from '../../utils/http.js'
import type { OrderItemRow, OrderRow, TaskRow } from '../../types/repository/rows.js'
@@ -210,6 +210,22 @@ export async function upsertOrderFromSource(
source: `${sourceLabel}_order_upsert`,
autoOnly: true,
})
logInfo('[order-consumer-router]', '来源订单消费者分流完成', {
source: sourceLabel,
orderId: order.id,
platformOrderId: order.platform_order_id,
fulfillment: {
result: 'blocked',
reason: readiness.blockReason,
detail: readiness.reason,
orderItemIds: configuredOrderItems.map((item) => item.id),
},
workerPlatform: {
createdCount: workOrderSync.createdCount,
skippedCount: workOrderSync.skippedCount,
decisions: workOrderSync.decisions,
},
})
publishSourceOrderRealtimeChanges(order.id, workOrderSync.created)
return {
ignoreReason: readiness.blockReason || 'fulfillment_not_ready',
@@ -231,6 +247,16 @@ export async function upsertOrderFromSource(
source: `${sourceLabel}_order_upsert`,
autoOnly: true,
})
logOrderConsumerOutcome({
sourceLabel,
order,
orderItems,
configuredOrderItems,
tasks,
workerDecisions: workOrderSync.decisions,
workerCreatedCount: workOrderSync.createdCount,
workerSkippedCount: workOrderSync.skippedCount,
})
publishSourceOrderRealtimeChanges(order.id, workOrderSync.created)
logIntegration('[order-service]', `${sourceLabel} 订单 upsert 完成`, {
@@ -262,6 +288,124 @@ function publishSourceOrderRealtimeChanges(
}
}
function logOrderConsumerOutcome({
sourceLabel,
order,
orderItems,
configuredOrderItems,
tasks,
workerDecisions,
workerCreatedCount,
workerSkippedCount,
}: {
sourceLabel: string
order: OrderRow
orderItems: OrderItemRow[]
configuredOrderItems: OrderItemRow[]
tasks: TaskRow[]
workerDecisions: Array<Record<string, unknown>>
workerCreatedCount: number
workerSkippedCount: number
}) {
const configuredItemIds = new Set(configuredOrderItems.map((item) => Number(item.id)))
const taskCountByOrderItemId = new Map<number, number>()
for (const task of tasks) {
const orderItemId = Number(task.order_item_id || 0)
taskCountByOrderItemId.set(orderItemId, (taskCountByOrderItemId.get(orderItemId) || 0) + 1)
}
const missingFulfillmentItems = configuredOrderItems
.filter(
(item) =>
(taskCountByOrderItemId.get(Number(item.id)) || 0) <
Math.max(1, Number(item.quantity || 1)),
)
.map((item) => ({
orderItemId: item.id,
skuCode: item.sku_code,
skuName: item.sku_name,
expectedTaskCount: Math.max(1, Number(item.quantity || 1)),
actualTaskCount: taskCountByOrderItemId.get(Number(item.id)) || 0,
route: resolveOrderItemRouteDiagnostic(item),
}))
const workerAcceptedItemIds = new Set(
workerDecisions
.filter((decision) => decision.result === 'created' || decision.result === 'existing')
.map((decision) => Number(decision.orderItemId || 0))
.filter((orderItemId) => orderItemId > 0),
)
const dualConsumedOrderItemIds = Array.from(workerAcceptedItemIds).filter((orderItemId) =>
configuredItemIds.has(orderItemId),
)
const unhandledOrderItems = orderItems
.filter((item) => {
const orderItemId = Number(item.id)
return !configuredItemIds.has(orderItemId) && !workerAcceptedItemIds.has(orderItemId)
})
.map((item) => ({
orderItemId: item.id,
skuCode: item.sku_code,
skuName: item.sku_name,
}))
const detail = {
source: sourceLabel,
orderId: order.id,
platformOrderId: order.platform_order_id,
fulfillment: {
configuredOrderItemIds: Array.from(configuredItemIds),
taskIds: tasks.map((task) => task.id),
taskCount: tasks.length,
missingItems: missingFulfillmentItems,
},
workerPlatform: {
createdCount: workerCreatedCount,
skippedCount: workerSkippedCount,
decisions: workerDecisions,
},
dualConsumedOrderItemIds,
unhandledOrderItems,
}
logInfo('[order-consumer-router]', '来源订单消费者分流完成', detail)
if (missingFulfillmentItems.length > 0) {
logWarn('[order-consumer-router]', '已配置履约的订单项存在任务缺口', detail)
}
if (dualConsumedOrderItemIds.length > 0) {
logWarn('[order-consumer-router]', '同一订单项被履约与发单平台消费者同时接收', detail)
}
if (unhandledOrderItems.length > 0) {
logWarn('[order-consumer-router]', '订单项未被履约或发单平台消费者接收', detail)
}
}
function resolveOrderItemRouteDiagnostic(item: OrderItemRow) {
const snapshot = parseOrderItemSnapshot(item.item_snapshot_json) || {}
const sendCode =
snapshot.kuaishouSendCode &&
typeof snapshot.kuaishouSendCode === 'object' &&
!Array.isArray(snapshot.kuaishouSendCode)
? (snapshot.kuaishouSendCode as Record<string, unknown>)
: {}
const route =
snapshot.fulfillmentRoute &&
typeof snapshot.fulfillmentRoute === 'object' &&
!Array.isArray(snapshot.fulfillmentRoute)
? (snapshot.fulfillmentRoute as Record<string, unknown>)
: {}
return {
snapshotIsConfigured: snapshot.isConfigured === true,
selectedExecutorKey: String(route.selectedExecutorKey || '').trim(),
selectedRuleId: String(route.selectedRuleId || '').trim(),
reason: String(route.reason || '').trim(),
sourceItemTitle: String(sendCode.itemTitle || '').trim(),
sourceSkuNick: String(sendCode.skuNick || '').trim(),
candidates: Array.isArray(route.candidates) ? route.candidates : [],
skipped: Array.isArray(route.skipped) ? route.skipped : [],
}
}
const ORDER_STATUS_PRIORITY = {
created: 0,
paid: 1,
@@ -480,23 +480,29 @@ async function syncOpen91OrderAfterSendCallbackSuccess(
}
const tasks = await listTasksByOrderId(order.id)
if (tasks.length > 0) {
await bindKuaishouIndustryVouchersToOrderTasks(order, tasks, {
source,
now: normalizeTimestampIso(new Date().toISOString()),
})
return
}
try {
const result = await retryOpen91Order(order.id)
const synchronizedTasks = result.tasks.length > 0 ? result.tasks : tasks
if (synchronizedTasks.length > 0) {
await bindKuaishouIndustryVouchersToOrderTasks(order, synchronizedTasks, {
source,
now: normalizeTimestampIso(new Date().toISOString()),
})
}
logIntegration('[kuaishou-industry/send-code]', '电子凭证发货回调成功后已触发 91 订单重试', {
source,
oid: normalizedOid,
orderId: order.id,
taskCount: result.tasks.length,
previousTaskCount: tasks.length,
})
} catch (error) {
if (tasks.length > 0) {
await bindKuaishouIndustryVouchersToOrderTasks(order, tasks, {
source,
now: normalizeTimestampIso(new Date().toISOString()),
})
}
if (isOpen91OrderAwaitingFulfillmentConfig(error)) {
return
}
@@ -19,6 +19,7 @@ import type { JsonObject } from '../../types/json.js'
import type { OrderItemRow, OrderRow } from '../../types/repository/rows.js'
import { randomId } from '../../utils/random.js'
import { nowIso } from '../../utils/time.js'
import { logInfo, logWarn } from '../../utils/logger.js'
import { safeParseJson } from '../admin/admin-query-utils.js'
import { createWorkOrderMaterialAdminNotification } from '../admin/admin-notification-service.js'
import {
@@ -44,6 +45,7 @@ export async function syncWorkerOrdersForSourceOrder(
const productMatchConfig = getWorkerProductMatchConfig()
const created: WorkOrderRow[] = []
const skipped: Array<{ orderItemId: number; reason: string }> = []
const decisions: Array<Record<string, unknown>> = []
for (const item of orderItems) {
const match = resolveMatchingProductRuleDecision(order, item, rules, mappings, {
@@ -75,19 +77,44 @@ export async function syncWorkerOrdersForSourceOrder(
})
}
if (!match.rule) {
const reason = match.reason === 'ambiguous' ? 'rule_match_ambiguous' : 'rule_not_matched'
skipped.push({
orderItemId: Number(item.id),
reason: match.reason === 'ambiguous' ? 'rule_match_ambiguous' : 'rule_not_matched',
reason,
})
decisions.push(
buildWorkerConsumerDecision(item, {
result: 'skipped',
reason,
matchReason: match.reason,
candidates: match.candidates,
}),
)
continue
}
const rule = match.rule
if (options.autoOnly && !rule.auto_create) {
skipped.push({ orderItemId: Number(item.id), reason: 'auto_create_disabled' })
decisions.push(
buildWorkerConsumerDecision(item, {
result: 'skipped',
reason: 'auto_create_disabled',
ruleId: rule.id,
ruleKey: rule.rule_key,
}),
)
continue
}
if (await getWorkOrderByOrderItemId(item.id)) {
skipped.push({ orderItemId: Number(item.id), reason: 'already_exists' })
decisions.push(
buildWorkerConsumerDecision(item, {
result: 'existing',
reason: 'already_exists',
ruleId: rule.id,
ruleKey: rule.rule_key,
}),
)
continue
}
// 订单号跨订单唯一:同一订单的多个 SKU 件允许共享订单号(按件建单是既有设计),
@@ -100,6 +127,14 @@ export async function syncWorkerOrdersForSourceOrder(
})
) {
skipped.push({ orderItemId: Number(item.id), reason: 'platform_order_exists' })
decisions.push(
buildWorkerConsumerDecision(item, {
result: 'skipped',
reason: 'platform_order_exists',
ruleId: rule.id,
ruleKey: rule.rule_key,
}),
)
continue
}
@@ -161,9 +196,31 @@ export async function syncWorkerOrdersForSourceOrder(
requirementJson: JSON.stringify({ fields }),
now,
})
if (!workOrder) continue
if (!workOrder) {
skipped.push({ orderItemId: Number(item.id), reason: 'create_returned_empty' })
decisions.push(
buildWorkerConsumerDecision(item, {
result: 'failed',
reason: 'create_returned_empty',
ruleId: rule.id,
ruleKey: rule.rule_key,
}),
)
continue
}
created.push(workOrder)
decisions.push(
buildWorkerConsumerDecision(item, {
result: 'created',
reason: 'rule_matched',
ruleId: rule.id,
ruleKey: rule.rule_key,
workOrderId: workOrder.id,
workOrderNo: workOrder.work_order_no,
workOrderStatus: workOrder.status,
}),
)
await createWorkOrderEvent({
workOrderId: workOrder.id,
actorType: 'system',
@@ -187,10 +244,54 @@ export async function syncWorkerOrdersForSourceOrder(
}
}
logInfo('[worker-order-consumer]', '发单平台消费者处理完成', {
source: options.source || 'source_order',
orderId: order.id,
platformOrderId: order.platform_order_id,
orderItemCount: orderItems.length,
createdCount: created.length,
skippedCount: skipped.length,
decisions,
})
const failedDecisions = decisions.filter((decision) => decision.result === 'failed')
if (failedDecisions.length > 0) {
logWarn('[worker-order-consumer]', '发单平台消费者存在未成功落库的订单项', {
source: options.source || 'source_order',
orderId: order.id,
platformOrderId: order.platform_order_id,
decisions: failedDecisions,
})
}
return {
created: created.map((workOrder) => mapWorkOrderAdmin(workOrder)),
skipped,
createdCount: created.length,
skippedCount: skipped.length,
decisions,
}
}
function buildWorkerConsumerDecision(
item: OrderItemRow,
detail: Record<string, unknown>,
): Record<string, unknown> {
const snapshot = safeParseJson(item.item_snapshot_json)
const sendCode =
snapshot.kuaishouSendCode &&
typeof snapshot.kuaishouSendCode === 'object' &&
!Array.isArray(snapshot.kuaishouSendCode)
? (snapshot.kuaishouSendCode as JsonObject)
: {}
return {
orderItemId: Number(item.id),
skuCode: item.sku_code,
skuName: item.sku_name,
quantity: Number(item.quantity || 1),
sourceItemTitle: String(sendCode.itemTitle || '').trim(),
sourceSkuNick: String(sendCode.skuNick || '').trim(),
...detail,
}
}
@@ -7,8 +7,12 @@ export type AdminViewerSessionInput = {
export type AdminOrderListQueryInput = {
page?: number | string
pageSize?: number | string
keyword?: string
provider?: string
executorKey?: string
platformOrderId?: string
payStatus?: string
taskStatus?: string
skuCode?: string
dateFrom?: string
dateTo?: string
@@ -17,6 +21,8 @@ export type AdminOrderListQueryInput = {
export type AdminTaskListQueryInput = {
page?: number | string
pageSize?: number | string
keyword?: string
executorKey?: string
status?: string
platformOrderId?: string
taskNo?: string
@@ -1,8 +1,12 @@
export type OrderListQueryInput = {
page?: number
pageSize?: number
keyword?: string
provider?: string
executorKey?: string
platformOrderId?: string
payStatus?: string
taskStatus?: string
skuCode?: string
dateFrom?: string
dateTo?: string
@@ -48,6 +52,8 @@ export type OrderItemReplaceInput = {
export type TaskListQueryInput = {
page?: number
pageSize?: number
keyword?: string
executorKey?: string
status?: string
platformOrderId?: string
taskNo?: string